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

XRootD transport handler. More...

#include <XrdClXRootDTransport.hh>

Inheritance diagram for XrdCl::XRootDTransport:
Collaboration diagram for XrdCl::XRootDTransport:

Public Member Functions

 XRootDTransport ()
 Constructor.
 ~XRootDTransport ()
 Destructor.
virtual void DecFileInstCnt (AnyObject &channelData)
 Decrement file object instance count bound to this channel.
virtual void Disconnect (AnyObject &channelData, uint16_t subStreamId)
 The stream has been disconnected, do the cleanups.
virtual void FinalizeChannel (AnyObject &channelData)
 Finalize channel.
virtual URL GetBindPreference (const URL &url, AnyObject &channelData)
 Get bind preference for the next data stream.
virtual XRootDStatus GetBody (Message &message, Socket *socket)
virtual XRootDStatus GetHeader (Message &message, Socket *socket)
virtual XRootDStatus GetMore (Message &message, Socket *socket)
virtual Status GetSignature (Message *toSign, Message *&sign, AnyObject &channelData)
 Get signature for given message.
virtual Status GetSignature (Message *toSign, Message *&sign, XRootDChannelInfo *info)
 Get signature for given message.
virtual XRootDStatus HandShake (HandShakeData *handShakeData, AnyObject &channelData)
 HandShake.
virtual bool HandShakeDone (HandShakeData *handShakeData, AnyObject &channelData)
virtual void InitializeChannel (const URL &url, AnyObject &channelData)
 Initialize channel.
virtual Status IsStreamBroken (time_t inactiveTime, AnyObject &channelData)
virtual bool IsStreamTTLElapsed (time_t time, AnyObject &channelData)
 Check if the stream should be disconnected.
virtual uint32_t MessageReceived (Message &msg, uint16_t subStream, AnyObject &channelData)
 Check if the message invokes a stream action.
virtual void MessageSent (Message *msg, uint16_t subStream, uint32_t bytesSent, AnyObject &channelData)
 Notify the transport about a message having been sent.
virtual PathID Multiplex (Message *msg, AnyObject &channelData, PathID *hint=0)
virtual PathID MultiplexSubStream (Message *msg, AnyObject &channelData, PathID *hint=0)
virtual bool NeedControlConnection ()
virtual bool NeedEncryption (HandShakeData *handShakeData, AnyObject &channelData)
virtual Status Query (uint16_t query, AnyObject &result, AnyObject &channelData)
 Query the channel.
virtual uint16_t SubStreamNumber (AnyObject &channelData)
 Return a number of substreams per stream that should be created.
virtual void WaitBeforeExit ()
 Wait until the program can safely exit.
Public Member Functions inherited from XrdCl::TransportHandler
virtual ~TransportHandler ()

Static Public Member Functions

static void GenerateDescription (char *msg, std::ostringstream &o)
 Get the description of a message.
static void LogErrorResponse (const Message &msg)
 Log server error response.
static XRootDStatus MarshallRequest (char *msg)
 Marshal the outgoing message.
static XRootDStatus MarshallRequest (Message *msg)
 Marshal the outgoing message.
static uint16_t NbConnectedStrm (AnyObject &channelData)
 Number of currently connected data streams.
static void SetDescription (Message *msg)
 Get the description of a message.
static XRootDStatus UnMarchalStatusMore (Message &msg)
 Unmarshall the correction-segment of the status response for pgwrite.
static XRootDStatus UnMarshallBody (Message *msg, uint16_t reqType)
 Unmarshall the body of the incoming message.
static void UnMarshallHeader (Message &msg)
 Unmarshall the header incoming message.
static XRootDStatus UnMarshallRequest (Message *msg)
static XRootDStatus UnMarshalStatusBody (Message &msg, uint16_t reqType)
 Unmarshall the body of the status response.

Friends

struct PluginUnloadHandler

Additional Inherited Members

Public Types inherited from XrdCl::TransportHandler
enum  StreamAction {
  NoAction = 0x0000 ,
  DigestMsg = 0x0001 ,
  AbortStream = 0x0002 ,
  CloseStream = 0x0004 ,
  ResumeStream = 0x0008 ,
  HoldStream = 0x0010 ,
  RequestClose = 0x0020
}
 Stream actions that may be triggered by incoming control messages. More...

Detailed Description

XRootD transport handler.

Definition at line 47 of file XrdClXRootDTransport.hh.

Constructor & Destructor Documentation

◆ XRootDTransport()

XrdCl::XRootDTransport::XRootDTransport ( )

Constructor.

Definition at line 292 of file XrdClXRootDTransport.cc.

292 :
293 pSecUnloadHandler( new PluginUnloadHandler() )
294 {
295 }

References PluginUnloadHandler.

Referenced by XrdCl::TransportManager::TransportManager().

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

◆ ~XRootDTransport()

XrdCl::XRootDTransport::~XRootDTransport ( )

Destructor.

Definition at line 300 of file XrdClXRootDTransport.cc.

301 {
302 delete pSecUnloadHandler; pSecUnloadHandler = 0;
303 }

Member Function Documentation

◆ DecFileInstCnt()

void XrdCl::XRootDTransport::DecFileInstCnt ( AnyObject & channelData)
virtual

Decrement file object instance count bound to this channel.

Implements XrdCl::TransportHandler.

Definition at line 1840 of file XrdClXRootDTransport.cc.

1841 {
1842 XRootDChannelInfo *info = 0;
1843 channelData.Get( info );
1844 if( info->finstcnt.load( std::memory_order_relaxed ) > 0 )
1845 info->finstcnt.fetch_sub( 1, std::memory_order_relaxed );
1846 }

References XrdCl::XRootDChannelInfo::finstcnt, and XrdCl::AnyObject::Get().

Here is the call graph for this function:

◆ Disconnect()

void XrdCl::XRootDTransport::Disconnect ( AnyObject & channelData,
uint16_t subStreamId )
virtual

The stream has been disconnected, do the cleanups.

Implements XrdCl::TransportHandler.

Definition at line 1571 of file XrdClXRootDTransport.cc.

1573 {
1574 XRootDChannelInfo *info = 0;
1575 channelData.Get( info );
1576
1577 if (!info) {
1578 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1579 return;
1580 }
1581
1582 XrdSysMutexHelper scopedLock( info->mutex );
1583
1584 if( !info->stream.empty() )
1585 {
1586 XRootDStreamInfo &sInfo = info->stream[subStreamId];
1587 sInfo.status = XRootDStreamInfo::Disconnected;
1588 }
1589
1590 if( subStreamId == 0 )
1591 {
1592 CleanUpProtection( info );
1593 info->sidManager->ReleaseAllTimedOut();
1594 info->sentOpens.clear();
1595 info->sentCloses.clear();
1596 info->openFiles = 0;
1597 info->waitBarrier = 0;
1598 }
1599 }
static Log * GetLog()
Get default log.
void Error(uint64_t topic, const char *format,...)
Report an error.
Definition XrdClLog.cc:231
const uint64_t XRootDTransportMsg

References XrdCl::XRootDStreamInfo::Disconnected, XrdCl::Log::Error(), XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetLog(), XrdCl::XRootDChannelInfo::mutex, XrdCl::XRootDChannelInfo::openFiles, XrdCl::XRootDChannelInfo::sentCloses, XrdCl::XRootDChannelInfo::sentOpens, XrdCl::XRootDChannelInfo::sidManager, XrdCl::XRootDStreamInfo::status, XrdCl::XRootDChannelInfo::stream, XrdCl::XRootDChannelInfo::waitBarrier, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ FinalizeChannel()

void XrdCl::XRootDTransport::FinalizeChannel ( AnyObject & channelData)
virtual

Finalize channel.

Implements XrdCl::TransportHandler.

Definition at line 472 of file XrdClXRootDTransport.cc.

473 {
474 }

◆ GenerateDescription()

void XrdCl::XRootDTransport::GenerateDescription ( char * msg,
std::ostringstream & o )
static

Get the description of a message.

Definition at line 3034 of file XrdClXRootDTransport.cc.

