XRootD
Loading...
Searching...
No Matches
XrdCl::CopyProcess Class Reference

Copy the data from one point to another. More...

#include <XrdClCopyProcess.hh>

Collaboration diagram for XrdCl::CopyProcess:

Public Member Functions

 CopyProcess ()
 Constructor.
virtual ~CopyProcess ()
 Destructor.
XRootDStatus AddJob (const PropertyList &properties, PropertyList *results)
XRootDStatus Prepare ()
XRootDStatus Run (CopyProgressHandler *handler)
 Run the copy jobs.

Detailed Description

Copy the data from one point to another.

Definition at line 107 of file XrdClCopyProcess.hh.

Constructor & Destructor Documentation

◆ CopyProcess()

XrdCl::CopyProcess::CopyProcess ( )

Constructor.

Definition at line 209 of file XrdClCopyProcess.cc.

209 : pImpl( new CopyProcessImpl() )
210 {
211 }

◆ ~CopyProcess()

XrdCl::CopyProcess::~CopyProcess ( )
virtual

Destructor.

Definition at line 216 of file XrdClCopyProcess.cc.

217 {
218 CleanUpJobs();
219 delete pImpl;
220 }

Member Function Documentation

◆ AddJob()

XRootDStatus XrdCl::CopyProcess::AddJob ( const PropertyList & properties,
PropertyList * results )

Add job

Parameters
propertiesjob configuration parameters
resultsplaceholder for the results

Configuration properties: source [string] - original source URL target [string] - target directory or file sourceLimit [uint32_t] - maximum number sources force [bool] - overwrite target if exists posc [bool] - persistify only on successful close coerce [bool] - ignore locking semantics on destination makeDir [bool] - create path to the file if it doesn't exist thirdParty [string] - "first" try third party copy, if it fails try normal copy; "only" only try third party copy checkSumMode [string] - "none" - no checksumming "end2end" - end to end checksumming "source" - calculate checksum at source "target" - calculate checksum at target checkSumType [string] - type of the checksum to be used checkSumPreset [string] - checksum preset chunkSize [uint32_t] - size of a copy chunks in bytes parallelChunks [uint8_t] - number of chunks that should be requested in parallel initTimeout [time_t] - time limit for successfull initialization of the copy job tpcTimeout [time_t] - time limit for the actual copy to finish dynamicSource [bool] - support for the case where the size source file may change during reading process

Configuration job - this is a job that that is supposed to configure the copy process as a whole instead of adding a copy job:

jobType [string] - "configuration" - for configuraion parallel [uint8_t] - nomber of copy jobs to be run in parallel

Results: sourceCheckSum [string] - checksum at source, if requested targetCheckSum [string] - checksum at target, if requested size [uint64_t] - file size status [XRootDStatus] - status of the copy operation sources [vector<string>] - all sources used realTarget [string] - the actual disk server target

Definition at line 225 of file XrdClCopyProcess.cc.

