46 QueuedCopyJob( XrdCl::CopyJob *job,
47 XrdCl::CopyProgressHandler *progress,
50 XrdSysSemaphore *sem = 0 ):
51 pJob(job), pProgress(progress), pCurrentJob(currentJob),
52 pTotalJobs(totalJobs), pSem(sem),
65 virtual void Run(
void * )
74 pProgress->BeginJob( pCurrentJob, pTotalJobs,
80 XrdCl::Monitor::CopyBInfo i;
86 gettimeofday( &bTOD, 0 );
91 XrdCl::XRootDStatus st;
94 st = pJob->Run( pProgress );
101 pJob->GetResults()->Get(
"LastURL", url );
102 XrdCl::URL lastURL( url );
104 auto itr = cgi.find(
"tried" );
105 if( itr != cgi.end() )
107 std::string tried = itr->second;
108 if( tried[tried.size() - 1] !=
',' ) tried +=
',';
109 tried += lastURL.GetHostName();
110 cgi[
"tried"] = tried;
113 cgi[
"tried"] = lastURL.GetHostName();
115 std::string recoveryRedir;
116 pJob->GetResults()->Get(
"WrtRecoveryRedir", recoveryRedir );
117 XrdCl::URL recRedirURL( recoveryRedir );
120 pJob->GetProperties()->Get(
"target", target );
121 XrdCl::URL trgURL( target );
122 trgURL.SetHostName( recRedirURL.GetHostName() );
123 trgURL.SetPort( recRedirURL.GetPort() );
124 trgURL.SetProtocol( recRedirURL.GetProtocol() );
125 trgURL.SetParams( cgi );
126 pJob->GetProperties()->Set(
"target", trgURL.GetURL() );
136 if( !st.
IsOK() && pRetryCnt > 0 &&
141 if( pRetryPolicy ==
"continue" )
143 pJob->GetProperties()->Set(
"force",
false );
144 pJob->GetProperties()->Set(
"continue",
true );
148 pJob->GetProperties()->Set(
"force",
true );
149 pJob->GetProperties()->Set(
"continue",
false );
159 pJob->GetResults()->Set(
"status", st );
166 std::vector<std::string> sources;
167 pJob->GetResults()->Get(
"sources", sources );
168 XrdCl::Monitor::CopyEInfo i;
173 gettimeofday( &i.
eTOD, 0 );
179 pProgress->EndJob( pCurrentJob, pJob->GetResults() );
186 XrdCl::CopyJob *pJob;
187 XrdCl::CopyProgressHandler *pProgress;
188 uint32_t pCurrentJob;
190 XrdSysSemaphore *pSem;
193 std::string pRetryPolicy;
234 properties.
Get<std::string>(
"jobType" ) ==
"configuration" )
236 if( pImpl->pJobProperties.size() > 0 &&
237 pImpl->pJobProperties.rbegin()->HasProperty(
"jobType" ) &&
238 pImpl->pJobProperties.rbegin()->Get<std::string>(
"jobType" ) ==
"configuration" )
241 PropertyList::PropertyMap::const_iterator it;
242 for( it = properties.
begin(); it != properties.
end(); ++it )
243 config.
Set( it->first, it->second );
246 pImpl->pJobProperties.push_back( properties );
259 pImpl->pJobProperties.push_back( properties );
262 const char *bools[] = {
"target",
"force",
"posc",
"coerce",
"makeDir",
263 "zipArchive",
"xcp",
"preserveXAttr",
"rmOnBadCksum",
264 "continue",
"zipAppend",
"doServer", 0};
265 for(
int i = 0; bools[i]; ++i )
267 p.
Set( bools[i],
false );
270 p.
Set(
"thirdParty",
"none" );
273 p.
Set(
"thirdPartyMode",
"pull" );
276 p.
Set(
"checkSumMode",
"none" );
281 pImpl->pJobProperties.pop_back();
283 "checkSumType not specified" );
290 std::string checkSumType;
291 p.
Get(
"checkSumType", checkSumType );
292 std::transform(checkSumType.begin(), checkSumType.end(),
293 checkSumType.begin(), ::tolower);
294 p.
Set(
"checkSumType", checkSumType );
301 env->
GetInt(
"CPParallelChunks", val );
302 p.
Set(
"parallelChunks", val );
308 env->
GetInt(
"CPChunkSize", val );
309 p.
Set(
"chunkSize", val );
315 env->
GetInt(
"XCpBlockSize", val );
316 p.
Set(
"xcpBlockSize", val );
322 env->
GetInt(
"CPInitTimeout", val );
323 p.
Set(
"initTimeout", val );
329 env->
GetInt(
"CPTPCTimeout", val );
330 p.
Set(
"tpcTimeout", val );
336 env->
GetInt(
"CPTimeout", val );
337 p.
Set(
"cpTimeout", val );
341 p.
Set(
"dynamicSource",
false );
346 if( !p.
HasProperty(
"xrateThreshold" ) || p.
Get<
long long>(
"xrateThreshold" ) == 0 )
349 env->
GetInt(
"XRateThreshold", val );
350 p.
Set(
"xrateThreshold", val );
359 pImpl->pJobResults.push_back( results );
369 std::vector<PropertyList>::iterator it;
372 pImpl->pJobProperties.size() );
374 std::map<std::string, uint32_t> targetFlags;
376 for( it = pImpl->pJobProperties.begin(); it != pImpl->pJobProperties.end(); ++it, ++i )
381 props.
Get<std::string>(
"jobType" ) ==
"configuration" )
387 props.
Get(
"source", tmp );
402 if( !st.
IsOK() )
return st;
407 URL::ParamsMap::const_iterator itr = cgi.find(
"xrdcl.unzip" );
408 if( itr != cgi.end() )
410 props.
Set(
"zipArchive",
true );
411 props.
Set(
"zipSource", itr->second );
414 props.
Get(
"target", tmp );
425 bool targetIsDir =
false;
426 props.
Get(
"targetIsDir", targetIsDir );
430 std::string path = target.
GetPath() +
'/';
434 props.
Get(
"zipArchive", isZip );
437 props.
Get(
"zipSource", fn );
450 size_t pos = fn.rfind(
'/' );
451 if( pos != std::string::npos )
452 fn = fn.substr( pos + 1 );
460 props.
Get(
"thirdParty", tmp );
473 res->
Set(
"status", st );
483 res->
Set(
"status", st );
500 pImpl->pJobs.push_back( job );
513 uint8_t parallelThreads = 1;
514 if( pImpl->pJobProperties.size() > 0 &&
515 pImpl->pJobProperties.rbegin()->HasProperty(
"jobType" ) &&
516 pImpl->pJobProperties.rbegin()->Get<std::string>(
"jobType" ) ==
"configuration" )
520 parallelThreads = (uint8_t)config.
Get<
int>(
"parallel" );
526 std::vector<CopyJob *>::iterator it;
527 uint32_t currentJob = 1;
528 uint32_t totalJobs = pImpl->pJobs.size();
533 if( parallelThreads == 1 )
537 for( it = pImpl->pJobs.begin(); it != pImpl->pJobs.end(); ++it )
539 QueuedCopyJob j( *it, progress, currentJob, totalJobs );
550 if( !err.
IsOK() )
return err;
557 uint32_t workers = std::min( (uint32_t)parallelThreads,
558 (uint32_t)pImpl->pJobs.size() );
563 "Unable to start job manager" );
566 std::vector<QueuedCopyJob*> queued;
567 for( it = pImpl->pJobs.begin(); it != pImpl->pJobs.end(); ++it )
569 QueuedCopyJob *j =
new QueuedCopyJob( *it, progress, currentJob,
572 queued.push_back( j );
577 std::vector<QueuedCopyJob*>::iterator itQ;
578 for( itQ = queued.begin(); itQ != queued.end(); ++itQ )
584 "Unable to stop job manager" );
586 for( itQ = queued.begin(); itQ != queued.end(); ++itQ )
589 for( it = pImpl->pJobs.begin(); it != pImpl->pJobs.end(); ++it )
592 if( !st.
IsOK() )
return st;
598 void CopyProcess::CleanUpJobs()
600 std::vector<CopyJob*>::iterator itJ;
601 for( itJ = pImpl->
pJobs.begin(); itJ != pImpl->
pJobs.end(); ++itJ )
612 pImpl->pJobs.clear();
std::string obfuscateAuth(const std::string &input)
ClassicCopyJob(uint32_t jobId, PropertyList *jobProperties, PropertyList *jobResults)
const URL & GetSource() const
Get source.
virtual ~CopyProcess()
Destructor.
CopyProcess()
Constructor.
XRootDStatus Run(CopyProgressHandler *handler)
Run the copy jobs.
XRootDStatus AddJob(const PropertyList &properties, PropertyList *results)
Interface for copy progress notification.
static Monitor * GetMonitor()
Get the monitor object.
static Log * GetLog()
Get default log.
static Env * GetEnv()
Get default client environment.
bool GetString(const std::string &key, std::string &value)
bool GetInt(const std::string &key, int &value)
bool Finalize()
Finalize the job manager, clear the queues.
bool Start()
Start the workers.
bool Initialize()
Initialize the job manager.
bool Stop()
Stop the workers.
void QueueJob(Job *job, void *arg=0)
Add a job to be run.
Interface for a job to be run by the job manager.
void Error(uint64_t topic, const char *format,...)
Report an error.
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
TransferInfo transfer
The transfer in question.
@ EvCopyBeg
CopyBInfo: Copy operation started.
@ EvCopyEnd
CopyEInfo: Copy operation ended.
virtual void Event(EventCode evCode, void *evData)=0
A key-value pair map storing both keys and values as strings.
void Set(const std::string &name, const Item &value)
PropertyMap::const_iterator end() const
Get the end iterator.
bool Get(const std::string &name, Item &item) const
bool HasProperty(const std::string &name) const
Check if we now about the given name.
PropertyMap::const_iterator begin() const
Get the begin iterator.
Singleton access to URL to virtual redirector mapping.
static RedirectorRegistry & Instance()
Returns reference to the single instance.
void Release(const URL &url)
Release the virtual redirector associated with the given URL.
XRootDStatus RegisterAndWait(const URL &url)
Creates a new virtual redirector and registers it (sync).
VirtualRedirector * Get(const URL &url) const
Get a virtual redirector associated with the given URL.
TPFallBackCopyJob(uint32_t jobId, PropertyList *jobProperties, PropertyList *jobResults)
Constructor.
const std::string & GetPath() const
Get the path.
bool IsMetalink() const
Is it a URL to a metalink.
std::map< std::string, std::string > ParamsMap
std::string GetURL() const
Get the URL.
void SetPath(const std::string &path)
Set the path.
const ParamsMap & GetParams() const
Get the URL params.
const std::string & GetProtocol() const
Get the protocol.
bool IsValid() const
Is the url valid.
static void LogPropertyList(Log *log, uint64_t topic, const char *format, const PropertyList &list)
Log property list.
An interface for metadata redirectors.
virtual std::string GetTargetName() const =0
Gets the file name as specified in the metalink.
XRootDStatus(uint16_t st=0, uint16_t code=0, uint32_t errN=0, const std::string &message="")
Constructor.
XrdSysSemaphore(int semval=1, const char *=0)
const int DefaultCPInitTimeout
const int DefaultXRateThreshold
const uint16_t errOperationExpired
const int DefaultCPChunkSize
const uint16_t stError
An error occurred that could potentially be retried.
const int DefaultRetryWrtAtLBLimit
std::vector< PropertyList * > pJobResults
const int DefaultCPParallelChunks
const uint16_t errOSError
const int DefaultXCpBlockSize
const uint64_t UtilityMsg
const int DefaultCPTimeout
const uint16_t errInvalidArgs
std::vector< PropertyList > pJobProperties
const uint16_t errRetry
Try again for whatever reason.
const uint16_t errThresholdExceeded
const char *const DefaultCpRetryPolicy
const int DefaultCPTPCTimeout
std::vector< CopyJob * > pJobs
TransferInfo transfer
The transfer in question.
int sources
Number of sources used for the copy.
timeval bTOD
Copy start time.
const XRootDStatus * status
Status of the copy.
timeval eTOD
Copy end time.
const URL * target
URL of the target.
const URL * origin
URL of the origin.
uint16_t code
Error type, or additional hints on what to do.
bool IsOK() const
We're fine.
static bool IsSocketError(uint16_t code)