3035 {
3036 Log *log = DefaultEnv::GetLog();
3037 if( log->GetLevel() < Log::ErrorMsg )
3038 return;
3039
3040 ClientRequestHdr *req = (ClientRequestHdr *)msg;
3041 switch( req->requestid )
3042 {
3043 //------------------------------------------------------------------------
3044 // kXR_open
3045 //------------------------------------------------------------------------
3046 case kXR_open:
3047 {
3048 ClientOpenRequest *sreq = (ClientOpenRequest *)msg;
3049 o << "kXR_open (";
3050 char *fn = GetDataAsString( msg );
3051 o << "file: " << fn << ", ";
3052 delete [] fn;
3053 o << "mode: 0" << std::setbase(8) << sreq->mode << ", ";
3054 o << std::setbase(10);
3055 o << "flags: ";
3056 if( sreq->options == 0 )
3057 o << "none ";
3058 else
3059 {
3060 if( sreq->options & kXR_compress )
3061 o << "kXR_compress ";
3062 if( sreq->options & kXR_delete )
3063 o << "kXR_delete ";
3064 if( sreq->options & kXR_force )
3065 o << "kXR_force ";
3066 if( sreq->options & kXR_mkpath )
3067 o << "kXR_mkpath ";
3068 if( sreq->options & kXR_new )
3069 o << "kXR_new ";
3070 if( sreq->options & kXR_nowait )
3071 o << "kXR_nowait ";
3072 if( sreq->options & kXR_open_apnd )
3073 o << "kXR_open_apnd ";
3074 if( sreq->options & kXR_open_read )
3075 o << "kXR_open_read ";
3076 if( sreq->options & kXR_open_updt )
3077 o << "kXR_open_updt ";
3078 if( sreq->options & kXR_open_wrto )
3079 o << "kXR_open_wrto ";
3080 if( sreq->options & kXR_posc )
3081 o << "kXR_posc ";
3082 if( sreq->options & kXR_prefname )
3083 o << "kXR_prefname ";
3084 if( sreq->options & kXR_refresh )
3085 o << "kXR_refresh ";
3086 if( sreq->options & kXR_4dirlist )
3087 o << "kXR_4dirlist ";
3088 if( sreq->options & kXR_replica )
3089 o << "kXR_replica ";
3090 if( sreq->options & kXR_seqio )
3091 o << "kXR_seqio ";
3092 if( sreq->options & kXR_async )
3093 o << "kXR_async ";
3094 if( sreq->options & kXR_retstat )
3095 o << "kXR_retstat ";
3096 }
3097 o << "flagt: ";
3098 if( sreq->optiont == 0 )
3099 o << "none ";
3100 else
3101 {
3102 if( sreq->optiont & kXR_dup )
3103 o << "kXR_dup ";
3104 if( sreq->options & kXR_samefs )
3105 o << "kXR_samefs ";
3106 }
3107 o << "fhtemplt: " << FileHandleToStr( sreq->fhtemplt );
3108 o << ")";
3109 break;
3110 }
3111
3112 //------------------------------------------------------------------------
3113 // kXR_clone
3114 //------------------------------------------------------------------------
3115 case kXR_clone:
3116 {
3117 ClientCloneRequest *sreq = (ClientCloneRequest *)msg;
3118 XrdProto::clone_list *dataChunk = (XrdProto::clone_list*)(msg + 24 );
3119 o << "kXR_clone ( ";
3120 o << "handle: " << FileHandleToStr( sreq->fhandle );
3121 o << std::setbase(10);
3122 o << " list [ ";
3123 for( size_t i = 0; i < req->dlen/sizeof(XrdProto::clone_list); ++i )
3124 {
3125 o << "(src_handle: ";
3126 o << FileHandleToStr( dataChunk[i].srcFH );
3127 o << ", ";
3128 o << std::setbase(10);
3129 o << "src_offset: " << dataChunk[i].srcOffs;
3130 o << ", src_length: " << dataChunk[i].srcLen;
3131 o << ", dst_offset: " << dataChunk[i].dstOffs << "); ";
3132 }
3133
3134 o << " ] )";
3135 break;
3136 }
3137
3138 //------------------------------------------------------------------------
3139 // kXR_close
3140 //------------------------------------------------------------------------
3141 case kXR_close:
3142 {
3143 ClientCloseRequest *sreq = (ClientCloseRequest *)msg;
3144 o << "kXR_close (";
3145 o << "handle: " << FileHandleToStr( sreq->fhandle );
3146 o << ")";
3147 break;
3148 }
3149
3150 //------------------------------------------------------------------------
3151 // kXR_stat
3152 //------------------------------------------------------------------------
3153 case kXR_stat:
3154 {
3155 ClientStatRequest *sreq = (ClientStatRequest *)msg;
3156 o << "kXR_stat (";
3157 if( sreq->dlen )
3158 {
3159 char *fn = GetDataAsString( msg );;
3160 o << "path: " << fn << ", ";
3161 delete [] fn;
3162 }
3163 else
3164 {
3165 o << "handle: " << FileHandleToStr( sreq->fhandle );
3166 o << ", ";
3167 }
3168 o << "flags: ";
3169 if( sreq->options == 0 )
3170 o << "none";
3171 else
3172 {
3173 if( sreq->options & kXR_vfs )
3174 o << "kXR_vfs";
3175 }
3176 o << ")";
3177 break;
3178 }
3179
3180 //------------------------------------------------------------------------
3181 // kXR_read
3182 //------------------------------------------------------------------------
3183 case kXR_read:
3184 {
3185 ClientReadRequest *sreq = (ClientReadRequest *)msg;
3186 o << "kXR_read (";
3187 o << "handle: " << FileHandleToStr( sreq->fhandle );
3188 o << std::setbase(10);
3189 o << ", ";
3190 o << "offset: " << sreq->offset << ", ";
3191 o << "size: " << sreq->rlen << ")";
3192 break;
3193 }
3194
3195 //------------------------------------------------------------------------
3196 // kXR_pgread
3197 //------------------------------------------------------------------------
3198 case kXR_pgread:
3199 {
3200 ClientPgReadRequest *sreq = (ClientPgReadRequest *)msg;
3201 o << "kXR_pgread (";
3202 o << "handle: " << FileHandleToStr( sreq->fhandle );
3203 o << std::setbase(10);
3204 o << ", ";
3205 o << "offset: " << sreq->offset << ", ";
3206 o << "size: " << sreq->rlen << ")";
3207 break;
3208 }
3209
3210 //------------------------------------------------------------------------
3211 // kXR_write
3212 //------------------------------------------------------------------------
3213 case kXR_write:
3214 {
3215 ClientWriteRequest *sreq = (ClientWriteRequest *)msg;
3216 o << "kXR_write (";
3217 o << "handle: " << FileHandleToStr( sreq->fhandle );
3218 o << std::setbase(10);
3219 o << ", ";
3220 o << "offset: " << sreq->offset << ", ";
3221 o << "size: " << sreq->dlen << ")";
3222 break;
3223 }
3224
3225 //------------------------------------------------------------------------
3226 // kXR_pgwrite
3227 //------------------------------------------------------------------------
3228 case kXR_pgwrite:
3229 {
3230 ClientPgWriteRequest *sreq = (ClientPgWriteRequest *)msg;
3231 o << "kXR_pgwrite (";
3232 o << "handle: " << FileHandleToStr( sreq->fhandle );
3233 o << std::setbase(10);
3234 o << ", ";
3235 o << "offset: " << sreq->offset << ", ";
3236 o << "size: " << sreq->dlen << ")";
3237 break;
3238 }
3239
3240 //------------------------------------------------------------------------
3241 // kXR_fattr
3242 //------------------------------------------------------------------------
3243 case kXR_fattr:
3244 {
3245 ClientFattrRequest *sreq = (ClientFattrRequest *)msg;
3246 int nattr = sreq->numattr;
3247 int options = sreq->options;
3248 o << "kXR_fattr";
3249 switch (sreq->subcode) {
3250 case kXR_fattrGet:
3251 o << "Get";
3252 break;
3253 case kXR_fattrSet:
3254 o << "Set";
3255 break;
3256 case kXR_fattrList:
3257 o << "List";
3258 break;
3259 case kXR_fattrDel:
3260 o << "Delete";
3261 break;
3262 default:
3263 o << " unknown subcode: " << sreq->subcode;
3264 break;
3265 }
3266 o << " (handle: " << FileHandleToStr( sreq->fhandle );
3267 o << std::setbase(10);
3268 if (nattr)
3269 o << ", numattr: " << nattr;
3270 if (options) {
3271 o << ", options: ";
3272 if (options & 0x01)
3273 o << "new";
3274 if (options & 0x10)
3275 o << "list values";
3276 }
3277 o << ", total size: " << req->dlen << ")";
3278 break;
3279 }
3280
3281 //------------------------------------------------------------------------
3282 // kXR_sync
3283 //------------------------------------------------------------------------
3284 case kXR_sync:
3285 {
3286 ClientSyncRequest *sreq = (ClientSyncRequest *)msg;
3287 o << "kXR_sync (";
3288 o << "handle: " << FileHandleToStr( sreq->fhandle );
3289 o << ")";
3290 break;
3291 }
3292
3293 //------------------------------------------------------------------------
3294 // kXR_truncate
3295 //------------------------------------------------------------------------
3296 case kXR_truncate:
3297 {
3298 ClientTruncateRequest *sreq = (ClientTruncateRequest *)msg;
3299 o << "kXR_truncate (";
3300 if( !sreq->dlen )
3301 o << "handle: " << FileHandleToStr( sreq->fhandle );
3302 else
3303 {
3304 char *fn = GetDataAsString( msg );
3305 o << "file: " << fn;
3306 delete [] fn;
3307 }
3308 o << std::setbase(10);
3309 o << ", ";
3310 o << "offset: " << sreq->offset;
3311 o << ")";
3312 break;
3313 }
3314
3315 //------------------------------------------------------------------------
3316 // kXR_readv
3317 //------------------------------------------------------------------------
3318 case kXR_readv:
3319 {
3320 unsigned char *fhandle = 0;
3321 o << "kXR_readv (";
3322
3323 o << "handle: ";
3324 readahead_list *dataChunk = (readahead_list*)(msg + 24 );
3325 fhandle = dataChunk[0].fhandle;
3326 if( fhandle )
3327 o << FileHandleToStr( fhandle );
3328 else
3329 o << "unknown";
3330 o << ", ";
3331 o << std::setbase(10);
3332 o << "chunks: [";
3333 uint64_t size = 0;
3334 for( size_t i = 0; i < req->dlen/sizeof(readahead_list); ++i )
3335 {
3336 size += dataChunk[i].rlen;
3337 o << "(offset: " << dataChunk[i].offset;
3338 o << ", size: " << dataChunk[i].rlen << "); ";
3339 }
3340 o << "], ";
3341 o << "total size: " << size << ")";
3342 break;
3343 }
3344
3345 //------------------------------------------------------------------------
3346 // kXR_writev
3347 //------------------------------------------------------------------------
3348 case kXR_writev:
3349 {
3350 unsigned char *fhandle = 0;
3351 o << "kXR_writev (";
3352
3353 XrdProto::write_list *wrtList =
3354 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
3355 uint64_t size = 0;
3356 uint32_t numChunks = 0;
3357 for( size_t i = 0; i < req->dlen/sizeof(XrdProto::write_list); ++i )
3358 {
3359 fhandle = wrtList[i].fhandle;
3360 size += wrtList[i].wlen;
3361 ++numChunks;
3362 }
3363 o << "handle: ";
3364 if( fhandle )
3365 o << FileHandleToStr( fhandle );
3366 else
3367 o << "unknown";
3368 o << ", ";
3369 o << std::setbase(10);
3370 o << "chunks: " << numChunks << ", ";
3371 o << "total size: " << size << ")";
3372 break;
3373 }
3374
3375 //------------------------------------------------------------------------
3376 // kXR_locate
3377 //------------------------------------------------------------------------
3378 case kXR_locate:
3379 {
3380 ClientLocateRequest *sreq = (ClientLocateRequest *)msg;
3381 char *fn = GetDataAsString( msg );;
3382 o << "kXR_locate (";
3383 o << "path: " << fn << ", ";
3384 delete [] fn;
3385 o << "flags: ";
3386 if( sreq->options == 0 )
3387 o << "none";
3388 else
3389 {
3390 if( sreq->options & kXR_refresh )
3391 o << "kXR_refresh ";
3392 if( sreq->options & kXR_prefname )
3393 o << "kXR_prefname ";
3394 if( sreq->options & kXR_nowait )
3395 o << "kXR_nowait ";
3396 if( sreq->options & kXR_force )
3397 o << "kXR_force ";
3398 if( sreq->options & kXR_compress )
3399 o << "kXR_compress ";
3400 }
3401 o << ")";
3402 break;
3403 }
3404
3405 //------------------------------------------------------------------------
3406 // kXR_mv
3407 //------------------------------------------------------------------------
3408 case kXR_mv:
3409 {
3410 ClientMvRequest *sreq = (ClientMvRequest *)msg;
3411 o << "kXR_mv (";
3412 o << "source: ";
3413 o.write( msg + sizeof( ClientMvRequest ), sreq->arg1len );
3414 o << ", ";
3415 o << "destination: ";
3416 o.write( msg + sizeof( ClientMvRequest ) + sreq->arg1len + 1, sreq->dlen - sreq->arg1len - 1 );
3417 o << ")";
3418 break;
3419 }
3420
3421 //------------------------------------------------------------------------
3422 // kXR_query
3423 //------------------------------------------------------------------------
3424 case kXR_query:
3425 {
3426 ClientQueryRequest *sreq = (ClientQueryRequest *)msg;
3427 o << "kXR_query (";
3428 o << "code: ";
3429 switch( sreq->infotype )
3430 {
3431 case kXR_Qconfig: o << "kXR_Qconfig"; break;
3432 case kXR_Qckscan: o << "kXR_Qckscan"; break;
3433 case kXR_Qcksum: o << "kXR_Qcksum"; break;
3434 case kXR_Qopaque: o << "kXR_Qopaque"; break;
3435 case kXR_Qopaquf: o << "kXR_Qopaquf"; break;
3436 case kXR_Qopaqug: o << "kXR_Qopaqug"; break;
3437 case kXR_QPrep: o << "kXR_QPrep"; break;
3438 case kXR_Qspace: o << "kXR_Qspace"; break;
3439 case kXR_QStats: o << "kXR_QStats"; break;
3440 case kXR_Qvisa: o << "kXR_Qvisa"; break;
3441 case kXR_Qxattr: o << "kXR_Qxattr"; break;
3442 default: o << sreq->infotype; break;
3443 }
3444 o << ", ";
3445
3446 if( sreq->infotype == kXR_Qopaqug || sreq->infotype == kXR_Qvisa )
3447 {
3448 o << "handle: " << FileHandleToStr( sreq->fhandle );
3449 o << ", ";
3450 }
3451
3452 o << "arg length: " << sreq->dlen << ")";
3453 break;
3454 }
3455
3456 //------------------------------------------------------------------------
3457 // kXR_rm
3458 //------------------------------------------------------------------------
3459 case kXR_rm:
3460 {
3461 o << "kXR_rm (";
3462 char *fn = GetDataAsString( msg );;
3463 o << "path: " << fn << ")";
3464 delete [] fn;
3465 break;
3466 }
3467
3468 //------------------------------------------------------------------------
3469 // kXR_mkdir
3470 //------------------------------------------------------------------------
3471 case kXR_mkdir:
3472 {
3473 ClientMkdirRequest *sreq = (ClientMkdirRequest *)msg;
3474 o << "kXR_mkdir (";
3475 char *fn = GetDataAsString( msg );
3476 o << "path: " << fn << ", ";
3477 delete [] fn;
3478 o << "mode: 0" << std::setbase(8) << sreq->mode << ", ";
3479 o << std::setbase(10);
3480 o << "flags: ";
3481 if( sreq->options[0] == 0 )
3482 o << "none";
3483 else
3484 {
3485 if( sreq->options[0] & kXR_mkdirpath )
3486 o << "kXR_mkdirpath";
3487 }
3488 o << ")";
3489 break;
3490 }
3491
3492 //------------------------------------------------------------------------
3493 // kXR_rmdir
3494 //------------------------------------------------------------------------
3495 case kXR_rmdir:
3496 {
3497 o << "kXR_rmdir (";
3498 char *fn = GetDataAsString( msg );
3499 o << "path: " << fn << ")";
3500 delete [] fn;
3501 break;
3502 }
3503
3504 //------------------------------------------------------------------------
3505 // kXR_chmod
3506 //------------------------------------------------------------------------
3507 case kXR_chmod:
3508 {
3509 ClientChmodRequest *sreq = (ClientChmodRequest *)msg;
3510 o << "kXR_chmod (";
3511 char *fn = GetDataAsString( msg );
3512 o << "path: " << fn << ", ";
3513 delete [] fn;
3514 o << "mode: 0" << std::setbase(8) << sreq->mode << ")";
3515 break;
3516 }
3517
3518 //------------------------------------------------------------------------
3519 // kXR_ping
3520 //------------------------------------------------------------------------
3521 case kXR_ping:
3522 {
3523 o << "kXR_ping ()";
3524 break;
3525 }
3526
3527 //------------------------------------------------------------------------
3528 // kXR_protocol
3529 //------------------------------------------------------------------------
3530 case kXR_protocol:
3531 {
3532 ClientProtocolRequest *sreq = (ClientProtocolRequest *)msg;
3533 o << "kXR_protocol (";
3534 o << "clientpv: 0x" << std::setbase(16) << sreq->clientpv << ")";
3535 break;
3536 }
3537
3538 //------------------------------------------------------------------------
3539 // kXR_dirlist
3540 //------------------------------------------------------------------------
3541 case kXR_dirlist:
3542 {
3543 o << "kXR_dirlist (";
3544 char *fn = GetDataAsString( msg );;
3545 o << "path: " << fn << ")";
3546 delete [] fn;
3547 break;
3548 }
3549
3550 //------------------------------------------------------------------------
3551 // kXR_set
3552 //------------------------------------------------------------------------
3553 case kXR_set:
3554 {
3555 o << "kXR_set (";
3556 char *fn = GetDataAsString( msg );;
3557 o << "data: " << fn << ")";
3558 delete [] fn;
3559 break;
3560 }
3561
3562 //------------------------------------------------------------------------
3563 // kXR_prepare
3564 //------------------------------------------------------------------------
3565 case kXR_prepare:
3566 {
3567 ClientPrepareRequest *sreq = (ClientPrepareRequest *)msg;
3568 o << "kXR_prepare (";
3569 o << "flags: ";
3570
3571 if( sreq->options == 0 )
3572 o << "none";
3573 else
3574 {
3575 if( sreq->options & kXR_stage )
3576 o << "kXR_stage ";
3577 if( sreq->options & kXR_wmode )
3578 o << "kXR_wmode ";
3579 if( sreq->options & kXR_coloc )
3580 o << "kXR_coloc ";
3581 if( sreq->options & kXR_fresh )
3582 o << "kXR_fresh ";
3583 }
3584
3585 o << ", priority: " << (int) sreq->prty << ", ";
3586
3587 char *fn = GetDataAsString( msg );
3588 char *cursor;
3589 for( cursor = fn; *cursor; ++cursor )
3590 if( *cursor == '\n' ) *cursor = ' ';
3591
3592 o << "paths: " << fn << ")";
3593 delete [] fn;
3594 break;
3595 }
3596
3597 case kXR_chkpoint:
3598 {
3599 ClientChkPointRequest *sreq = (ClientChkPointRequest*)msg;
3600 o << "kXR_chkpoint (";
3601 o << "opcode: ";
3602 if( sreq->opcode == kXR_ckpBegin ) o << "kXR_ckpBegin)";
3603 else if( sreq->opcode == kXR_ckpCommit ) o << "kXR_ckpCommit)";
3604 else if( sreq->opcode == kXR_ckpQuery ) o << "kXR_ckpQuery)";
3605 else if( sreq->opcode == kXR_ckpRollback ) o << "kXR_ckpRollback)";
3606 else if( sreq->opcode == kXR_ckpXeq )
3607 {
3608 o << "kXR_ckpXeq) ";
3609 // In this case our request body will be one of kXR_pgwrite,
3610 // kXR_truncate, kXR_write, or kXR_writev request.
3611 GenerateDescription( msg + sizeof( ClientChkPointRequest ), o );
3612 }
3613
3614 break;
3615 }
3616
3617 //------------------------------------------------------------------------
3618 // Default
3619 //------------------------------------------------------------------------
3620 default:
3621 {
3622 o << "kXR_unknown (length: " << req->dlen << ")";
3623 break;
3624 }
3625 };
3626 }
kXR_int16 arg1len
Definition XProtocol.hh:460
@ kXR_fattrDel
Definition XProtocol.hh:300
@ kXR_fattrSet
Definition XProtocol.hh:303
@ kXR_fattrList
Definition XProtocol.hh:302
@ kXR_fattrGet
Definition XProtocol.hh:301
kXR_char fhandle[4]
Definition XProtocol.hh:565
kXR_char fhandle[4]
Definition XProtocol.hh:823
kXR_char fhandle[4]
Definition XProtocol.hh:848
kXR_char fhandle[4]
Definition XProtocol.hh:812
kXR_int32 dlen
Definition XProtocol.hh:461
kXR_char fhtemplt[4]
Definition XProtocol.hh:516
kXR_unt16 options
Definition XProtocol.hh:513
@ kXR_open_wrto
Definition XProtocol.hh:499
@ kXR_compress
Definition XProtocol.hh:482
@ kXR_async
Definition XProtocol.hh:488
@ kXR_delete
Definition XProtocol.hh:483
@ kXR_prefname
Definition XProtocol.hh:491
@ kXR_nowait
Definition XProtocol.hh:497
@ kXR_open_read
Definition XProtocol.hh:486
@ kXR_open_updt
Definition XProtocol.hh:487
@ kXR_mkpath
Definition XProtocol.hh:490
@ kXR_seqio
Definition XProtocol.hh:498
@ kXR_replica
Definition XProtocol.hh:495
@ kXR_posc
Definition XProtocol.hh:496
@ kXR_refresh
Definition XProtocol.hh:489
@ kXR_new
Definition XProtocol.hh:485
@ kXR_force
Definition XProtocol.hh:484
@ kXR_4dirlist
Definition XProtocol.hh:494
@ kXR_open_apnd
Definition XProtocol.hh:492
@ kXR_retstat
Definition XProtocol.hh:493
kXR_char fhandle[4]
Definition XProtocol.hh:543
kXR_unt16 optiont
Definition XProtocol.hh:514
kXR_char fhandle[4]
Definition XProtocol.hh:681
kXR_char fhandle[4]
Definition XProtocol.hh:695
kXR_char fhandle[4]
Definition XProtocol.hh:258
kXR_unt16 requestid
Definition XProtocol.hh:159
kXR_char fhandle[4]
Definition XProtocol.hh:669
@ kXR_read
Definition XProtocol.hh:126
@ kXR_open
Definition XProtocol.hh:123
@ kXR_writev
Definition XProtocol.hh:144
@ kXR_clone
Definition XProtocol.hh:145
@ kXR_readv
Definition XProtocol.hh:138
@ kXR_mkdir
Definition XProtocol.hh:121
@ kXR_sync
Definition XProtocol.hh:129
@ kXR_chmod
Definition XProtocol.hh:115
@ kXR_dirlist
Definition XProtocol.hh:117
@ kXR_fattr
Definition XProtocol.hh:133
@ kXR_rm
Definition XProtocol.hh:127
@ kXR_query
Definition XProtocol.hh:114
@ kXR_write
Definition XProtocol.hh:132
@ kXR_set
Definition XProtocol.hh:131
@ kXR_rmdir
Definition XProtocol.hh:128
@ kXR_truncate
Definition XProtocol.hh:141
@ kXR_protocol
Definition XProtocol.hh:119
@ kXR_mv
Definition XProtocol.hh:122
@ kXR_ping
Definition XProtocol.hh:124
@ kXR_stat
Definition XProtocol.hh:130
@ kXR_pgread
Definition XProtocol.hh:143
@ kXR_chkpoint
Definition XProtocol.hh:125
@ kXR_locate
Definition XProtocol.hh:140
@ kXR_close
Definition XProtocol.hh:116
@ kXR_pgwrite
Definition XProtocol.hh:139
@ kXR_prepare
Definition XProtocol.hh:134
kXR_int32 rlen
Definition XProtocol.hh:696
kXR_char options[1]
Definition XProtocol.hh:446
kXR_int64 offset
Definition XProtocol.hh:697
@ kXR_vfs
Definition XProtocol.hh:799
@ kXR_mkdirpath
Definition XProtocol.hh:440
@ kXR_wmode
Definition XProtocol.hh:625
@ kXR_fresh
Definition XProtocol.hh:627
@ kXR_coloc
Definition XProtocol.hh:626
@ kXR_stage
Definition XProtocol.hh:624
@ kXR_dup
Definition XProtocol.hh:503
@ kXR_samefs
Definition XProtocol.hh:504
@ kXR_QPrep
Definition XProtocol.hh:650
@ kXR_Qopaqug
Definition XProtocol.hh:661
@ kXR_Qconfig
Definition XProtocol.hh:655
@ kXR_Qopaquf
Definition XProtocol.hh:660
@ kXR_Qckscan
Definition XProtocol.hh:654
@ kXR_Qxattr
Definition XProtocol.hh:652
@ kXR_Qspace
Definition XProtocol.hh:653
@ kXR_Qvisa
Definition XProtocol.hh:656
@ kXR_QStats
Definition XProtocol.hh:649
@ kXR_Qcksum
Definition XProtocol.hh:651
@ kXR_Qopaque
Definition XProtocol.hh:659
kXR_char fhandle[4]
Definition XProtocol.hh:231
@ ErrorMsg
report errors
Definition XrdClLog.hh:109
static void GenerateDescription(char *msg, std::ostringstream &o)
Get the description of a message.
XrdSysError Log
Definition XrdConfig.cc:113
kXR_char fhandle[4]
Definition XProtocol.hh:873
kXR_char fhandle[4]
Definition XProtocol.hh:318