227 {
228 Env *env = DefaultEnv::GetEnv();
229
230 //--------------------------------------------------------------------------
231 // Process a configuraion job
232 //--------------------------------------------------------------------------
233 if( properties.HasProperty( "jobType" ) &&
234 properties.Get<std::string>( "jobType" ) == "configuration" )
235 {
236 if( pImpl->pJobProperties.size() > 0 &&
237 pImpl->pJobProperties.rbegin()->HasProperty( "jobType" ) &&
238 pImpl->pJobProperties.rbegin()->Get<std::string>( "jobType" ) == "configuration" )
239 {
240 PropertyList &config = *pImpl->pJobProperties.rbegin();
241 PropertyList::PropertyMap::const_iterator it;
242 for( it = properties.begin(); it != properties.end(); ++it )
243 config.Set( it->first, it->second );
244 }
245 else
246 pImpl->pJobProperties.push_back( properties );
247 return XRootDStatus();
248 }
249
250 //--------------------------------------------------------------------------
251 // Validate properties
252 //--------------------------------------------------------------------------
253 if( !properties.HasProperty( "source" ) )
254 return XRootDStatus( stError, errInvalidArgs, 0, "source not specified" );
255
256 if( !properties.HasProperty( "target" ) )
257 return XRootDStatus( stError, errInvalidArgs, 0, "target not specified" );
258
259 pImpl->pJobProperties.push_back( properties );
260 PropertyList &p = pImpl->pJobProperties.back();
261
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 )
266 if( !p.HasProperty( bools[i] ) )
267 p.Set( bools[i], false );
268
269 if( !p.HasProperty( "thirdParty" ) )
270 p.Set( "thirdParty", "none" );
271
272 if( !p.HasProperty( "thirdPartyMode" ) )
273 p.Set( "thirdPartyMode", "pull" );
274
275 if( !p.HasProperty( "checkSumMode" ) )
276 p.Set( "checkSumMode", "none" );
277 else
278 {
279 if( !p.HasProperty( "checkSumType" ) )
280 {
281 pImpl->pJobProperties.pop_back();
283 "checkSumType not specified" );
284 }
285 else
286 {
287 //----------------------------------------------------------------------
288 // Checksum type has to be case insensitive
289 //----------------------------------------------------------------------
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 );
295 }
296 }
297
298 if( !p.HasProperty( "parallelChunks" ) )
299 {
300 int val = DefaultCPParallelChunks;
301 env->GetInt( "CPParallelChunks", val );
302 p.Set( "parallelChunks", val );
303 }
304
305 if( !p.HasProperty( "chunkSize" ) )
306 {
307 int val = DefaultCPChunkSize;
308 env->GetInt( "CPChunkSize", val );
309 p.Set( "chunkSize", val );
310 }
311
312 if( !p.HasProperty( "xcpBlockSize" ) )
313 {
314 int val = DefaultXCpBlockSize;
315 env->GetInt( "XCpBlockSize", val );
316 p.Set( "xcpBlockSize", val );
317 }
318
319 if( !p.HasProperty( "initTimeout" ) )
320 {
321 int val = DefaultCPInitTimeout;
322 env->GetInt( "CPInitTimeout", val );
323 p.Set( "initTimeout", val );
324 }
325
326 if( !p.HasProperty( "tpcTimeout" ) )
327 {
328 int val = DefaultCPTPCTimeout;
329 env->GetInt( "CPTPCTimeout", val );
330 p.Set( "tpcTimeout", val );
331 }
332
333 if( !p.HasProperty( "cpTimeout" ) )
334 {
335 int val = DefaultCPTimeout;
336 env->GetInt( "CPTimeout", val );
337 p.Set( "cpTimeout", val );
338 }
339
340 if( !p.HasProperty( "dynamicSource" ) )
341 p.Set( "dynamicSource", false );
342
343 if( !p.HasProperty( "xrate" ) )
344 p.Set( "xrate", 0 );
345
346 if( !p.HasProperty( "xrateThreshold" ) || p.Get<long long>( "xrateThreshold" ) == 0 )
347 {
348 int val = DefaultXRateThreshold;
349 env->GetInt( "XRateThreshold", val );
350 p.Set( "xrateThreshold", val );
351 }
352
353 //--------------------------------------------------------------------------
354 // Insert the properties
355 //--------------------------------------------------------------------------
356 Log *log = DefaultEnv::GetLog();
357 Utils::LogPropertyList( log, UtilityMsg, "Adding job with properties: %s",
358 p );
359 pImpl->pJobResults.push_back( results );
360 return XRootDStatus();
361 }
static Log * GetLog()
Get default log.
static Env * GetEnv()
Get default client environment.
static void LogPropertyList(Log *log, uint64_t topic, const char *format, const PropertyList &list)
Log property list.
XRootDStatus(uint16_t st=0, uint16_t code=0, uint32_t errN=0, const std::string &message="")
Constructor.
const int DefaultCPInitTimeout
const int DefaultXRateThreshold
const int DefaultCPChunkSize
const uint16_t stError
An error occurred that could potentially be retried.
const int DefaultCPParallelChunks
const int DefaultXCpBlockSize
const uint64_t UtilityMsg
const int DefaultCPTimeout
const uint16_t errInvalidArgs
const int DefaultCPTPCTimeout
XrdSysError Log
Definition XrdConfig.cc:113