References ClientMvRequest::arg1len, ClientProtocolRequest::clientpv, ClientMvRequest::dlen, ClientPgWriteRequest::dlen, ClientQueryRequest::dlen, ClientRequestHdr::dlen, ClientStatRequest::dlen, ClientTruncateRequest::dlen, ClientWriteRequest::dlen, XrdProto::clone_list::dstOffs, XrdCl::Log::ErrorMsg, ClientCloneRequest::fhandle, ClientCloseRequest::fhandle, ClientFattrRequest::fhandle, ClientPgReadRequest::fhandle, ClientPgWriteRequest::fhandle, ClientQueryRequest::fhandle, ClientReadRequest::fhandle, ClientStatRequest::fhandle, ClientSyncRequest::fhandle, ClientTruncateRequest::fhandle, ClientWriteRequest::fhandle, readahead_list::fhandle, XrdProto::write_list::fhandle, ClientOpenRequest::fhtemplt, GenerateDescription(), XrdCl::Log::GetLevel(), XrdCl::DefaultEnv::GetLog(), ClientQueryRequest::infotype, kXR_4dirlist, kXR_async, kXR_chkpoint, kXR_chmod, kXR_clone, kXR_close, kXR_coloc, kXR_compress, kXR_delete, kXR_dirlist, kXR_dup, kXR_fattr, kXR_fattrDel, kXR_fattrGet, kXR_fattrList, kXR_fattrSet, kXR_force, kXR_fresh, kXR_locate, kXR_mkdir, kXR_mkdirpath, kXR_mkpath, kXR_mv, kXR_new, kXR_nowait, kXR_open, kXR_open_apnd, kXR_open_read, kXR_open_updt, kXR_open_wrto, kXR_pgread, kXR_pgwrite, kXR_ping, kXR_posc, kXR_prefname, kXR_prepare, kXR_protocol, kXR_Qckscan, kXR_Qcksum, kXR_Qconfig, kXR_Qopaque, kXR_Qopaquf, kXR_Qopaqug, kXR_QPrep, kXR_Qspace, kXR_QStats, kXR_query, kXR_Qvisa, kXR_Qxattr, kXR_read, kXR_readv, kXR_refresh, kXR_replica, kXR_retstat, kXR_rm, kXR_rmdir, kXR_samefs, kXR_seqio, kXR_set, kXR_stage, kXR_stat, kXR_sync, kXR_truncate, kXR_vfs, kXR_wmode, kXR_write, kXR_writev, ClientChmodRequest::mode, ClientMkdirRequest::mode, ClientOpenRequest::mode, ClientFattrRequest::numattr, ClientPgReadRequest::offset, ClientPgWriteRequest::offset, ClientReadRequest::offset, ClientTruncateRequest::offset, ClientWriteRequest::offset, readahead_list::offset, ClientChkPointRequest::opcode, ClientFattrRequest::options, ClientLocateRequest::options, ClientMkdirRequest::options, ClientOpenRequest::options, ClientPrepareRequest::options, ClientStatRequest::options, ClientOpenRequest::optiont, ClientPrepareRequest::prty, ClientRequestHdr::requestid, ClientPgReadRequest::rlen, ClientReadRequest::rlen, readahead_list::rlen, XrdProto::clone_list::srcLen, XrdProto::clone_list::srcOffs, ClientFattrRequest::subcode, and XrdProto::write_list::wlen.

Referenced by GenerateDescription(), and SetDescription().

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

◆ GetBindPreference()

URL XrdCl::XRootDTransport::GetBindPreference ( const URL & url,
AnyObject & channelData )
virtual

Get bind preference for the next data stream.

Implements XrdCl::TransportHandler.

Definition at line 1938 of file XrdClXRootDTransport.cc.

1940 {
1941 XRootDChannelInfo *info = 0;
1942 channelData.Get( info );
1943
1944 if(!info || !info->bindSelector)
1945 return url;
1946
1947 return URL( info->bindSelector->Get() );
1948 }
URL()
Default constructor.
Definition XrdClURL.cc:36

References XrdCl::URL::URL(), XrdCl::XRootDChannelInfo::bindSelector, and XrdCl::AnyObject::Get().

Here is the call graph for this function:

◆ GetBody()

XRootDStatus XrdCl::XRootDTransport::GetBody ( Message & message,
Socket * socket )
virtual

Read the message body from the socket, the socket is non-blocking, the method may be called multiple times - see GetHeader for details

Parameters
messagethe message buffer containing the header
socketthe socket
Returns
stOK & suDone if the whole message has been processed stOK & suRetry if more data is needed stError on failure

Implements XrdCl::TransportHandler.

Definition at line 348 of file XrdClXRootDTransport.cc.

349 {
350 //--------------------------------------------------------------------------
351 // Retrieve the body
352 //--------------------------------------------------------------------------
353 size_t leftToBeRead = 0;
354 uint32_t bodySize = 0;
355 ServerResponseHeader* rsphdr = (ServerResponseHeader*)message.GetBuffer();
356 bodySize = rsphdr->dlen;
357
358 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
360 "Response body too large." );
361
362 if( message.GetSize() < bodySize + 8 )
363 message.ReAllocate( bodySize + 8 );
364
365 leftToBeRead = bodySize-(message.GetCursor()-8);
366 while( leftToBeRead )
367 {
368 int bytesRead = 0;
369 XRootDStatus status = socket->Read( message.GetBufferAtCursor(), leftToBeRead, bytesRead );
370
371 if( !status.IsOK() || status.code == suRetry )
372 return status;
373
374 leftToBeRead -= bytesRead;
375 message.AdvanceCursor( bytesRead );
376 }
377
378 return XRootDStatus( stOK, suDone );
379 }
XRootDStatus(uint16_t st=0, uint16_t code=0, uint32_t errN=0, const std::string &message="")
Constructor.
const uint16_t suRetry
const uint16_t stError
An error occurred that could potentially be retried.
const uint16_t stOK
Everything went OK.
const uint16_t suDone
const uint16_t errInvalidMessage

References XrdCl::XRootDStatus::XRootDStatus(), XrdCl::Buffer::AdvanceCursor(), XrdCl::Status::code, ServerResponseHeader::dlen, XrdCl::errInvalidMessage, XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::GetBufferAtCursor(), XrdCl::Buffer::GetCursor(), XrdCl::Buffer::GetSize(), XrdCl::Status::IsOK(), XrdCl::Socket::Read(), XrdCl::Buffer::ReAllocate(), XrdCl::stError, XrdCl::stOK, XrdCl::suDone, and XrdCl::suRetry.

Here is the call graph for this function:

◆ GetHeader()

XRootDStatus XrdCl::XRootDTransport::GetHeader ( Message & message,
Socket * socket )
virtual

Read a message header from the socket, the socket is non-blocking, so if there is not enough data the function should return suRetry in which case it will be called again when more data arrives, with the data previously read stored in the message buffer

Parameters
messagethe message buffer
socketthe socket
Returns
stOK & suDone if the whole message has been processed stOK & suRetry if more data is needed stError on failure

Implements XrdCl::TransportHandler.

Definition at line 308 of file XrdClXRootDTransport.cc.

309 {
310 //--------------------------------------------------------------------------
311 // A new message - allocate the space needed for the header
312 //--------------------------------------------------------------------------
313 if( message.GetCursor() == 0 && message.GetSize() < 8 )
314 message.Allocate( 8 );
315
316 //--------------------------------------------------------------------------
317 // Read the message header
318 //--------------------------------------------------------------------------
319 if( message.GetCursor() < 8 )
320 {
321 size_t leftToBeRead = 8 - message.GetCursor();
322 while( leftToBeRead )
323 {
324 int bytesRead = 0;
325 XRootDStatus status = socket->Read( message.GetBufferAtCursor(),
326 leftToBeRead, bytesRead );
327 if( !status.IsOK() || status.code == suRetry )
328 return status;
329
330 leftToBeRead -= bytesRead;
331 message.AdvanceCursor( bytesRead );
332 }
333 UnMarshallHeader( message );
334
335 uint32_t bodySize = *(uint32_t*)(message.GetBuffer(4));
336 Log *log = DefaultEnv::GetLog();
337 log->Dump( XRootDTransportMsg, "[msg: %p] Expecting %d bytes of message "
338 "body", (void*)&message, bodySize );
339
340 return XRootDStatus( stOK, suDone );
341 }
343 }
static void UnMarshallHeader(Message &msg)
Unmarshall the header incoming message.
const uint16_t errInternal
Internal error.

References XrdCl::XRootDStatus::XRootDStatus(), XrdCl::Buffer::AdvanceCursor(), XrdCl::Buffer::Allocate(), XrdCl::Status::code, XrdCl::Log::Dump(), XrdCl::errInternal, XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::GetBufferAtCursor(), XrdCl::Buffer::GetCursor(), XrdCl::DefaultEnv::GetLog(), XrdCl::Buffer::GetSize(), XrdCl::Status::IsOK(), XrdCl::Socket::Read(), XrdCl::stError, XrdCl::stOK, XrdCl::suDone, XrdCl::suRetry, UnMarshallHeader(), and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ GetMore()

XRootDStatus XrdCl::XRootDTransport::GetMore ( Message & message,
Socket * socket )
virtual

Read more of the message body from the socket, the socket is non-blocking the method may be called multiple times - see GetHeader for details

Parameters
messagethe message buffer containing the header
socketthe socket
Returns
stOK & suDone if the whole message has been processed stOK & suRetry if more data is needed stError on failure

Implements XrdCl::TransportHandler.

Definition at line 384 of file XrdClXRootDTransport.cc.

385 {
386 ServerResponseHeader* rsphdr = (ServerResponseHeader*)message.GetBuffer();
387 if( rsphdr->status != kXR_status )
389
390 //--------------------------------------------------------------------------
391 // In case of non kXR_status responses we read all the response, including
392 // data. For kXR_status responses we first read only the remainder of the
393 // header. The header must then be unmarshalled, and then a second call to
394 // GetMore (repeated for suRetry as needed) will read the data.
395 //--------------------------------------------------------------------------
396
397 uint32_t bodySize = rsphdr->dlen;
398 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
400 "kXR_status: response body too large." );
401 if( bodySize+8 < sizeof( ServerResponseStatus ) )
403 "kXR_status: invalid message size." );
404
405 ServerResponseStatus *rspst = (ServerResponseStatus*)message.GetBuffer();
406 uint32_t moreSize = static_cast<uint32_t>( rspst->bdy.dlen );
407 if( moreSize > std::numeric_limits<uint32_t>::max() - 8 - bodySize )
409 "kXR_status: response body too large." );
410 bodySize += moreSize;
411
412 if( message.GetSize() < bodySize + 8 )
413 message.ReAllocate( bodySize + 8 );
414
415 size_t leftToBeRead = bodySize-(message.GetCursor()-8);
416 while( leftToBeRead )
417 {
418 int bytesRead = 0;
419 XRootDStatus status = socket->Read( message.GetBufferAtCursor(), leftToBeRead, bytesRead );
420
421 if( !status.IsOK() || status.code == suRetry )
422 return status;
423
424 leftToBeRead -= bytesRead;
425 message.AdvanceCursor( bytesRead );
426 }
427
428 // Unmarchal to message body
429 Log *log = DefaultEnv::GetLog();
430 XRootDStatus st = XRootDTransport::UnMarchalStatusMore( message );
431 if( !st.IsOK() && st.code == errDataError )
432 {
433 log->Error( XRootDTransportMsg, "[msg: %p] %s", (void*)&message,
434 st.GetErrorMessage().c_str() );
435 return st;
436 }
437
438 if( !st.IsOK() )
439 {
440 log->Error( XRootDTransportMsg, "[msg: %p] Failed to unmarshall status body.",
441 (void*)&message );
442 return st;
443 }
444
445 return XRootDStatus( stOK, suDone );
446 }
@ kXR_status
Definition XProtocol.hh:949
struct ServerResponseBody_Status bdy
static XRootDStatus UnMarchalStatusMore(Message &msg)
Unmarshall the correction-segment of the status response for pgwrite.
const uint16_t errDataError
data is corrupted
const uint16_t errInvalidOp

References XrdCl::XRootDStatus::XRootDStatus(), XrdCl::Buffer::AdvanceCursor(), ServerResponseStatus::bdy, XrdCl::Status::code, ServerResponseBody_Status::dlen, ServerResponseHeader::dlen, XrdCl::errDataError, XrdCl::errInvalidMessage, XrdCl::errInvalidOp, XrdCl::Log::Error(), XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::GetBufferAtCursor(), XrdCl::Buffer::GetCursor(), XrdCl::XRootDStatus::GetErrorMessage(), XrdCl::DefaultEnv::GetLog(), XrdCl::Buffer::GetSize(), XrdCl::Status::IsOK(), kXR_status, XrdCl::Socket::Read(), XrdCl::Buffer::ReAllocate(), ServerResponseHeader::status, XrdCl::stError, XrdCl::stOK, XrdCl::suDone, XrdCl::suRetry, UnMarchalStatusMore(), and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ GetSignature() [1/2]

Status XrdCl::XRootDTransport::GetSignature ( Message * toSign,
Message *& sign,
AnyObject & channelData )
virtual

Get signature for given message.

Implements XrdCl::TransportHandler.

Definition at line 1800 of file XrdClXRootDTransport.cc.

1801 {
1802 XRootDChannelInfo *info = 0;
1803 channelData.Get( info );
1804 return GetSignature( toSign, sign, info );
1805 }
virtual Status GetSignature(Message *toSign, Message *&sign, AnyObject &channelData)
Get signature for given message.

References XrdCl::AnyObject::Get(), and GetSignature().

Referenced by GetSignature().

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

◆ GetSignature() [2/2]

Status XrdCl::XRootDTransport::GetSignature ( Message * toSign,
Message *& sign,
XRootDChannelInfo * info )
virtual

Get signature for given message.

Definition at line 1810 of file XrdClXRootDTransport.cc.

1813 {
1814 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
1815 if( pSecUnloadHandler->unloaded ) return Status( stError, errInvalidOp );
1816
1817 ClientRequest *thereq = reinterpret_cast<ClientRequest*>( toSign->GetBuffer() );
1818 if( !info ) return Status( stError, errInternal );
1819 if( info->protection )
1820 {
1821 SecurityRequest *newreq = 0;
1822 // check if we have to secure the request in the first place
1823 if( !( NEED2SECURE ( info->protection )( *thereq ) ) ) return Status();
1824 // secure (sign/encrypt) the request
1825 int rc = info->protection->Secure( newreq, *thereq, 0 );
1826 // there was an error
1827 if( rc < 0 )
1828 return Status( stError, errInternal, -rc );
1829
1830 sign = new Message();
1831 sign->Grab( reinterpret_cast<char*>( newreq ), rc );
1832 }
1833
1834 return Status();
1835 }
#define NEED2SECURE(protP)
This class implements the XRootD protocol security protection.
Message(uint32_t size=0)
Constructor.
Status(uint16_t st=stOK, uint16_t cod=errNone, uint32_t errN=0)
Constructor.

References XrdCl::Message::Message(), XrdCl::Status::Status(), XrdCl::errInternal, XrdCl::errInvalidOp, XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::Grab(), NEED2SECURE, XrdCl::XRootDChannelInfo::protection, XrdSecProtect::Secure(), and XrdCl::stError.

Here is the call graph for this function:

◆ HandShake()

XRootDStatus XrdCl::XRootDTransport::HandShake ( HandShakeData * handShakeData,
AnyObject & channelData )
virtual

HandShake.

Implements XrdCl::TransportHandler.

Definition at line 479 of file XrdClXRootDTransport.cc.

481 {
482 XRootDChannelInfo *info = 0;
483 channelData.Get( info );
484
485 if (!info)
487
488 XrdSysMutexHelper scopedLock( info->mutex );
489
490 if( info->stream.size() <= handShakeData->subStreamId )
491 {
492 Log *log = DefaultEnv::GetLog();
493 log->Error( XRootDTransportMsg,
494 "[%s] Internal error: not enough substreams",
495 handShakeData->streamName.c_str() );
497 }
498
499 if( handShakeData->subStreamId == 0 )
500 {
501 info->streamName = handShakeData->streamName;
502 return HandShakeMain( handShakeData, channelData );
503 }
504 return HandShakeParallel( handShakeData, channelData );
505 }
const uint16_t stFatal
Fatal error, it's still an error.

References XrdCl::XRootDStatus::XRootDStatus(), XrdCl::errInternal, XrdCl::Log::Error(), XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetLog(), XrdCl::XRootDChannelInfo::mutex, XrdCl::stFatal, XrdCl::XRootDChannelInfo::stream, XrdCl::HandShakeData::streamName, XrdCl::XRootDChannelInfo::streamName, XrdCl::HandShakeData::subStreamId, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ HandShakeDone()

bool XrdCl::XRootDTransport::HandShakeDone ( HandShakeData * handShakeData,
AnyObject & channelData )
virtual

Implements XrdCl::TransportHandler.

Definition at line 758 of file XrdClXRootDTransport.cc.

760 {
761 XRootDChannelInfo *info = 0;
762 channelData.Get( info );
763
764 if (!info) {
766 "[%s] Internal error: no channel info",
767 handShakeData->streamName.c_str());
768 return false;
769 }
770
771 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
772 return ( sInfo.status == XRootDStreamInfo::Connected );
773 }

References XrdCl::XRootDStreamInfo::Connected, XrdCl::Log::Error(), XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetLog(), XrdCl::XRootDStreamInfo::status, XrdCl::XRootDChannelInfo::stream, XrdCl::HandShakeData::streamName, XrdCl::HandShakeData::subStreamId, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ InitializeChannel()

void XrdCl::XRootDTransport::InitializeChannel ( const URL & url,
AnyObject & channelData )
virtual

Initialize channel.

Implements XrdCl::TransportHandler.

Definition at line 451 of file XrdClXRootDTransport.cc.

453 {
454 XRootDChannelInfo *info = new XRootDChannelInfo( url );
455 XrdSysMutexHelper scopedLock( info->mutex );
456 channelData.Set( info );
457
458 Env *env = DefaultEnv::GetEnv();
459 int streams = DefaultSubStreamsPerChannel;
460 env->GetInt( "SubStreamsPerChannel", streams );
461 if( streams < 1 ) streams = 1;
462 info->stream.resize( streams );
463 info->strmSelector.reset( new StreamSelector( streams ) );
464 info->encrypted = url.IsSecure();
465 info->istpc = url.IsTPC();
466 info->logintoken = url.GetLoginToken();
467 }
static Env * GetEnv()
Get default client environment.
const int DefaultSubStreamsPerChannel

References XrdCl::StreamSelector::StreamSelector(), XrdCl::XRootDChannelInfo::XRootDChannelInfo(), XrdCl::DefaultSubStreamsPerChannel, XrdCl::XRootDChannelInfo::encrypted, XrdCl::DefaultEnv::GetEnv(), XrdCl::Env::GetInt(), XrdCl::URL::GetLoginToken(), XrdCl::URL::IsSecure(), XrdCl::URL::IsTPC(), XrdCl::XRootDChannelInfo::istpc, XrdCl::XRootDChannelInfo::logintoken, XrdCl::XRootDChannelInfo::mutex, XrdCl::AnyObject::Set(), XrdCl::XRootDChannelInfo::stream, and XrdCl::XRootDChannelInfo::strmSelector.

Here is the call graph for this function:

◆ IsStreamBroken()

Status XrdCl::XRootDTransport::IsStreamBroken ( time_t inactiveTime,
AnyObject & channelData )
virtual

Check the stream is broken - ie. TCP connection got broken and went undetected by the TCP stack

Implements XrdCl::TransportHandler.

Definition at line 831 of file XrdClXRootDTransport.cc.

833 {
834 XRootDChannelInfo *info = 0;
835 channelData.Get( info );
836 Env *env = DefaultEnv::GetEnv();
837 Log *log = DefaultEnv::GetLog();
838
839 if (!info) {
840 log->Error(XRootDTransportMsg,
841 "Internal error: no channel info, behaving as if stream is broken");
842 return true;
843 }
844
845 int streamTimeout = DefaultStreamTimeout;
846 env->GetInt( "StreamTimeout", streamTimeout );
847
848 XrdSysMutexHelper scopedLock( info->mutex );
849
850 const time_t now = time(0);
851 const bool anySID =
852 info->sidManager->IsAnySIDOldAs( now - streamTimeout );
853
854 log->Dump( XRootDTransportMsg, "[%s] Stream inactive since %lld seconds, "
855 "stream timeout: %d, any SID: %d, wait barrier: %s",
856 info->streamName.c_str(), (long long) inactiveTime, streamTimeout,
857 anySID, Utils::TimeToString(info->waitBarrier).c_str() );
858
859 if( inactiveTime < streamTimeout )
860 return Status();
861
862 if( now < info->waitBarrier )
863 return Status();
864
865 if( !anySID )
866 return Status();
867
869 }
static std::string TimeToString(time_t timestamp)
Convert timestamp to a string.
const uint16_t errSocketTimeout
const int DefaultStreamTimeout

References XrdCl::Status::Status(), XrdCl::DefaultStreamTimeout, XrdCl::Log::Dump(), XrdCl::Log::Error(), XrdCl::errSocketTimeout, XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetEnv(), XrdCl::Env::GetInt(), XrdCl::DefaultEnv::GetLog(), XrdCl::XRootDChannelInfo::mutex, XrdCl::XRootDChannelInfo::sidManager, XrdCl::stError, XrdCl::XRootDChannelInfo::streamName, XrdCl::Utils::TimeToString(), XrdCl::XRootDChannelInfo::waitBarrier, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ IsStreamTTLElapsed()

bool XrdCl::XRootDTransport::IsStreamTTLElapsed ( time_t time,
AnyObject & channelData )
virtual

Check if the stream should be disconnected.

Implements XrdCl::TransportHandler.

Definition at line 778 of file XrdClXRootDTransport.cc.