References XrdCl::XRootDStatus::XRootDStatus(), XrdCl::PropertyList::begin(), XrdCl::DefaultCPChunkSize, XrdCl::DefaultCPInitTimeout, XrdCl::DefaultCPParallelChunks, XrdCl::DefaultCPTimeout, XrdCl::DefaultCPTPCTimeout, XrdCl::DefaultXCpBlockSize, XrdCl::DefaultXRateThreshold, XrdCl::PropertyList::end(), XrdCl::errInvalidArgs, XrdCl::PropertyList::Get(), XrdCl::DefaultEnv::GetEnv(), XrdCl::Env::GetInt(), XrdCl::DefaultEnv::GetLog(), XrdCl::PropertyList::HasProperty(), XrdCl::Utils::LogPropertyList(), XrdCl::PropertyList::Set(), XrdCl::stError, and XrdCl::UtilityMsg.

Referenced by DoCat(), and main().

Here is the call graph for this function:
Here is the caller graph for this function:

◆ Prepare()

XRootDStatus XrdCl::CopyProcess::Prepare ( )

Definition at line 366 of file XrdClCopyProcess.cc.

367 {
368 Log *log = DefaultEnv::GetLog();
369 std::vector<PropertyList>::iterator it;
370
371 log->Debug( UtilityMsg, "CopyProcess: %zu jobs to prepare",
372 pImpl->pJobProperties.size() );
373
374 std::map<std::string, uint32_t> targetFlags;
375 int i = 0;
376 for( it = pImpl->pJobProperties.begin(); it != pImpl->pJobProperties.end(); ++it, ++i )
377 {
378 PropertyList &props = *it;
379
380 if( props.HasProperty( "jobType" ) &&
381 props.Get<std::string>( "jobType" ) == "configuration" )
382 continue;
383
384 PropertyList *res = pImpl->pJobResults[i];
385 std::string tmp;
386
387 props.Get( "source", tmp );
388 URL source = tmp;
389 if( !source.IsValid() )
390 {
391 log->Error( UtilityMsg, "Invalid copy source URL: %s", obfuscateAuth( tmp ).c_str() );
392 return XRootDStatus( stError, errInvalidArgs, 0, "invalid source" );
393 }
394
395 //--------------------------------------------------------------------------
396 // Create a virtual redirector if it is a Metalink file
397 //--------------------------------------------------------------------------
398 if( source.IsMetalink() )
399 {
400 RedirectorRegistry &registry = RedirectorRegistry::Instance();
401 XRootDStatus st = registry.RegisterAndWait( source );
402 if( !st.IsOK() ) return st;
403 }
404
405 // handle UNZIP CGI
406 const URL::ParamsMap &cgi = source.GetParams();
407 URL::ParamsMap::const_iterator itr = cgi.find( "xrdcl.unzip" );
408 if( itr != cgi.end() )
409 {
410 props.Set( "zipArchive", true );
411 props.Set( "zipSource", itr->second );
412 }
413
414 props.Get( "target", tmp );
415 URL target = tmp;
416 if( !target.IsValid() )
417 {
418 log->Error( UtilityMsg, "Invalid copy target URL: %s", obfuscateAuth( tmp ).c_str() );
419 return XRootDStatus( stError, errInvalidArgs, 0, "invalid target" );
420 }
421
422 if( target.GetProtocol() != "stdio" )
423 {
424 // handle directories
425 bool targetIsDir = false;
426 props.Get( "targetIsDir", targetIsDir );
427
428 if( targetIsDir )
429 {
430 std::string path = target.GetPath() + '/';
431 std::string fn;
432
433 bool isZip = false;
434 props.Get( "zipArchive", isZip );
435 if( isZip )
436 {
437 props.Get( "zipSource", fn );
438 }
439 else if( source.IsMetalink() )
440 {
441 RedirectorRegistry &registry = XrdCl::RedirectorRegistry::Instance();
442 VirtualRedirector *redirector = registry.Get( source );
443 fn = redirector->GetTargetName();
444 }
445 else
446 {
447 fn = source.GetPath();
448 }
449
450 size_t pos = fn.rfind( '/' );
451 if( pos != std::string::npos )
452 fn = fn.substr( pos + 1 );
453 path += fn;
454 target.SetPath( path );
455 props.Set( "target", target.GetURL() );
456 }
457 }
458
459 bool tpc = false;
460 props.Get( "thirdParty", tmp );
461 if( tmp != "none" )
462 tpc = true;
463
464 //------------------------------------------------------------------------
465 // Check if we have all we need
466 //------------------------------------------------------------------------
467 if( source.GetProtocol() != "stdio" && source.GetPath().empty() )
468 {
469 log->Debug( UtilityMsg, "CopyProcess (job #%d): no source specified.",
470 i );
471 CleanUpJobs();
472 XRootDStatus st = XRootDStatus( stError, errInvalidArgs );
473 res->Set( "status", st );
474 return st;
475 }
476
477 if( target.GetProtocol() != "stdio" && target.GetPath().empty() )
478 {
479 log->Debug( UtilityMsg, "CopyProcess (job #%d): no target specified.",
480 i );
481 CleanUpJobs();
482 XRootDStatus st = XRootDStatus( stError, errInvalidArgs );
483 res->Set( "status", st );
484 return st;
485 }
486
487 //------------------------------------------------------------------------
488 // Check what kind of job we should do
489 //------------------------------------------------------------------------
490 CopyJob *job = 0;
491
492 if( tpc == true )
493 {
494 MarkTPC( props );
495 job = new TPFallBackCopyJob( i+1, &props, res );
496 }
497 else
498 job = new ClassicCopyJob( i+1, &props, res );
499
500 pImpl->pJobs.push_back( job );
501 }
502 return XRootDStatus();
503 }
std::string obfuscateAuth(const std::string &input)
ClassicCopyJob(uint32_t jobId, PropertyList *jobProperties, PropertyList *jobResults)
static RedirectorRegistry & Instance()
Returns reference to the single instance.
TPFallBackCopyJob(uint32_t jobId, PropertyList *jobProperties, PropertyList *jobResults)
Constructor.
std::map< std::string, std::string > ParamsMap
Definition XrdClURL.hh:33