780 {
781 XRootDChannelInfo *info = 0;
782 channelData.Get( info );
783
784 Env *env = DefaultEnv::GetEnv();
785 Log *log = DefaultEnv::GetLog();
786
787 if (!info) {
788 log->Error(XRootDTransportMsg,
789 "Internal error: no channel info, behaving as if TTL has elapsed");
790 return true;
791 }
792
793 //--------------------------------------------------------------------------
794 // Check the TTL settings for the current server
795 //--------------------------------------------------------------------------
796 int ttl;
797 if( info->serverFlags & kXR_isServer )
798 {
800 env->GetInt( "DataServerTTL", ttl );
801 }
802 else
803 {
805 env->GetInt( "LoadBalancerTTL", ttl );
806 }
807
808 //--------------------------------------------------------------------------
809 // See whether we can give a go-ahead for the disconnection
810 //--------------------------------------------------------------------------
811 XrdSysMutexHelper scopedLock( info->mutex );
812 uint16_t allocatedSIDs = info->sidManager->GetNumberOfAllocatedSIDs();
813 log->Dump( XRootDTransportMsg, "[%s] Stream inactive since %lld seconds, "
814 "TTL: %d, allocated SIDs: %d, open files: %d, bound file objects: %d",
815 info->streamName.c_str(), (long long) inactiveTime, ttl, allocatedSIDs,
816 info->openFiles, info->finstcnt.load( std::memory_order_relaxed ) );
817
818 if( info->openFiles != 0 && info->finstcnt.load( std::memory_order_relaxed ) != 0 )
819 return false;
820
821 if( !allocatedSIDs && inactiveTime > ttl )
822 return true;
823
824 return false;
825 }
#define kXR_isServer
const int DefaultLoadBalancerTTL
const int DefaultDataServerTTL

References XrdCl::DefaultDataServerTTL, XrdCl::DefaultLoadBalancerTTL, XrdCl::Log::Dump(), XrdCl::Log::Error(), XrdCl::XRootDChannelInfo::finstcnt, XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetEnv(), XrdCl::Env::GetInt(), XrdCl::DefaultEnv::GetLog(), kXR_isServer, XrdCl::XRootDChannelInfo::mutex, XrdCl::XRootDChannelInfo::openFiles, XrdCl::XRootDChannelInfo::serverFlags, XrdCl::XRootDChannelInfo::sidManager, XrdCl::XRootDChannelInfo::streamName, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ LogErrorResponse()

void XrdCl::XRootDTransport::LogErrorResponse ( const Message & msg)
static

Log server error response.

Definition at line 1534 of file XrdClXRootDTransport.cc.

1535 {
1536 Log *log = DefaultEnv::GetLog();
1537 ServerResponse *rsp = (ServerResponse *)msg.GetBuffer();
1538 char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
1539 memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
1540 log->Error( XRootDTransportMsg, "Server responded with an error [%d]: %s",
1541 rsp->body.error.errnum, errmsg );
1542 delete [] errmsg;
1543 }
union ServerResponse::@040373375333017131300127053271011057331004327334 body
ServerResponseHeader hdr

References ServerResponse::body, ServerResponseHeader::dlen, XrdCl::Log::Error(), XrdCl::Buffer::GetBuffer(), XrdCl::DefaultEnv::GetLog(), ServerResponse::hdr, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ MarshallRequest() [1/2]

XRootDStatus XrdCl::XRootDTransport::MarshallRequest ( char * msg)
static

Marshal the outgoing message.

Definition at line 1115 of file XrdClXRootDTransport.cc.

1116 {
1117 ClientRequest *req = (ClientRequest*)msg;
1118 switch( req->header.requestid )
1119 {
1120 //------------------------------------------------------------------------
1121 // kXR_protocol
1122 //------------------------------------------------------------------------
1123 case kXR_protocol:
1124 req->protocol.clientpv = htonl( req->protocol.clientpv );
1125 break;
1126
1127 //------------------------------------------------------------------------
1128 // kXR_login
1129 //------------------------------------------------------------------------
1130 case kXR_login:
1131 req->login.pid = htonl( req->login.pid );
1132 break;
1133
1134 //------------------------------------------------------------------------
1135 // kXR_locate
1136 //------------------------------------------------------------------------
1137 case kXR_locate:
1138 req->locate.options = htons( req->locate.options );
1139 break;
1140
1141 //------------------------------------------------------------------------
1142 // kXR_query
1143 //------------------------------------------------------------------------
1144 case kXR_query:
1145 req->query.infotype = htons( req->query.infotype );
1146 break;
1147
1148 //------------------------------------------------------------------------
1149 // kXR_truncate
1150 //------------------------------------------------------------------------
1151 case kXR_truncate:
1152 req->truncate.offset = htonll( req->truncate.offset );
1153 break;
1154
1155 //------------------------------------------------------------------------
1156 // kXR_mkdir
1157 //------------------------------------------------------------------------
1158 case kXR_mkdir:
1159 req->mkdir.mode = htons( req->mkdir.mode );
1160 break;
1161
1162 //------------------------------------------------------------------------
1163 // kXR_chmod
1164 //------------------------------------------------------------------------
1165 case kXR_chmod:
1166 req->chmod.mode = htons( req->chmod.mode );
1167 break;
1168
1169 //------------------------------------------------------------------------
1170 // kXR_open
1171 //------------------------------------------------------------------------
1172 case kXR_open:
1173 req->open.mode = htons( req->open.mode );
1174 req->open.options = htons( req->open.options );
1175 req->open.optiont = htons( req->open.optiont );
1176 break;
1177
1178 //------------------------------------------------------------------------
1179 // kXR_read
1180 //------------------------------------------------------------------------
1181 case kXR_read:
1182 req->read.offset = htonll( req->read.offset );
1183 req->read.rlen = htonl( req->read.rlen );
1184 break;
1185
1186 //------------------------------------------------------------------------
1187 // kXR_write
1188 //------------------------------------------------------------------------
1189 case kXR_write:
1190 req->write.offset = htonll( req->write.offset );
1191 break;
1192
1193 //------------------------------------------------------------------------
1194 // kXR_mv
1195 //------------------------------------------------------------------------
1196 case kXR_mv:
1197 req->mv.arg1len = htons( req->mv.arg1len );
1198 break;
1199
1200 //------------------------------------------------------------------------
1201 // kXR_readv
1202 //------------------------------------------------------------------------
1203 case kXR_readv:
1204 {
1205 uint16_t numChunks = (req->readv.dlen)/16;
1206 readahead_list *dataChunk = (readahead_list*)( msg + 24 );
1207 for( size_t i = 0; i < numChunks; ++i )
1208 {
1209 dataChunk[i].rlen = htonl( dataChunk[i].rlen );
1210 dataChunk[i].offset = htonll( dataChunk[i].offset );
1211 }
1212 break;
1213 }
1214
1215 case kXR_clone:
1216 {
1217 uint32_t numChunks = (req->clone.dlen)/sizeof(XrdProto::clone_list);
1218 XrdProto::clone_list *dataChunk =
1219 (XrdProto::clone_list*)( msg + sizeof( ClientRequestHdr ) );
1220 for( size_t i = 0; i < numChunks; ++i )
1221 {
1222 dataChunk[i].srcOffs = htonll( dataChunk[i].srcOffs );
1223 dataChunk[i].srcLen = htonll( dataChunk[i].srcLen );
1224 dataChunk[i].dstOffs = htonll( dataChunk[i].dstOffs );
1225 }
1226 break;
1227 }
1228
1229 //------------------------------------------------------------------------
1230 // kXR_writev
1231 //------------------------------------------------------------------------
1232 case kXR_writev:
1233 {
1234 uint16_t numChunks = (req->writev.dlen)/16;
1235 XrdProto::write_list *wrtList =
1236 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
1237 for( size_t i = 0; i < numChunks; ++i )
1238 {
1239 wrtList[i].wlen = htonl( wrtList[i].wlen );
1240 wrtList[i].offset = htonll( wrtList[i].offset );
1241 }
1242
1243 break;
1244 }
1245
1246 case kXR_pgread:
1247 {
1248 req->pgread.offset = htonll( req->pgread.offset );
1249 req->pgread.rlen = htonl( req->pgread.rlen );
1250 break;
1251 }
1252
1253 case kXR_pgwrite:
1254 {
1255 req->pgwrite.offset = htonll( req->pgwrite.offset );
1256 break;
1257 }
1258
1259 //------------------------------------------------------------------------
1260 // kXR_prepare
1261 //------------------------------------------------------------------------
1262 case kXR_prepare:
1263 {
1264 req->prepare.optionX = htons( req->prepare.optionX );
1265 req->prepare.port = htons( req->prepare.port );
1266 break;
1267 }
1268
1269 case kXR_chkpoint:
1270 {
1271 if( req->chkpoint.opcode == kXR_ckpXeq )
1272 MarshallRequest( msg + 24 );
1273 break;
1274 }
1275 };
1276
1277 req->header.requestid = htons( req->header.requestid );
1278 req->header.dlen = htonl( req->header.dlen );
1279 return XRootDStatus();
1280 }
struct ClientTruncateRequest truncate
Definition XProtocol.hh:917
struct ClientPgReadRequest pgread
Definition XProtocol.hh:903
struct ClientMkdirRequest mkdir
Definition XProtocol.hh:900
struct ClientPgWriteRequest pgwrite
Definition XProtocol.hh:904
struct ClientReadVRequest readv
Definition XProtocol.hh:910
struct ClientOpenRequest open
Definition XProtocol.hh:902
struct ClientRequestHdr header
Definition XProtocol.hh:887
struct ClientWriteVRequest writev
Definition XProtocol.hh:919
struct ClientLoginRequest login
Definition XProtocol.hh:899
@ kXR_login
Definition XProtocol.hh:120
struct ClientChmodRequest chmod
Definition XProtocol.hh:891
struct ClientQueryRequest query
Definition XProtocol.hh:908
struct ClientReadRequest read
Definition XProtocol.hh:909
struct ClientMvRequest mv
Definition XProtocol.hh:901
struct ClientChkPointRequest chkpoint
Definition XProtocol.hh:890
struct ClientPrepareRequest prepare
Definition XProtocol.hh:906
struct ClientWriteRequest write
Definition XProtocol.hh:918
struct ClientProtocolRequest protocol
Definition XProtocol.hh:907
struct ClientLocateRequest locate
Definition XProtocol.hh:898
struct ClientCloneRequest clone
Definition XProtocol.hh:892
static XRootDStatus MarshallRequest(Message *msg)
Marshal the outgoing message.

References XrdCl::XRootDStatus::XRootDStatus(), ClientMvRequest::arg1len, ClientRequest::chkpoint, ClientRequest::chmod, ClientProtocolRequest::clientpv, ClientRequest::clone, ClientCloneRequest::dlen, ClientReadVRequest::dlen, ClientRequestHdr::dlen, ClientWriteVRequest::dlen, XrdProto::clone_list::dstOffs, ClientRequest::header, ClientQueryRequest::infotype, kXR_chkpoint, kXR_chmod, kXR_clone, kXR_locate, kXR_login, kXR_mkdir, kXR_mv, kXR_open, kXR_pgread, kXR_pgwrite, kXR_prepare, kXR_protocol, kXR_query, kXR_read, kXR_readv, kXR_truncate, kXR_write, kXR_writev, ClientRequest::locate, ClientRequest::login, MarshallRequest(), ClientRequest::mkdir, ClientChmodRequest::mode, ClientMkdirRequest::mode, ClientOpenRequest::mode, ClientRequest::mv, ClientPgReadRequest::offset, ClientPgWriteRequest::offset, ClientReadRequest::offset, ClientTruncateRequest::offset, ClientWriteRequest::offset, readahead_list::offset, XrdProto::write_list::offset, ClientChkPointRequest::opcode, ClientRequest::open, ClientLocateRequest::options, ClientOpenRequest::options, ClientOpenRequest::optiont, ClientPrepareRequest::optionX, ClientRequest::pgread, ClientRequest::pgwrite, ClientLoginRequest::pid, ClientPrepareRequest::port, ClientRequest::prepare, ClientRequest::protocol, ClientRequest::query, ClientRequest::read, ClientRequest::readv, ClientRequestHdr::requestid, ClientPgReadRequest::rlen, ClientReadRequest::rlen, readahead_list::rlen, XrdProto::clone_list::srcLen, XrdProto::clone_list::srcOffs, ClientRequest::truncate, XrdProto::write_list::wlen, ClientRequest::write, and ClientRequest::writev.

Here is the call graph for this function:

◆ MarshallRequest() [2/2]

XRootDStatus XrdCl::XRootDTransport::MarshallRequest ( Message * msg)
inlinestatic

Marshal the outgoing message.

Definition at line 175 of file XrdClXRootDTransport.hh.

176 {
177 MarshallRequest( msg->GetBuffer() );
178 msg->SetIsMarshalled( true );
179 return XRootDStatus();
180 }

References XrdCl::XRootDStatus::XRootDStatus(), XrdCl::Buffer::GetBuffer(), MarshallRequest(), and XrdCl::Message::SetIsMarshalled().

Referenced by MarshallRequest(), MarshallRequest(), MultiplexSubStream(), XrdCl::MessageUtils::RedirectMessage(), XrdCl::MessageUtils::SendMessage(), and UnMarshallRequest().

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

◆ MessageReceived()

uint32_t XrdCl::XRootDTransport::MessageReceived ( Message & msg,
uint16_t subStream,
AnyObject & channelData )
virtual

Check if the message invokes a stream action.

Implements XrdCl::TransportHandler.

Definition at line 1656 of file XrdClXRootDTransport.cc.

1659 {
1660 XRootDChannelInfo *info = 0;
1661 channelData.Get( info );
1662 if( !info ) return NoAction;
1663 XrdSysMutexHelper scopedLock( info->mutex );
1664 Log *log = DefaultEnv::GetLog();
1665
1666 //--------------------------------------------------------------------------
1667 // Update the substream queues
1668 //--------------------------------------------------------------------------
1669 info->strmSelector->MsgReceived( subStream );
1670
1671 //--------------------------------------------------------------------------
1672 // Check whether this message is a response to a request that has
1673 // timed out, and if so, drop it
1674 //--------------------------------------------------------------------------
1675 ServerResponse *rsp = (ServerResponse*)msg.GetBuffer();
1676 if( rsp->hdr.status == kXR_attn )
1677 {
1678 return NoAction;
1679 }
1680
1681 if( info->sidManager->IsTimedOut( rsp->hdr.streamid ) )
1682 {
1683 log->Error( XRootDTransportMsg, "Message %p, stream [%d, %d] is a "
1684 "response that we're no longer interested in (timed out)",
1685 (void*)&msg, rsp->hdr.streamid[0], rsp->hdr.streamid[1] );
1686 //------------------------------------------------------------------------
1687 // If it is kXR_waitresp there will be another one,
1688 // so we don't release the sid yet
1689 //------------------------------------------------------------------------
1690 if( rsp->hdr.status != kXR_waitresp )
1691 info->sidManager->ReleaseTimedOut( rsp->hdr.streamid );
1692 //------------------------------------------------------------------------
1693 // If it is a successful response to an open request
1694 // that timed out, we need to send a close
1695 //------------------------------------------------------------------------
1696 uint16_t sid; memcpy( &sid, rsp->hdr.streamid, 2 );
1697 std::set<uint16_t>::iterator sidIt = info->sentOpens.find( sid );
1698 if( sidIt != info->sentOpens.end() )
1699 {
1700 info->sentOpens.erase( sidIt );
1701 if( rsp->hdr.status == kXR_ok ) return RequestClose;
1702 }
1703 return DigestMsg;
1704 }
1705
1706 //--------------------------------------------------------------------------
1707 // If we have a wait or waitresp
1708 //--------------------------------------------------------------------------
1709 uint32_t seconds = 0;
1710 if( rsp->hdr.status == kXR_wait )
1711 seconds = ntohl( rsp->body.wait.seconds ) + 5; // we need extra time
1712 // to re-send the request
1713 else if( rsp->hdr.status == kXR_waitresp )
1714 {
1715 seconds = ntohl( rsp->body.waitresp.seconds );
1716
1717 log->Dump( XRootDMsg, "[%s] Got kXR_waitresp response of %u seconds, "
1718 "setting up wait barrier.",
1719 info->streamName.c_str(),
1720 seconds );
1721 }
1722
1723 time_t barrier = time(0) + seconds;
1724 if( info->waitBarrier < barrier )
1725 info->waitBarrier = barrier;
1726
1727 //--------------------------------------------------------------------------
1728 // If we got a response to an open request, we may need to bump the counter
1729 // of open files
1730 //--------------------------------------------------------------------------
1731 uint16_t sid; memcpy( &sid, rsp->hdr.streamid, 2 );
1732 std::set<uint16_t>::iterator sidIt = info->sentOpens.find( sid );
1733 if( sidIt != info->sentOpens.end() )
1734 {
1735 if( rsp->hdr.status == kXR_waitresp )
1736 return NoAction;
1737 info->sentOpens.erase( sidIt );
1738 if( rsp->hdr.status == kXR_ok )
1739 {
1740 ++info->openFiles;
1741 info->finstcnt.fetch_add( 1, std::memory_order_relaxed ); // another file File object instance has been bound with this connection
1742 }
1743 return NoAction;
1744 }
1745
1746 //--------------------------------------------------------------------------
1747 // If we got a response to a close, we may need to decrement the counter of
1748 // open files
1749 //--------------------------------------------------------------------------
1750 sidIt = info->sentCloses.find( sid );
1751 if( sidIt != info->sentCloses.end() )
1752 {
1753 if( rsp->hdr.status == kXR_waitresp )
1754 return NoAction;
1755 info->sentCloses.erase( sidIt );
1756 --info->openFiles;
1757 return NoAction;
1758 }
1759 return NoAction;
1760 }
kXR_char streamid[2]
Definition XProtocol.hh:956
@ kXR_waitresp
Definition XProtocol.hh:948
@ kXR_ok
Definition XProtocol.hh:941
@ kXR_attn
Definition XProtocol.hh:943
@ kXR_wait
Definition XProtocol.hh:947
@ RequestClose
Send a close request.
const uint64_t XRootDMsg

References ServerResponse::body, XrdCl::TransportHandler::DigestMsg, XrdCl::Log::Dump(), XrdCl::Log::Error(), XrdCl::XRootDChannelInfo::finstcnt, XrdCl::AnyObject::Get(), XrdCl::Buffer::GetBuffer(), XrdCl::DefaultEnv::GetLog(), ServerResponse::hdr, kXR_attn, kXR_ok, kXR_wait, kXR_waitresp, XrdCl::XRootDChannelInfo::mutex, XrdCl::TransportHandler::NoAction, XrdCl::XRootDChannelInfo::openFiles, XrdCl::TransportHandler::RequestClose, XrdCl::XRootDChannelInfo::sentCloses, XrdCl::XRootDChannelInfo::sentOpens, XrdCl::XRootDChannelInfo::sidManager, ServerResponseHeader::status, ServerResponseHeader::streamid, XrdCl::XRootDChannelInfo::streamName, XrdCl::XRootDChannelInfo::strmSelector, XrdCl::XRootDChannelInfo::waitBarrier, XrdCl::XRootDMsg, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ MessageSent()

void XrdCl::XRootDTransport::MessageSent ( Message * msg,
uint16_t subStream,
uint32_t bytesSent,
AnyObject & channelData )
virtual

Notify the transport about a message having been sent.

Implements XrdCl::TransportHandler.

Definition at line 1765 of file XrdClXRootDTransport.cc.

1769 {
1770 // Called when a message has been sent. For messages that return on a
1771 // different pathid (and hence may use a different poller) it is possible
1772 // that the server has already replied and the reply will trigger
1773 // MessageReceived() before this method has been called. However for open
1774 // and close this is never the case and this method is used for tracking
1775 // only those.
1776 XRootDChannelInfo *info = 0;
1777 channelData.Get( info );
1778 if( !info ) return;
1779 XrdSysMutexHelper scopedLock( info->mutex );
1780 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1781 uint16_t reqid = ntohs( req->header.requestid );
1782
1783
1784 //--------------------------------------------------------------------------
1785 // We need to track opens to know if we can close streams due to idleness
1786 //--------------------------------------------------------------------------
1787 uint16_t sid;
1788 memcpy( &sid, req->header.streamid, 2 );
1789
1790 if( reqid == kXR_open )
1791 info->sentOpens.insert( sid );
1792 else if( reqid == kXR_close )
1793 info->sentCloses.insert( sid );
1794 }
kXR_char streamid[2]
Definition XProtocol.hh:158

References XrdCl::AnyObject::Get(), XrdCl::Buffer::GetBuffer(), ClientRequest::header, kXR_close, kXR_open, XrdCl::XRootDChannelInfo::mutex, ClientRequestHdr::requestid, XrdCl::XRootDChannelInfo::sentCloses, XrdCl::XRootDChannelInfo::sentOpens, and ClientRequestHdr::streamid.

Here is the call graph for this function:

◆ Multiplex()

PathID XrdCl::XRootDTransport::Multiplex ( Message * msg,
AnyObject & channelData,
PathID * hint = 0 )
virtual

Return the ID for the up stream this message should be sent by and the down stream which the answer should be expected at. Modify the message itself if necessary. If hint is non-zero then the message should be modified such that the answer will be returned via the hinted stream.

Implements XrdCl::TransportHandler.

Definition at line 874 of file XrdClXRootDTransport.cc.

875 {
876 return PathID( 0, 0 );
877 }
PathID(uint16_t u=0, uint16_t d=0)

References XrdCl::PathID::PathID().

Here is the call graph for this function:

◆ MultiplexSubStream()

PathID XrdCl::XRootDTransport::MultiplexSubStream ( Message * msg,
AnyObject & channelData,
PathID * hint = 0 )
virtual

Return the ID for the up substream this message should be sent by and the down substream which the answer should be expected at. Modify the message itself if necessary. If hint is non-zero then the message should be modified such that the answer will be returned via the hinted stream.

Implements XrdCl::TransportHandler.

Definition at line 882 of file XrdClXRootDTransport.cc.

885 {
886 XRootDChannelInfo *info = 0;
887 channelData.Get( info );
888
889 if (!info) {
891 "Internal error: no channel info, cannot multiplex");
892 return PathID(0,0);
893 }
894
895 XrdSysMutexHelper scopedLock( info->mutex );
896
897 //--------------------------------------------------------------------------
898 // If we're not connected to a data server or we don't know that yet
899 // we stream through 0
900 //--------------------------------------------------------------------------
901 if( !(info->serverFlags & kXR_isServer) || info->stream.size() == 0 )
902 return PathID( 0, 0 );
903
904 //--------------------------------------------------------------------------
905 // Select the streams
906 //--------------------------------------------------------------------------
907 Log *log = DefaultEnv::GetLog();
908 uint16_t upStream = 0;
909 uint16_t downStream = 0;
910
911 if( hint )
912 {
913 upStream = hint->up;
914 downStream = hint->down;
915 }
916 else
917 {
918 upStream = 0;
919 std::vector<bool> connected;
920 connected.reserve( info->stream.size() - 1 );
921 size_t nbConnected = 0;
922 for( size_t i = 1; i < info->stream.size(); ++i )
923 if( info->stream[i].status == XRootDStreamInfo::Connected )
924 {
925 connected.push_back( true );
926 ++nbConnected;
927 }
928 else
929 connected.push_back( false );
930
931 if( nbConnected == 0 )
932 downStream = 0;
933 else
934 downStream = info->strmSelector->Select( connected );
935 }
936
937 if( upStream >= info->stream.size() )
938 {
939 log->Debug( XRootDTransportMsg,
940 "[%s] Up link stream %d does not exist, using 0",
941 info->streamName.c_str(), upStream );
942 upStream = 0;
943 }
944
945 if( downStream >= info->stream.size() )
946 {
947 log->Debug( XRootDTransportMsg,
948 "[%s] Down link stream %d does not exist, using 0",
949 info->streamName.c_str(), downStream );
950 downStream = 0;
951 }
952
953 //--------------------------------------------------------------------------
954 // Modify the message
955 //--------------------------------------------------------------------------
956 UnMarshallRequest( msg );
957 ClientRequestHdr *hdr = (ClientRequestHdr*)msg->GetBuffer();
958 switch( hdr->requestid )
959 {
960 //------------------------------------------------------------------------
961 // Read - we update the path id to tell the server where we want to
962 // get the response, but we still send the request through stream 0
963 // We need to allocate space for read_args if we don't have it
964 // included yet
965 //------------------------------------------------------------------------
966 case kXR_read:
967 {
968 if( msg->GetSize() < sizeof(ClientReadRequest) + 8 )
969 {
970 msg->ReAllocate( sizeof(ClientReadRequest) + 8 );
971 void *newBuf = msg->GetBuffer(sizeof(ClientReadRequest));
972 memset( newBuf, 0, 8 );
973 ClientReadRequest *req = (ClientReadRequest*)msg->GetBuffer();
974 req->dlen += 8;
975 }
976 read_args *args = (read_args*)msg->GetBuffer(sizeof(ClientReadRequest));
977 args->pathid = info->stream[downStream].pathId;
978 break;
979 }
980
981
982 //------------------------------------------------------------------------
983 // PgRead - we update the path id to tell the server where we want to
984 // get the response, but we still send the request through stream 0
985 // We need to allocate space for ClientPgReadReqArgs if we don't have it
986 // included yet
987 //------------------------------------------------------------------------
988 case kXR_pgread:
989 {
990 if( msg->GetSize() < sizeof( ClientPgReadRequest ) + sizeof( ClientPgReadReqArgs ) )
991 {
992 msg->ReAllocate( sizeof( ClientPgReadRequest ) + sizeof( ClientPgReadReqArgs ) );
993 void *newBuf = msg->GetBuffer( sizeof( ClientPgReadRequest ) );
994 memset( newBuf, 0, sizeof( ClientPgReadReqArgs ) );
995 ClientPgReadRequest *req = (ClientPgReadRequest*)msg->GetBuffer();
996 req->dlen += sizeof( ClientPgReadReqArgs );
997 }
998 ClientPgReadReqArgs *args = reinterpret_cast<ClientPgReadReqArgs*>(
999 msg->GetBuffer( sizeof( ClientPgReadRequest ) ) );
1000 args->pathid = info->stream[downStream].pathId;
1001 break;
1002 }
1003
1004 //------------------------------------------------------------------------
1005 // ReadV - the situation is identical to read but we don't need any
1006 // additional structures to specify the return path
1007 //------------------------------------------------------------------------
1008 case kXR_readv:
1009 {
1010 ClientReadVRequest *req = (ClientReadVRequest*)msg->GetBuffer();
1011 req->pathid = info->stream[downStream].pathId;
1012 break;
1013 }
1014
1015 //------------------------------------------------------------------------
1016 // Write - multiplexing writes doesn't work properly in the server
1017 //------------------------------------------------------------------------
1018 case kXR_write:
1019 {
1020// ClientWriteRequest *req = (ClientWriteRequest*)msg->GetBuffer();
1021// req->pathid = info->stream[downStream].pathId;
1022 break;
1023 }
1024
1025 //------------------------------------------------------------------------
1026 // WriteV - multiplexing writes doesn't work properly in the server
1027 //------------------------------------------------------------------------
1028 case kXR_writev:
1029 {
1030// ClientWriteVRequest *req = (ClientWriteVRequest*)msg->GetBuffer();
1031// req->pathid = info->stream[downStream].pathId;
1032 break;
1033 }
1034
1035 //------------------------------------------------------------------------
1036 // PgWrite - multiplexing writes doesn't work properly in the server
1037 //------------------------------------------------------------------------
1038 case kXR_pgwrite:
1039 {
1040// ClientWriteVRequest *req = (ClientWriteVRequest*)msg->GetBuffer();
1041// req->pathid = info->stream[downStream].pathId;
1042 break;
1043 }
1044 };
1045 MarshallRequest( msg );
1046 return PathID( upStream, downStream );
1047 }
kXR_char pathid
Definition XProtocol.hh:689
static XRootDStatus UnMarshallRequest(Message *msg)