References XrdCl::ClassicCopyJob::ClassicCopyJob(), XrdCl::TPFallBackCopyJob::TPFallBackCopyJob(), XrdCl::XRootDStatus::XRootDStatus(), XrdCl::Log::Debug(), XrdCl::errInvalidArgs, XrdCl::Log::Error(), XrdCl::PropertyList::Get(), XrdCl::RedirectorRegistry::Get(), XrdCl::DefaultEnv::GetLog(), XrdCl::URL::GetParams(), XrdCl::URL::GetPath(), XrdCl::URL::GetProtocol(), XrdCl::VirtualRedirector::GetTargetName(), XrdCl::URL::GetURL(), XrdCl::PropertyList::HasProperty(), XrdCl::RedirectorRegistry::Instance(), XrdCl::URL::IsMetalink(), XrdCl::Status::IsOK(), XrdCl::URL::IsValid(), obfuscateAuth(), XrdCl::RedirectorRegistry::RegisterAndWait(), XrdCl::PropertyList::Set(), XrdCl::URL::SetPath(), XrdCl::stError, and XrdCl::UtilityMsg.

Referenced by DoCat(), and main().

Here is the call graph for this function:
Here is the caller graph for this function:

◆ Run()

XRootDStatus XrdCl::CopyProcess::Run ( CopyProgressHandler * handler)