References XrdCl::PathID::PathID(), XrdCl::XRootDStreamInfo::Connected, XrdCl::Log::Debug(), ClientPgReadRequest::dlen, ClientReadRequest::dlen, XrdCl::PathID::down, XrdCl::Log::Error(), XrdCl::AnyObject::Get(), XrdCl::Buffer::GetBuffer(), XrdCl::DefaultEnv::GetLog(), XrdCl::Buffer::GetSize(), kXR_isServer, kXR_pgread, kXR_pgwrite, kXR_read, kXR_readv, kXR_write, kXR_writev, MarshallRequest(), XrdCl::XRootDChannelInfo::mutex, ClientPgReadReqArgs::pathid, ClientReadVRequest::pathid, read_args::pathid, XrdCl::Buffer::ReAllocate(), ClientRequestHdr::requestid, XrdCl::XRootDChannelInfo::serverFlags, XrdCl::XRootDChannelInfo::stream, XrdCl::XRootDChannelInfo::streamName, XrdCl::XRootDChannelInfo::strmSelector, UnMarshallRequest(), XrdCl::PathID::up, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ NbConnectedStrm()

uint16_t XrdCl::XRootDTransport::NbConnectedStrm ( AnyObject & channelData)
static

Number of currently connected data streams.

Definition at line 1548 of file XrdClXRootDTransport.cc.

1549 {
1550 XRootDChannelInfo *info = 0;
1551 channelData.Get( info );
1552
1553 if (!info) {
1554 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1555 return 0;
1556 }
1557
1558 XrdSysMutexHelper scopedLock( info->mutex );
1559
1560 uint16_t nbConnected = 0;
1561 for( size_t i = 1; i < info->stream.size(); ++i )
1562 if( info->stream[i].status == XRootDStreamInfo::Connected )
1563 ++nbConnected;
1564
1565 return nbConnected;
1566 }

References XrdCl::XRootDStreamInfo::Connected, XrdCl::Log::Error(), XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetLog(), XrdCl::XRootDChannelInfo::mutex, XrdCl::XRootDChannelInfo::stream, and XrdCl::XRootDTransportMsg.

Referenced by XrdCl::Channel::NbConnectedStrm().

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

◆ NeedControlConnection()

virtual bool XrdCl::XRootDTransport::NeedControlConnection ( )
inlinevirtual

Return the information whether a control connection needs to be valid before establishing other connections

Definition at line 167 of file XrdClXRootDTransport.hh.

168 {
169 return true;
170 }

◆ NeedEncryption()

bool XrdCl::XRootDTransport::NeedEncryption ( HandShakeData * handShakeData,
AnyObject & channelData )
virtual
Returns
: true if encryption should be turned on, false otherwise

Implements XrdCl::TransportHandler.

Definition at line 1860 of file XrdClXRootDTransport.cc.

1862 {
1863 XRootDChannelInfo *info = 0;
1864 channelData.Get( info );
1865
1866 XrdCl::Env *env = XrdCl::DefaultEnv::GetEnv();
1867 int notlsok = DefaultNoTlsOK;
1868 env->GetInt( "NoTlsOK", notlsok );
1869
1870
1871 if( notlsok )
1872 return info->encrypted;
1873
1874 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
1875
1876 // Did the server instructed us to switch to TLS right away?
1877 if( sInfo.serverFlags & kXR_gotoTLS )
1878 {
1879 if( handShakeData->subStreamId == 0 ) info->encrypted = true;
1880 return true ;
1881 }
1882
1883 //--------------------------------------------------------------------------
1884 // The control stream (sub-stream 0) might need to switch to TLS before
1885 // login or after login
1886 //--------------------------------------------------------------------------
1887 if( handShakeData->subStreamId == 0 )
1888 {
1889 //------------------------------------------------------------------------
1890 // We are about to login and the server asked to start encrypting
1891 // before login
1892 //------------------------------------------------------------------------
1893 if( ( sInfo.status == XRootDStreamInfo::LoginSent ) &&
1894 ( info->serverFlags & kXR_tlsLogin ) )
1895 {
1896 info->encrypted = true;
1897 return true;
1898 }
1899
1900 //--------------------------------------------------------------------
1901 // The hand-shake is done and the server requested to encrypt the session
1902 //--------------------------------------------------------------------
1903 if( (sInfo.status == XRootDStreamInfo::Connected ||
1904 //--------------------------------------------------------------------
1905 // we really need to turn on TLS before we sent kXR_endsess and we
1906 // are about to do so (1st enable encryption, then send kXR_endsess)
1907 //--------------------------------------------------------------------
1908 sInfo.status == XRootDStreamInfo::EndSessionSent ) &&
1909 ( info->serverFlags & kXR_tlsSess ) )
1910 {
1911 info->encrypted = true;
1912 return true;
1913 }
1914 }
1915 //--------------------------------------------------------------------------
1916 // A data stream (sub-stream > 0) if need be will be switched to TLS before
1917 // bind.
1918 //--------------------------------------------------------------------------
1919 else
1920 {
1921 //------------------------------------------------------------------------
1922 // We are about to bind a data stream and the server asked to start
1923 // encrypting before bind
1924 //------------------------------------------------------------------------
1925 if( ( sInfo.status == XRootDStreamInfo::BindSent ) &&
1926 ( info->serverFlags & kXR_tlsData ) )
1927 {
1928 return true;
1929 }
1930 }
1931
1932 return false;
1933 }
#define kXR_tlsLogin
#define kXR_gotoTLS
#define kXR_tlsSess
#define kXR_tlsData
bool GetInt(const std::string &key, int &value)
Definition XrdClEnv.cc:115
const int DefaultNoTlsOK

References XrdCl::XRootDStreamInfo::BindSent, XrdCl::XRootDStreamInfo::Connected, XrdCl::DefaultNoTlsOK, XrdCl::XRootDChannelInfo::encrypted, XrdCl::XRootDStreamInfo::EndSessionSent, XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetEnv(), XrdCl::Env::GetInt(), kXR_gotoTLS, kXR_tlsData, kXR_tlsLogin, kXR_tlsSess, XrdCl::XRootDStreamInfo::LoginSent, XrdCl::XRootDChannelInfo::serverFlags, XrdCl::XRootDStreamInfo::serverFlags, XrdCl::XRootDStreamInfo::status, XrdCl::XRootDChannelInfo::stream, and XrdCl::HandShakeData::subStreamId.

Here is the call graph for this function:

◆ Query()

Status XrdCl::XRootDTransport::Query ( uint16_t query,
AnyObject & result,
AnyObject & channelData )
virtual

Query the channel.

Implements XrdCl::TransportHandler.

Definition at line 1604 of file XrdClXRootDTransport.cc.

1607 {
1608 XRootDChannelInfo *info = 0;
1609 channelData.Get( info );
1610
1611 if (!info)
1613
1614 XrdSysMutexHelper scopedLock( info->mutex );
1615
1616 switch( query )
1617 {
1618 //------------------------------------------------------------------------
1619 // Protocol name
1620 //------------------------------------------------------------------------
1622 result.Set( (const char*)"XRootD", false );
1623 return Status();
1624
1625 //------------------------------------------------------------------------
1626 // Authentication
1627 //------------------------------------------------------------------------
1629 result.Set( new std::string( info->authProtocolName ), false );
1630 return Status();
1631
1632 //------------------------------------------------------------------------
1633 // Server flags
1634 //------------------------------------------------------------------------
1636 result.Set( new int( info->serverFlags ), false );
1637 return Status();
1638
1639 //------------------------------------------------------------------------
1640 // Protocol version
1641 //------------------------------------------------------------------------
1643 result.Set( new int( info->protocolVersion ), false );
1644 return Status();
1645
1647 result.Set( new bool( info->encrypted ), false );
1648 return Status();
1649 };
1651 }
const uint16_t errQueryNotSupported
static const uint16_t Name
Transport name, returns const char *.
static const uint16_t Auth
Transport name, returns std::string *.
static const uint16_t ServerFlags
returns server flags
static const uint16_t ProtocolVersion
returns the protocol version
static const uint16_t IsEncrypted
returns true if the channel is encrypted

References XrdCl::Status::Status(), XrdCl::XRootDStatus::XRootDStatus(), XrdCl::TransportQuery::Auth, XrdCl::XRootDChannelInfo::authProtocolName, XrdCl::XRootDChannelInfo::encrypted, XrdCl::errInternal, XrdCl::errQueryNotSupported, XrdCl::AnyObject::Get(), XrdCl::XRootDQuery::IsEncrypted, XrdCl::XRootDChannelInfo::mutex, XrdCl::TransportQuery::Name, XrdCl::XRootDQuery::ProtocolVersion, XrdCl::XRootDChannelInfo::protocolVersion, XrdCl::XRootDQuery::ServerFlags, XrdCl::XRootDChannelInfo::serverFlags, XrdCl::AnyObject::Set(), XrdCl::stError, and XrdCl::stFatal.

Here is the call graph for this function:

◆ SetDescription()

void XrdCl::XRootDTransport::SetDescription ( Message * msg)
inlinestatic

Get the description of a message.

Definition at line 245 of file XrdClXRootDTransport.hh.

246 {
247 std::ostringstream o;
248 GenerateDescription( msg->GetBuffer(), o );
249 msg->SetDescription( o.str() );
250 }

References GenerateDescription(), XrdCl::Buffer::GetBuffer(), and XrdCl::Message::SetDescription().

Referenced by XrdCl::FileStateHandler::Checkpoint(), XrdCl::FileStateHandler::ChkptWrt(), XrdCl::FileStateHandler::ChkptWrtV(), XrdCl::FileSystem::ChMod(), XrdCl::FileStateHandler::Clone(), XrdCl::FileStateHandler::Close(), XrdCl::FileSystem::DirList(), XrdCl::FileStateHandler::Fcntl(), XrdCl::FileSystem::Locate(), XrdCl::FileSystem::MkDir(), XrdCl::FileSystem::Mv(), XrdCl::FileStateHandler::PgReadImpl(), XrdCl::FileStateHandler::PgWriteImpl(), XrdCl::FileSystem::Ping(), XrdCl::FileSystem::Prepare(), XrdCl::FileStateHandler::PreRead(), XrdCl::FileSystem::Protocol(), XrdCl::FileSystem::Query(), XrdCl::FileStateHandler::Read(), XrdCl::FileStateHandler::ReadV(), XrdCl::MessageUtils::RewriteCGIAndPath(), XrdCl::FileSystem::Rm(), XrdCl::FileSystem::RmDir(), XrdCl::FileStateHandler::Stat(), XrdCl::FileSystem::Stat(), XrdCl::FileSystem::StatVFS(), XrdCl::FileStateHandler::Sync(), XrdCl::FileStateHandler::Truncate(), XrdCl::FileSystem::Truncate(), XrdCl::FileStateHandler::VectorRead(), XrdCl::FileStateHandler::VectorWrite(), XrdCl::FileStateHandler::Visa(), XrdCl::FileStateHandler::Write(), and XrdCl::FileStateHandler::WriteV().

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

◆ SubStreamNumber()

uint16_t XrdCl::XRootDTransport::SubStreamNumber ( AnyObject & channelData)
virtual

Return a number of substreams per stream that should be created.

Implements XrdCl::TransportHandler.

Definition at line 1054 of file XrdClXRootDTransport.cc.

1055 {
1056 XRootDChannelInfo *info = 0;
1057 channelData.Get( info );
1058
1059 if (!info) {
1060 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1061 return 1;
1062 }
1063
1064 XrdSysMutexHelper scopedLock( info->mutex );
1065
1066 //--------------------------------------------------------------------------
1067 // If the connection has been opened in order to orchestrate a TPC or
1068 // the remote server is a Manager or Metamanager we will need only one
1069 // (control) stream.
1070 //--------------------------------------------------------------------------
1071 if( info->istpc || !(info->serverFlags & kXR_isServer ) ) return 1;
1072
1073 //--------------------------------------------------------------------------
1074 // Number of streams requested by user
1075 //--------------------------------------------------------------------------
1076 uint16_t ret = info->stream.size();
1077
1078 XrdCl::Env *env = XrdCl::DefaultEnv::GetEnv();
1079 int nodata = DefaultTlsNoData;
1080 env->GetInt( "TlsNoData", nodata );
1081
1082 // Does the server require the stream 0 to be encrypted?
1083 bool srvTlsStrm0 = ( info->serverFlags & kXR_gotoTLS ) ||
1084 ( info->serverFlags & kXR_tlsLogin ) ||
1085 ( info->serverFlags & kXR_tlsSess );
1086 // Does the server NOT require the data streams to be encrypted?
1087 bool srvNoTlsData = !( info->serverFlags & kXR_tlsData );
1088 // Does the user require the stream 0 to be encrypted?
1089 bool usrTlsStrm0 = info->encrypted;
1090 // Does the user NOT require the data streams to be encrypted?
1091 bool usrNoTlsData = !info->encrypted || ( info->encrypted && nodata );
1092
1093 if( ( usrTlsStrm0 && usrNoTlsData && srvNoTlsData ) ||
1094 ( srvTlsStrm0 && srvNoTlsData && usrNoTlsData ) )
1095 {
1096 //------------------------------------------------------------------------
1097 // The server or user asked us to encrypt stream 0, but to send the data
1098 // (read/write) using a plain TCP connection
1099 //------------------------------------------------------------------------
1100 if( ret == 1 ) ++ret;
1101 }
1102
1103 if( ret > info->stream.size() )
1104 {
1105 info->stream.resize( ret );
1106 info->strmSelector->AdjustQueues( ret );
1107 }
1108
1109 return ret;
1110 }
const int DefaultTlsNoData

References XrdCl::DefaultTlsNoData, XrdCl::XRootDChannelInfo::encrypted, XrdCl::Log::Error(), XrdCl::AnyObject::Get(), XrdCl::DefaultEnv::GetEnv(), XrdCl::Env::GetInt(), XrdCl::DefaultEnv::GetLog(), XrdCl::XRootDChannelInfo::istpc, kXR_gotoTLS, kXR_isServer, kXR_tlsData, kXR_tlsLogin, kXR_tlsSess, XrdCl::XRootDChannelInfo::mutex, XrdCl::XRootDChannelInfo::serverFlags, XrdCl::XRootDChannelInfo::stream, XrdCl::XRootDChannelInfo::strmSelector, and XrdCl::XRootDTransportMsg.

Here is the call graph for this function:

◆ UnMarchalStatusMore()

XRootDStatus XrdCl::XRootDTransport::UnMarchalStatusMore ( Message & msg)
static

Unmarshall the correction-segment of the status response for pgwrite.

Definition at line 1461 of file XrdClXRootDTransport.cc.

1462 {
1463 ServerResponseV2 *rsp = (ServerResponseV2*)msg.GetBuffer();
1464 uint16_t reqType = rsp->status.bdy.requestid + kXR_1stRequest;
1465
1466 switch( reqType )
1467 {
1468 case kXR_pgwrite:
1469 {
1470 //--------------------------------------------------------------------------
1471 // If there's no additional data there's nothing to unmarshal
1472 //--------------------------------------------------------------------------
1473 if( rsp->status.bdy.dlen == 0 ) return XRootDStatus();
1474 //--------------------------------------------------------------------------
1475 // If there's not enough data to form correction-segment report an error
1476 //--------------------------------------------------------------------------
1477 if( size_t( rsp->status.bdy.dlen ) < sizeof( ServerResponseBody_pgWrCSE ) )
1479 "kXR_status: invalid message size." );
1480
1481 //--------------------------------------------------------------------------
1482 // Calculate the crc32c for the additional data
1483 //--------------------------------------------------------------------------
1484 ServerResponseBody_pgWrCSE *cse = (ServerResponseBody_pgWrCSE*)msg.GetBuffer( sizeof( ServerResponseV2 ) );
1485 cse->cseCRC = ntohl( cse->cseCRC );
1486 size_t length = rsp->status.bdy.dlen - sizeof( uint32_t );
1487 void* buffer = msg.GetBuffer( sizeof( ServerResponseV2 ) + sizeof( uint32_t ) );
1488 uint32_t crcval = XrdOucCRC::Calc32C( buffer, length );
1489
1490 //--------------------------------------------------------------------------
1491 // Do the integrity checks
1492 //--------------------------------------------------------------------------
1493 if( crcval != cse->cseCRC )
1494 {
1495 return XRootDStatus( stError, errDataError, 0, "kXR_status response header "
1496 "corrupted (crc32c integrity check failed)." );
1497 }
1498
1499 cse->dlFirst = ntohs( cse->dlFirst );
1500 cse->dlLast = ntohs( cse->dlLast );
1501
1502 size_t pgcnt = ( rsp->status.bdy.dlen - sizeof( ServerResponseBody_pgWrCSE ) ) /
1503 sizeof( kXR_int64 );
1504 kXR_int64 *pgoffs = (kXR_int64*)msg.GetBuffer( sizeof( ServerResponseV2 ) +
1505 sizeof( ServerResponseBody_pgWrCSE ) );
1506
1507 for( size_t i = 0; i < pgcnt; ++i )
1508 pgoffs[i] = ntohll( pgoffs[i] );
1509
1510 return XRootDStatus();
1511 break;
1512 }
1513
1514 default:
1515 break;
1516 }
1517
1519 }
ServerResponseStatus status
@ kXR_1stRequest
Definition XProtocol.hh:112
long long kXR_int64
Definition XPtypes.hh:98
static uint32_t Calc32C(const void *data, size_t count, uint32_t prevcs=0)
Definition XrdOucCRC.cc:190
const uint16_t errNotSupported

References XrdCl::XRootDStatus::XRootDStatus(), ServerResponseStatus::bdy, XrdOucCRC::Calc32C(), ServerResponseBody_pgWrCSE::cseCRC, ServerResponseBody_Status::dlen, ServerResponseBody_pgWrCSE::dlFirst, ServerResponseBody_pgWrCSE::dlLast, XrdCl::errDataError, XrdCl::errInvalidMessage, XrdCl::errNotSupported, XrdCl::Buffer::GetBuffer(), kXR_1stRequest, kXR_pgwrite, ServerResponseBody_Status::requestid, ServerResponseV2::status, and XrdCl::stError.

Referenced by GetMore().

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

◆ UnMarshallBody()

XRootDStatus XrdCl::XRootDTransport::UnMarshallBody ( Message * msg,
uint16_t reqType )
static

Unmarshall the body of the incoming message.

Definition at line 1307 of file XrdClXRootDTransport.cc.

1308 {
1309 ServerResponse *m = (ServerResponse *)msg->GetBuffer();
1310
1311 //--------------------------------------------------------------------------
1312 // kXR_ok
1313 //--------------------------------------------------------------------------
1314 if( m->hdr.status == kXR_ok )
1315 {
1316 switch( reqType )
1317 {
1318 //----------------------------------------------------------------------
1319 // kXR_protocol
1320 //----------------------------------------------------------------------
1321 case kXR_protocol:
1322 if( m->hdr.dlen < 8 )
1323 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_protocol: body too short." );
1324 m->body.protocol.pval = ntohl( m->body.protocol.pval );
1325 m->body.protocol.flags = ntohl( m->body.protocol.flags );
1326 break;
1327 }
1328 }
1329 //--------------------------------------------------------------------------
1330 // kXR_error
1331 //--------------------------------------------------------------------------
1332 else if( m->hdr.status == kXR_error )
1333 {
1334 if( m->hdr.dlen < 4 )
1335 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_error: body too short." );
1336 m->body.error.errnum = ntohl( m->body.error.errnum );
1337 }
1338
1339 //--------------------------------------------------------------------------
1340 // kXR_wait
1341 //--------------------------------------------------------------------------
1342 else if( m->hdr.status == kXR_wait )
1343 {
1344 if( m->hdr.dlen < 4 )
1345 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_wait: body too short." );
1346 m->body.wait.seconds = htonl( m->body.wait.seconds );
1347 }
1348
1349 //--------------------------------------------------------------------------
1350 // kXR_redirect
1351 //--------------------------------------------------------------------------
1352 else if( m->hdr.status == kXR_redirect )
1353 {
1354 if( m->hdr.dlen < 4 )
1355 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_redirect: body too short." );
1356 m->body.redirect.port = htonl( m->body.redirect.port );
1357 }
1358
1359 //--------------------------------------------------------------------------
1360 // kXR_waitresp
1361 //--------------------------------------------------------------------------
1362 else if( m->hdr.status == kXR_waitresp )
1363 {
1364 if( m->hdr.dlen < 4 )
1365 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_waitresp: body too short." );
1366 m->body.waitresp.seconds = htonl( m->body.waitresp.seconds );
1367 }
1368
1369 //--------------------------------------------------------------------------
1370 // kXR_attn
1371 //--------------------------------------------------------------------------
1372 else if( m->hdr.status == kXR_attn )
1373 {
1374 if( m->hdr.dlen < 4 )
1375 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_attn: body too short." );
1376 m->body.attn.actnum = htonl( m->body.attn.actnum );
1377 }
1378
1379 return XRootDStatus();
1380 }
@ kXR_redirect
Definition XProtocol.hh:946
@ kXR_error
Definition XProtocol.hh:945

References XrdCl::XRootDStatus::XRootDStatus(), ServerResponse::body, ServerResponseHeader::dlen, XrdCl::errInvalidMessage, XrdCl::Buffer::GetBuffer(), ServerResponse::hdr, kXR_attn, kXR_error, kXR_ok, kXR_protocol, kXR_redirect, kXR_wait, kXR_waitresp, ServerResponseHeader::status, and XrdCl::stError.

Referenced by XrdCl::XRootDMsgHandler::Process().

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

◆ UnMarshallHeader()

void XrdCl::XRootDTransport::UnMarshallHeader ( Message & msg)
static

Unmarshall the header incoming message.

Definition at line 1524 of file XrdClXRootDTransport.cc.

1525 {
1526 ServerResponseHeader *header = (ServerResponseHeader *)msg.GetBuffer();
1527 header->status = ntohs( header->status );
1528 header->dlen = ntohl( header->dlen );
1529 }

References ServerResponseHeader::dlen, XrdCl::Buffer::GetBuffer(), and ServerResponseHeader::status.

Referenced by GetHeader().

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

◆ UnMarshallRequest()

XRootDStatus XrdCl::XRootDTransport::UnMarshallRequest ( Message * msg)
static

Unmarshall the request - sometimes the requests need to be rewritten, so we need to unmarshall them

Definition at line 1286 of file XrdClXRootDTransport.cc.

1287 {
1288 if( !msg->IsMarshalled() ) return XRootDStatus( stOK, suAlreadyDone );
1289 // We rely on the marshaling process to be symmetric!
1290 // First we unmarshall the request ID and the length because
1291 // MarshallRequest() relies on these, and then we need to unmarshall these
1292 // two again, because they get marshalled in MarshallRequest().
1293 // All this is pretty damn ugly and should be rewritten.
1294 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1295 req->header.requestid = htons( req->header.requestid );
1296 req->header.dlen = htonl( req->header.dlen );
1297 XRootDStatus st = MarshallRequest( msg );
1298 req->header.requestid = htons( req->header.requestid );
1299 req->header.dlen = htonl( req->header.dlen );
1300 msg->SetIsMarshalled( false );
1301 return st;
1302 }
const uint16_t suAlreadyDone

References XrdCl::XRootDStatus::XRootDStatus(), ClientRequestHdr::dlen, XrdCl::Buffer::GetBuffer(), ClientRequest::header, XrdCl::Message::IsMarshalled(), MarshallRequest(), ClientRequestHdr::requestid, XrdCl::Message::SetIsMarshalled(), XrdCl::stOK, and XrdCl::suAlreadyDone.

Referenced by MultiplexSubStream(), XrdCl::MessageUtils::RedirectMessage(), and XrdCl::MessageUtils::SendMessage().

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

◆ UnMarshalStatusBody()

XRootDStatus XrdCl::XRootDTransport::UnMarshalStatusBody ( Message & msg,
uint16_t reqType )
static

Unmarshall the body of the status response.

Definition at line 1385 of file XrdClXRootDTransport.cc.

1386 {
1387 //--------------------------------------------------------------------------
1388 // Calculate the crc32c before the unmarshaling the body!
1389 //--------------------------------------------------------------------------
1390 ServerResponseStatus *rspst = (ServerResponseStatus*)msg.GetBuffer();
1391 char *buffer = msg.GetBuffer( 8 + sizeof( rspst->bdy.crc32c ) );
1392 size_t length = rspst->hdr.dlen - sizeof( rspst->bdy.crc32c );
1393 uint32_t crcval = XrdOucCRC::Calc32C( buffer, length );
1394
1395 size_t stlen = sizeof( ServerResponseStatus );
1396 switch( reqType )
1397 {
1398 case kXR_pgread:
1399 {
1400 stlen += sizeof( ServerResponseBody_pgRead );
1401 break;
1402 }
1403
1404 case kXR_pgwrite:
1405 {
1406 stlen += sizeof( ServerResponseBody_pgWrite );
1407 break;
1408 }
1409 }
1410
1411 if( msg.GetSize() < stlen ) return XRootDStatus( stError, errInvalidMessage, 0,
1412 "kXR_status: invalid message size." );
1413
1414 rspst->bdy.crc32c = ntohl( rspst->bdy.crc32c );
1415 rspst->bdy.dlen = ntohl( rspst->bdy.dlen );
1416
1417 switch( reqType )
1418 {
1419 case kXR_pgread:
1420 {
1421 ServerResponseBody_pgRead *pgrdbdy = (ServerResponseBody_pgRead*)msg.GetBuffer( sizeof( ServerResponseStatus ) );
1422 pgrdbdy->offset = ntohll( pgrdbdy->offset );
1423 break;
1424 }
1425
1426 case kXR_pgwrite:
1427 {
1428 ServerResponseBody_pgWrite *pgwrtbdy = (ServerResponseBody_pgWrite*)msg.GetBuffer( sizeof( ServerResponseStatus ) );
1429 pgwrtbdy->offset = ntohll( pgwrtbdy->offset );
1430 break;
1431 }
1432 }
1433
1434 //--------------------------------------------------------------------------
1435 // Do the integrity checks
1436 //--------------------------------------------------------------------------
1437 if( crcval != rspst->bdy.crc32c )
1438 {
1439 return XRootDStatus( stError, errDataError, 0, "kXR_status response header "
1440 "corrupted (crc32c integrity check failed)." );
1441 }
1442
1443 if( rspst->hdr.streamid[0] != rspst->bdy.streamID[0] ||
1444 rspst->hdr.streamid[1] != rspst->bdy.streamID[1] )
1445 {
1446 return XRootDStatus( stError, errDataError, 0, "response header corrupted "
1447 "(stream ID mismatch)." );
1448 }
1449
1450
1451
1452 if( rspst->bdy.requestid + kXR_1stRequest != reqType )
1453 {
1454 return XRootDStatus( stError, errDataError, 0, "kXR_status response header corrupted "
1455 "(request ID mismatch)." );
1456 }
1457
1458 return XRootDStatus();
1459 }
struct ServerResponseHeader hdr

References XrdCl::XRootDStatus::XRootDStatus(), ServerResponseStatus::bdy, XrdOucCRC::Calc32C(), ServerResponseBody_Status::crc32c, ServerResponseBody_Status::dlen, ServerResponseHeader::dlen, XrdCl::errDataError, XrdCl::errInvalidMessage, XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::GetSize(), ServerResponseStatus::hdr, kXR_1stRequest, kXR_pgread, kXR_pgwrite, ServerResponseBody_pgRead::offset, ServerResponseBody_pgWrite::offset, ServerResponseBody_Status::requestid, XrdCl::stError, ServerResponseBody_Status::streamID, and ServerResponseHeader::streamid.

Referenced by XrdCl::XRootDMsgHandler::InspectStatusRsp().

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

◆ WaitBeforeExit()

void XrdCl::XRootDTransport::WaitBeforeExit ( )
virtual

Wait until the program can safely exit.

Implements XrdCl::TransportHandler.

Definition at line 1851 of file XrdClXRootDTransport.cc.

1852 {
1853 XrdSysRWLockHelper scope( pSecUnloadHandler->lock, false ); // obtain write lock
1854 pSecUnloadHandler->unloaded = true;
1855 }

◆ PluginUnloadHandler

friend struct PluginUnloadHandler
friend

Definition at line 432 of file XrdClXRootDTransport.hh.

References PluginUnloadHandler.

Referenced by XRootDTransport(), and PluginUnloadHandler.


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