Run the copy jobs.

Definition at line 508 of file XrdClCopyProcess.cc.

509 {
510 //--------------------------------------------------------------------------
511 // Get the configuration
512 //--------------------------------------------------------------------------
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" )
517 {
518 PropertyList &config = *pImpl->pJobProperties.rbegin();
519 if( config.HasProperty( "parallel" ) )
520 parallelThreads = (uint8_t)config.Get<int>( "parallel" );
521 }
522
523 //--------------------------------------------------------------------------
524 // Run the show
525 //--------------------------------------------------------------------------
526 std::vector<CopyJob *>::iterator it;
527 uint32_t currentJob = 1;
528 uint32_t totalJobs = pImpl->pJobs.size();
529
530 //--------------------------------------------------------------------------
531 // Single thread
532 //--------------------------------------------------------------------------
533 if( parallelThreads == 1 )
534 {
535 XRootDStatus err;
536
537 for( it = pImpl->pJobs.begin(); it != pImpl->pJobs.end(); ++it )
538 {
539 QueuedCopyJob j( *it, progress, currentJob, totalJobs );
540 j.Run(0);
541
542 XRootDStatus st = (*it)->GetResults()->Get<XRootDStatus>( "status" );
543 if( err.IsOK() && !st.IsOK() )
544 {
545 err = st;
546 }
547 ++currentJob;
548 }
549
550 if( !err.IsOK() ) return err;
551 }
552 //--------------------------------------------------------------------------
553 // Multiple threads
554 //--------------------------------------------------------------------------
555 else
556 {
557 uint32_t workers = std::min( (uint32_t)parallelThreads,
558 (uint32_t)pImpl->pJobs.size() );
559 JobManager jm( workers );
560 jm.Initialize();
561 if( !jm.Start() )
562 return XRootDStatus( stError, errOSError, 0,
563 "Unable to start job manager" );
564
565 XrdSysSemaphore *sem = new XrdSysSemaphore(0);
566 std::vector<QueuedCopyJob*> queued;
567 for( it = pImpl->pJobs.begin(); it != pImpl->pJobs.end(); ++it )
568 {
569 QueuedCopyJob *j = new QueuedCopyJob( *it, progress, currentJob,
570 totalJobs, sem );
571
572 queued.push_back( j );
573 jm.QueueJob(j, 0);
574 ++currentJob;
575 }
576
577 std::vector<QueuedCopyJob*>::iterator itQ;
578 for( itQ = queued.begin(); itQ != queued.end(); ++itQ )
579 sem->Wait();
580 delete sem;
581
582 if( !jm.Stop() )
583 return XRootDStatus( stError, errOSError, 0,
584 "Unable to stop job manager" );
585 jm.Finalize();
586 for( itQ = queued.begin(); itQ != queued.end(); ++itQ )
587 delete *itQ;
588
589 for( it = pImpl->pJobs.begin(); it != pImpl->pJobs.end(); ++it )
590 {
591 XRootDStatus st = (*it)->GetResults()->Get<XRootDStatus>( "status" );
592 if( !st.IsOK() ) return st;
593 }
594 };
595 return XRootDStatus();
596 }
XrdSysSemaphore(int semval=1, const char *=0)
const uint16_t errOSError

References XrdSysSemaphore::XrdSysSemaphore(), XrdCl::XRootDStatus::XRootDStatus(), XrdCl::errOSError, XrdCl::JobManager::Finalize(), XrdCl::PropertyList::Get(), XrdCl::PropertyList::HasProperty(), XrdCl::JobManager::Initialize(), XrdCl::Status::IsOK(), XrdCl::JobManager::QueueJob(), XrdCl::JobManager::Start(), XrdCl::stError, XrdCl::JobManager::Stop(), and XrdSysSemaphore::Wait().

Referenced by DoCat(), and main().

Here is the call graph for this function:
Here is the caller graph for this function:

The documentation for this class was generated from the following files: