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 290 of file XrdClXRootDTransport.cc.

290 :
291 pSecUnloadHandler( new PluginUnloadHandler() )
292 {
293 }

References PluginUnloadHandler.

+ Here is the call graph for this function:

◆ ~XRootDTransport()

XrdCl::XRootDTransport::~XRootDTransport ( )

Destructor.

Definition at line 298 of file XrdClXRootDTransport.cc.

299 {
300 delete pSecUnloadHandler; pSecUnloadHandler = 0;
301 }

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 1824 of file XrdClXRootDTransport.cc.

1825 {
1826 XRootDChannelInfo *info = 0;
1827 channelData.Get( info );
1828 if( info->finstcnt.load( std::memory_order_relaxed ) > 0 )
1829 info->finstcnt.fetch_sub( 1, std::memory_order_relaxed );
1830 }

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 1554 of file XrdClXRootDTransport.cc.

1556 {
1557 XRootDChannelInfo *info = 0;
1558 channelData.Get( info );
1559
1560 if (!info) {
1561 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1562 return;
1563 }
1564
1565 XrdSysMutexHelper scopedLock( info->mutex );
1566
1567 CleanUpProtection( info );
1568
1569 if( !info->stream.empty() )
1570 {
1571 XRootDStreamInfo &sInfo = info->stream[subStreamId];
1572 sInfo.status = XRootDStreamInfo::Disconnected;
1573 }
1574
1575 if( subStreamId == 0 )
1576 {
1577 info->sidManager->ReleaseAllTimedOut();
1578 info->sentOpens.clear();
1579 info->sentCloses.clear();
1580 info->openFiles = 0;
1581 info->waitBarrier = 0;
1582 }
1583 }
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 470 of file XrdClXRootDTransport.cc.

471 {
472 }

◆ GenerateDescription()

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

Get the description of a message.

Definition at line 3003 of file XrdClXRootDTransport.cc.

3004 {
3005 Log *log = DefaultEnv::GetLog();
3006 if( log->GetLevel() < Log::ErrorMsg )
3007 return;
3008
3009 ClientRequestHdr *req = (ClientRequestHdr *)msg;
3010 switch( req->requestid )
3011 {
3012 //------------------------------------------------------------------------
3013 // kXR_open
3014 //------------------------------------------------------------------------
3015 case kXR_open:
3016 {
3017 ClientOpenRequest *sreq = (ClientOpenRequest *)msg;
3018 o << "kXR_open (";
3019 char *fn = GetDataAsString( msg );
3020 o << "file: " << fn << ", ";
3021 delete [] fn;
3022 o << "mode: 0" << std::setbase(8) << sreq->mode << ", ";
3023 o << std::setbase(10);
3024 o << "flags: ";
3025 if( sreq->options == 0 )
3026 o << "none";
3027 else
3028 {
3029 if( sreq->options & kXR_compress )
3030 o << "kXR_compress ";
3031 if( sreq->options & kXR_delete )
3032 o << "kXR_delete ";
3033 if( sreq->options & kXR_force )
3034 o << "kXR_force ";
3035 if( sreq->options & kXR_mkpath )
3036 o << "kXR_mkpath ";
3037 if( sreq->options & kXR_new )
3038 o << "kXR_new ";
3039 if( sreq->options & kXR_nowait )
3040 o << "kXR_nowait ";
3041 if( sreq->options & kXR_open_apnd )
3042 o << "kXR_open_apnd ";
3043 if( sreq->options & kXR_open_read )
3044 o << "kXR_open_read ";
3045 if( sreq->options & kXR_open_updt )
3046 o << "kXR_open_updt ";
3047 if( sreq->options & kXR_open_wrto )
3048 o << "kXR_open_wrto ";
3049 if( sreq->options & kXR_posc )
3050 o << "kXR_posc ";
3051 if( sreq->options & kXR_prefname )
3052 o << "kXR_prefname ";
3053 if( sreq->options & kXR_refresh )
3054 o << "kXR_refresh ";
3055 if( sreq->options & kXR_4dirlist )
3056 o << "kXR_4dirlist ";
3057 if( sreq->options & kXR_replica )
3058 o << "kXR_replica ";
3059 if( sreq->options & kXR_seqio )
3060 o << "kXR_seqio ";
3061 if( sreq->options & kXR_async )
3062 o << "kXR_async ";
3063 if( sreq->options & kXR_retstat )
3064 o << "kXR_retstat ";
3065 }
3066 o << ")";
3067 break;
3068 }
3069
3070 //------------------------------------------------------------------------
3071 // kXR_close
3072 //------------------------------------------------------------------------
3073 case kXR_close:
3074 {
3075 ClientCloseRequest *sreq = (ClientCloseRequest *)msg;
3076 o << "kXR_close (";
3077 o << "handle: " << FileHandleToStr( sreq->fhandle );
3078 o << ")";
3079 break;
3080 }
3081
3082 //------------------------------------------------------------------------
3083 // kXR_stat
3084 //------------------------------------------------------------------------
3085 case kXR_stat:
3086 {
3087 ClientStatRequest *sreq = (ClientStatRequest *)msg;
3088 o << "kXR_stat (";
3089 if( sreq->dlen )
3090 {
3091 char *fn = GetDataAsString( msg );;
3092 o << "path: " << fn << ", ";
3093 delete [] fn;
3094 }
3095 else
3096 {
3097 o << "handle: " << FileHandleToStr( sreq->fhandle );
3098 o << ", ";
3099 }
3100 o << "flags: ";
3101 if( sreq->options == 0 )
3102 o << "none";
3103 else
3104 {
3105 if( sreq->options & kXR_vfs )
3106 o << "kXR_vfs";
3107 }
3108 o << ")";
3109 break;
3110 }
3111
3112 //------------------------------------------------------------------------
3113 // kXR_read
3114 //------------------------------------------------------------------------
3115 case kXR_read:
3116 {
3117 ClientReadRequest *sreq = (ClientReadRequest *)msg;
3118 o << "kXR_read (";
3119 o << "handle: " << FileHandleToStr( sreq->fhandle );
3120 o << std::setbase(10);
3121 o << ", ";
3122 o << "offset: " << sreq->offset << ", ";
3123 o << "size: " << sreq->rlen << ")";
3124 break;
3125 }
3126
3127 //------------------------------------------------------------------------
3128 // kXR_pgread
3129 //------------------------------------------------------------------------
3130 case kXR_pgread:
3131 {
3132 ClientPgReadRequest *sreq = (ClientPgReadRequest *)msg;
3133 o << "kXR_pgread (";
3134 o << "handle: " << FileHandleToStr( sreq->fhandle );
3135 o << std::setbase(10);
3136 o << ", ";
3137 o << "offset: " << sreq->offset << ", ";
3138 o << "size: " << sreq->rlen << ")";
3139 break;
3140 }
3141
3142 //------------------------------------------------------------------------
3143 // kXR_write
3144 //------------------------------------------------------------------------
3145 case kXR_write:
3146 {
3147 ClientWriteRequest *sreq = (ClientWriteRequest *)msg;
3148 o << "kXR_write (";
3149 o << "handle: " << FileHandleToStr( sreq->fhandle );
3150 o << std::setbase(10);
3151 o << ", ";
3152 o << "offset: " << sreq->offset << ", ";
3153 o << "size: " << sreq->dlen << ")";
3154 break;
3155 }
3156
3157 //------------------------------------------------------------------------
3158 // kXR_pgwrite
3159 //------------------------------------------------------------------------
3160 case kXR_pgwrite:
3161 {
3162 ClientPgWriteRequest *sreq = (ClientPgWriteRequest *)msg;
3163 o << "kXR_pgwrite (";
3164 o << "handle: " << FileHandleToStr( sreq->fhandle );
3165 o << std::setbase(10);
3166 o << ", ";
3167 o << "offset: " << sreq->offset << ", ";
3168 o << "size: " << sreq->dlen << ")";
3169 break;
3170 }
3171
3172 //------------------------------------------------------------------------
3173 // kXR_fattr
3174 //------------------------------------------------------------------------
3175 case kXR_fattr:
3176 {
3177 ClientFattrRequest *sreq = (ClientFattrRequest *)msg;
3178 int nattr = sreq->numattr;
3179 int options = sreq->options;
3180 o << "kXR_fattr";
3181 switch (sreq->subcode) {
3182 case kXR_fattrGet:
3183 o << "Get";
3184 break;
3185 case kXR_fattrSet:
3186 o << "Set";
3187 break;
3188 case kXR_fattrList:
3189 o << "List";
3190 break;
3191 case kXR_fattrDel:
3192 o << "Delete";
3193 break;
3194 default:
3195 o << " unknown subcode: " << sreq->subcode;
3196 break;
3197 }
3198 o << " (handle: " << FileHandleToStr( sreq->fhandle );
3199 o << std::setbase(10);
3200 if (nattr)
3201 o << ", numattr: " << nattr;
3202 if (options) {
3203 o << ", options: ";
3204 if (options & 0x01)
3205 o << "new";
3206 if (options & 0x10)
3207 o << "list values";
3208 }
3209 o << ", total size: " << req->dlen << ")";
3210 break;
3211 }
3212
3213 //------------------------------------------------------------------------
3214 // kXR_sync
3215 //------------------------------------------------------------------------
3216 case kXR_sync:
3217 {
3218 ClientSyncRequest *sreq = (ClientSyncRequest *)msg;
3219 o << "kXR_sync (";
3220 o << "handle: " << FileHandleToStr( sreq->fhandle );
3221 o << ")";
3222 break;
3223 }
3224
3225 //------------------------------------------------------------------------
3226 // kXR_truncate
3227 //------------------------------------------------------------------------
3228 case kXR_truncate:
3229 {
3230 ClientTruncateRequest *sreq = (ClientTruncateRequest *)msg;
3231 o << "kXR_truncate (";
3232 if( !sreq->dlen )
3233 o << "handle: " << FileHandleToStr( sreq->fhandle );
3234 else
3235 {
3236 char *fn = GetDataAsString( msg );
3237 o << "file: " << fn;
3238 delete [] fn;
3239 }
3240 o << std::setbase(10);
3241 o << ", ";
3242 o << "offset: " << sreq->offset;
3243 o << ")";
3244 break;
3245 }
3246
3247 //------------------------------------------------------------------------
3248 // kXR_readv
3249 //------------------------------------------------------------------------
3250 case kXR_readv:
3251 {
3252 unsigned char *fhandle = 0;
3253 o << "kXR_readv (";
3254
3255 o << "handle: ";
3256 readahead_list *dataChunk = (readahead_list*)(msg + 24 );
3257 fhandle = dataChunk[0].fhandle;
3258 if( fhandle )
3259 o << FileHandleToStr( fhandle );
3260 else
3261 o << "unknown";
3262 o << ", ";
3263 o << std::setbase(10);
3264 o << "chunks: [";
3265 uint64_t size = 0;
3266 for( size_t i = 0; i < req->dlen/sizeof(readahead_list); ++i )
3267 {
3268 size += dataChunk[i].rlen;
3269 o << "(offset: " << dataChunk[i].offset;
3270 o << ", size: " << dataChunk[i].rlen << "); ";
3271 }
3272 o << "], ";
3273 o << "total size: " << size << ")";
3274 break;
3275 }
3276
3277 //------------------------------------------------------------------------
3278 // kXR_writev
3279 //------------------------------------------------------------------------
3280 case kXR_writev:
3281 {
3282 unsigned char *fhandle = 0;
3283 o << "kXR_writev (";
3284
3285 XrdProto::write_list *wrtList =
3286 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
3287 uint64_t size = 0;
3288 uint32_t numChunks = 0;
3289 for( size_t i = 0; i < req->dlen/sizeof(XrdProto::write_list); ++i )
3290 {
3291 fhandle = wrtList[i].fhandle;
3292 size += wrtList[i].wlen;
3293 ++numChunks;
3294 }
3295 o << "handle: ";
3296 if( fhandle )
3297 o << FileHandleToStr( fhandle );
3298 else
3299 o << "unknown";
3300 o << ", ";
3301 o << std::setbase(10);
3302 o << "chunks: " << numChunks << ", ";
3303 o << "total size: " << size << ")";
3304 break;
3305 }
3306
3307 //------------------------------------------------------------------------
3308 // kXR_locate
3309 //------------------------------------------------------------------------
3310 case kXR_locate:
3311 {
3312 ClientLocateRequest *sreq = (ClientLocateRequest *)msg;
3313 char *fn = GetDataAsString( msg );;
3314 o << "kXR_locate (";
3315 o << "path: " << fn << ", ";
3316 delete [] fn;
3317 o << "flags: ";
3318 if( sreq->options == 0 )
3319 o << "none";
3320 else
3321 {
3322 if( sreq->options & kXR_refresh )
3323 o << "kXR_refresh ";
3324 if( sreq->options & kXR_prefname )
3325 o << "kXR_prefname ";
3326 if( sreq->options & kXR_nowait )
3327 o << "kXR_nowait ";
3328 if( sreq->options & kXR_force )
3329 o << "kXR_force ";
3330 if( sreq->options & kXR_compress )
3331 o << "kXR_compress ";
3332 }
3333 o << ")";
3334 break;
3335 }
3336
3337 //------------------------------------------------------------------------
3338 // kXR_mv
3339 //------------------------------------------------------------------------
3340 case kXR_mv:
3341 {
3342 ClientMvRequest *sreq = (ClientMvRequest *)msg;
3343 o << "kXR_mv (";
3344 o << "source: ";
3345 o.write( msg + sizeof( ClientMvRequest ), sreq->arg1len );
3346 o << ", ";
3347 o << "destination: ";
3348 o.write( msg + sizeof( ClientMvRequest ) + sreq->arg1len + 1, sreq->dlen - sreq->arg1len - 1 );
3349 o << ")";
3350 break;
3351 }
3352
3353 //------------------------------------------------------------------------
3354 // kXR_query
3355 //------------------------------------------------------------------------
3356 case kXR_query:
3357 {
3358 ClientQueryRequest *sreq = (ClientQueryRequest *)msg;
3359 o << "kXR_query (";
3360 o << "code: ";
3361 switch( sreq->infotype )
3362 {
3363 case kXR_Qconfig: o << "kXR_Qconfig"; break;
3364 case kXR_Qckscan: o << "kXR_Qckscan"; break;
3365 case kXR_Qcksum: o << "kXR_Qcksum"; break;
3366 case kXR_Qopaque: o << "kXR_Qopaque"; break;
3367 case kXR_Qopaquf: o << "kXR_Qopaquf"; break;
3368 case kXR_Qopaqug: o << "kXR_Qopaqug"; break;
3369 case kXR_QPrep: o << "kXR_QPrep"; break;
3370 case kXR_Qspace: o << "kXR_Qspace"; break;
3371 case kXR_QStats: o << "kXR_QStats"; break;
3372 case kXR_Qvisa: o << "kXR_Qvisa"; break;
3373 case kXR_Qxattr: o << "kXR_Qxattr"; break;
3374 default: o << sreq->infotype; break;
3375 }
3376 o << ", ";
3377
3378 if( sreq->infotype == kXR_Qopaqug || sreq->infotype == kXR_Qvisa )
3379 {
3380 o << "handle: " << FileHandleToStr( sreq->fhandle );
3381 o << ", ";
3382 }
3383
3384 o << "arg length: " << sreq->dlen << ")";
3385 break;
3386 }
3387
3388 //------------------------------------------------------------------------
3389 // kXR_rm
3390 //------------------------------------------------------------------------
3391 case kXR_rm:
3392 {
3393 o << "kXR_rm (";
3394 char *fn = GetDataAsString( msg );;
3395 o << "path: " << fn << ")";
3396 delete [] fn;
3397 break;
3398 }
3399
3400 //------------------------------------------------------------------------
3401 // kXR_mkdir
3402 //------------------------------------------------------------------------
3403 case kXR_mkdir:
3404 {
3405 ClientMkdirRequest *sreq = (ClientMkdirRequest *)msg;
3406 o << "kXR_mkdir (";
3407 char *fn = GetDataAsString( msg );
3408 o << "path: " << fn << ", ";
3409 delete [] fn;
3410 o << "mode: 0" << std::setbase(8) << sreq->mode << ", ";
3411 o << std::setbase(10);
3412 o << "flags: ";
3413 if( sreq->options[0] == 0 )
3414 o << "none";
3415 else
3416 {
3417 if( sreq->options[0] & kXR_mkdirpath )
3418 o << "kXR_mkdirpath";
3419 }
3420 o << ")";
3421 break;
3422 }
3423
3424 //------------------------------------------------------------------------
3425 // kXR_rmdir
3426 //------------------------------------------------------------------------
3427 case kXR_rmdir:
3428 {
3429 o << "kXR_rmdir (";
3430 char *fn = GetDataAsString( msg );
3431 o << "path: " << fn << ")";
3432 delete [] fn;
3433 break;
3434 }
3435
3436 //------------------------------------------------------------------------
3437 // kXR_chmod
3438 //------------------------------------------------------------------------
3439 case kXR_chmod:
3440 {
3441 ClientChmodRequest *sreq = (ClientChmodRequest *)msg;
3442 o << "kXR_chmod (";
3443 char *fn = GetDataAsString( msg );
3444 o << "path: " << fn << ", ";
3445 delete [] fn;
3446 o << "mode: 0" << std::setbase(8) << sreq->mode << ")";
3447 break;
3448 }
3449
3450 //------------------------------------------------------------------------
3451 // kXR_ping
3452 //------------------------------------------------------------------------
3453 case kXR_ping:
3454 {
3455 o << "kXR_ping ()";
3456 break;
3457 }
3458
3459 //------------------------------------------------------------------------
3460 // kXR_protocol
3461 //------------------------------------------------------------------------
3462 case kXR_protocol:
3463 {
3464 ClientProtocolRequest *sreq = (ClientProtocolRequest *)msg;
3465 o << "kXR_protocol (";
3466 o << "clientpv: 0x" << std::setbase(16) << sreq->clientpv << ")";
3467 break;
3468 }
3469
3470 //------------------------------------------------------------------------
3471 // kXR_dirlist
3472 //------------------------------------------------------------------------
3473 case kXR_dirlist:
3474 {
3475 o << "kXR_dirlist (";
3476 char *fn = GetDataAsString( msg );;
3477 o << "path: " << fn << ")";
3478 delete [] fn;
3479 break;
3480 }
3481
3482 //------------------------------------------------------------------------
3483 // kXR_set
3484 //------------------------------------------------------------------------
3485 case kXR_set:
3486 {
3487 o << "kXR_set (";
3488 char *fn = GetDataAsString( msg );;
3489 o << "data: " << fn << ")";
3490 delete [] fn;
3491 break;
3492 }
3493
3494 //------------------------------------------------------------------------
3495 // kXR_prepare
3496 //------------------------------------------------------------------------
3497 case kXR_prepare:
3498 {
3499 ClientPrepareRequest *sreq = (ClientPrepareRequest *)msg;
3500 o << "kXR_prepare (";
3501 o << "flags: ";
3502
3503 if( sreq->options == 0 )
3504 o << "none";
3505 else
3506 {
3507 if( sreq->options & kXR_stage )
3508 o << "kXR_stage ";
3509 if( sreq->options & kXR_wmode )
3510 o << "kXR_wmode ";
3511 if( sreq->options & kXR_coloc )
3512 o << "kXR_coloc ";
3513 if( sreq->options & kXR_fresh )
3514 o << "kXR_fresh ";
3515 }
3516
3517 o << ", priority: " << (int) sreq->prty << ", ";
3518
3519 char *fn = GetDataAsString( msg );
3520 char *cursor;
3521 for( cursor = fn; *cursor; ++cursor )
3522 if( *cursor == '\n' ) *cursor = ' ';
3523
3524 o << "paths: " << fn << ")";
3525 delete [] fn;
3526 break;
3527 }
3528
3529 case kXR_chkpoint:
3530 {
3531 ClientChkPointRequest *sreq = (ClientChkPointRequest*)msg;
3532 o << "kXR_chkpoint (";
3533 o << "opcode: ";
3534 if( sreq->opcode == kXR_ckpBegin ) o << "kXR_ckpBegin)";
3535 else if( sreq->opcode == kXR_ckpCommit ) o << "kXR_ckpCommit)";
3536 else if( sreq->opcode == kXR_ckpQuery ) o << "kXR_ckpQuery)";
3537 else if( sreq->opcode == kXR_ckpRollback ) o << "kXR_ckpRollback)";
3538 else if( sreq->opcode == kXR_ckpXeq )
3539 {
3540 o << "kXR_ckpXeq) ";
3541 // In this case our request body will be one of kXR_pgwrite,
3542 // kXR_truncate, kXR_write, or kXR_writev request.
3543 GenerateDescription( msg + sizeof( ClientChkPointRequest ), o );
3544 }
3545
3546 break;
3547 }
3548
3549 //------------------------------------------------------------------------
3550 // Default
3551 //------------------------------------------------------------------------
3552 default:
3553 {
3554 o << "kXR_unknown (length: " << req->dlen << ")";
3555 break;
3556 }
3557 };
3558 }
static const int kXR_ckpRollback
Definition XProtocol.hh:215
kXR_int16 arg1len
Definition XProtocol.hh:430
@ kXR_fattrDel
Definition XProtocol.hh:270
@ kXR_fattrSet
Definition XProtocol.hh:273
@ kXR_fattrList
Definition XProtocol.hh:272
@ kXR_fattrGet
Definition XProtocol.hh:271
kXR_char fhandle[4]
Definition XProtocol.hh:531
kXR_char fhandle[4]
Definition XProtocol.hh:782
kXR_char fhandle[4]
Definition XProtocol.hh:807
kXR_char fhandle[4]
Definition XProtocol.hh:771
kXR_int32 dlen
Definition XProtocol.hh:431
kXR_unt16 options
Definition XProtocol.hh:481
static const int kXR_ckpXeq
Definition XProtocol.hh:216
@ kXR_open_wrto
Definition XProtocol.hh:469
@ kXR_compress
Definition XProtocol.hh:452
@ kXR_async
Definition XProtocol.hh:458
@ kXR_delete
Definition XProtocol.hh:453
@ kXR_prefname
Definition XProtocol.hh:461
@ kXR_nowait
Definition XProtocol.hh:467
@ kXR_open_read
Definition XProtocol.hh:456
@ kXR_open_updt
Definition XProtocol.hh:457
@ kXR_mkpath
Definition XProtocol.hh:460
@ kXR_seqio
Definition XProtocol.hh:468
@ kXR_replica
Definition XProtocol.hh:465
@ kXR_posc
Definition XProtocol.hh:466
@ kXR_refresh
Definition XProtocol.hh:459
@ kXR_new
Definition XProtocol.hh:455
@ kXR_force
Definition XProtocol.hh:454
@ kXR_4dirlist
Definition XProtocol.hh:464
@ kXR_open_apnd
Definition XProtocol.hh:462
@ kXR_retstat
Definition XProtocol.hh:463
kXR_char fhandle[4]
Definition XProtocol.hh:509
kXR_char fhandle[4]
Definition XProtocol.hh:645
kXR_char fhandle[4]
Definition XProtocol.hh:659
kXR_char fhandle[4]
Definition XProtocol.hh:229
kXR_unt16 requestid
Definition XProtocol.hh:157
kXR_char fhandle[4]
Definition XProtocol.hh:633
@ kXR_read
Definition XProtocol.hh:125
@ kXR_open
Definition XProtocol.hh:122
@ kXR_writev
Definition XProtocol.hh:143
@ kXR_readv
Definition XProtocol.hh:137
@ kXR_mkdir
Definition XProtocol.hh:120
@ kXR_sync
Definition XProtocol.hh:128
@ kXR_chmod
Definition XProtocol.hh:114
@ kXR_dirlist
Definition XProtocol.hh:116
@ kXR_fattr
Definition XProtocol.hh:132
@ kXR_rm
Definition XProtocol.hh:126
@ kXR_query
Definition XProtocol.hh:113
@ kXR_write
Definition XProtocol.hh:131
@ kXR_set
Definition XProtocol.hh:130
@ kXR_rmdir
Definition XProtocol.hh:127
@ kXR_truncate
Definition XProtocol.hh:140
@ kXR_protocol
Definition XProtocol.hh:118
@ kXR_mv
Definition XProtocol.hh:121
@ kXR_ping
Definition XProtocol.hh:123
@ kXR_stat
Definition XProtocol.hh:129
@ kXR_pgread
Definition XProtocol.hh:142
@ kXR_chkpoint
Definition XProtocol.hh:124
@ kXR_locate
Definition XProtocol.hh:139
@ kXR_close
Definition XProtocol.hh:115
@ kXR_pgwrite
Definition XProtocol.hh:138
@ kXR_prepare
Definition XProtocol.hh:133
kXR_int32 rlen
Definition XProtocol.hh:660
kXR_char options[1]
Definition XProtocol.hh:416
static const int kXR_ckpCommit
Definition XProtocol.hh:213
kXR_int64 offset
Definition XProtocol.hh:661
@ kXR_vfs
Definition XProtocol.hh:763
@ kXR_mkdirpath
Definition XProtocol.hh:410
@ kXR_wmode
Definition XProtocol.hh:591
@ kXR_fresh
Definition XProtocol.hh:593
@ kXR_coloc
Definition XProtocol.hh:592
@ kXR_stage
Definition XProtocol.hh:590
static const int kXR_ckpQuery
Definition XProtocol.hh:214
@ kXR_QPrep
Definition XProtocol.hh:616
@ kXR_Qopaqug
Definition XProtocol.hh:625
@ kXR_Qconfig
Definition XProtocol.hh:621
@ kXR_Qopaquf
Definition XProtocol.hh:624
@ kXR_Qckscan
Definition XProtocol.hh:620
@ kXR_Qxattr
Definition XProtocol.hh:618
@ kXR_Qspace
Definition XProtocol.hh:619
@ kXR_Qvisa
Definition XProtocol.hh:622
@ kXR_QStats
Definition XProtocol.hh:615
@ kXR_Qcksum
Definition XProtocol.hh:617
@ kXR_Qopaque
Definition XProtocol.hh:623
static const int kXR_ckpBegin
Definition XProtocol.hh:212
@ 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:832
kXR_char fhandle[4]
Definition XProtocol.hh:288

References ClientMvRequest::arg1len, ClientProtocolRequest::clientpv, ClientMvRequest::dlen, ClientPgWriteRequest::dlen, ClientQueryRequest::dlen, ClientRequestHdr::dlen, ClientStatRequest::dlen, ClientTruncateRequest::dlen, ClientWriteRequest::dlen, XrdCl::Log::ErrorMsg, 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, GenerateDescription(), XrdCl::Log::GetLevel(), XrdCl::DefaultEnv::GetLog(), ClientQueryRequest::infotype, kXR_4dirlist, kXR_async, kXR_chkpoint, kXR_chmod, kXR_ckpBegin, kXR_ckpCommit, kXR_ckpQuery, kXR_ckpRollback, kXR_ckpXeq, kXR_close, kXR_coloc, kXR_compress, kXR_delete, kXR_dirlist, 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_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, ClientPrepareRequest::prty, ClientRequestHdr::requestid, ClientPgReadRequest::rlen, ClientReadRequest::rlen, readahead_list::rlen, 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 1923 of file XrdClXRootDTransport.cc.

1925 {
1926 XRootDChannelInfo *info = 0;
1927 channelData.Get( info );
1928
1929 if(!info || !info->bindSelector)
1930 return url;
1931
1932 return URL( info->bindSelector->Get() );
1933 }

References 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 346 of file XrdClXRootDTransport.cc.

347 {
348 //--------------------------------------------------------------------------
349 // Retrieve the body
350 //--------------------------------------------------------------------------
351 size_t leftToBeRead = 0;
352 uint32_t bodySize = 0;
353 ServerResponseHeader* rsphdr = (ServerResponseHeader*)message.GetBuffer();
354 bodySize = rsphdr->dlen;
355
356 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
357 return XRootDStatus( stError, errInvalidMessage, 0,
358 "Response body too large." );
359
360 if( message.GetSize() < bodySize + 8 )
361 message.ReAllocate( bodySize + 8 );
362
363 leftToBeRead = bodySize-(message.GetCursor()-8);
364 while( leftToBeRead )
365 {
366 int bytesRead = 0;
367 XRootDStatus status = socket->Read( message.GetBufferAtCursor(), leftToBeRead, bytesRead );
368
369 if( !status.IsOK() || status.code == suRetry )
370 return status;
371
372 leftToBeRead -= bytesRead;
373 message.AdvanceCursor( bytesRead );
374 }
375
376 return XRootDStatus( stOK, suDone );
377 }
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::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 306 of file XrdClXRootDTransport.cc.

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

References 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 382 of file XrdClXRootDTransport.cc.

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

1785 {
1786 XRootDChannelInfo *info = 0;
1787 channelData.Get( info );
1788 return GetSignature( toSign, sign, info );
1789 }
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 1794 of file XrdClXRootDTransport.cc.

1797 {
1798 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
1799 if( pSecUnloadHandler->unloaded ) return Status( stError, errInvalidOp );
1800
1801 ClientRequest *thereq = reinterpret_cast<ClientRequest*>( toSign->GetBuffer() );
1802 if( !info ) return Status( stError, errInternal );
1803 if( info->protection )
1804 {
1805 SecurityRequest *newreq = 0;
1806 // check if we have to secure the request in the first place
1807 if( !( NEED2SECURE ( info->protection )( *thereq ) ) ) return Status();
1808 // secure (sign/encrypt) the request
1809 int rc = info->protection->Secure( newreq, *thereq, 0 );
1810 // there was an error
1811 if( rc < 0 )
1812 return Status( stError, errInternal, -rc );
1813
1814 sign = new Message();
1815 sign->Grab( reinterpret_cast<char*>( newreq ), rc );
1816 }
1817
1818 return Status();
1819 }
#define NEED2SECURE(protP)
This class implements the XRootD protocol security protection.

References 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 477 of file XrdClXRootDTransport.cc.

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

References 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 756 of file XrdClXRootDTransport.cc.

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

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 449 of file XrdClXRootDTransport.cc.

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

References 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 829 of file XrdClXRootDTransport.cc.

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

References 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 776 of file XrdClXRootDTransport.cc.

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

1518 {
1519 Log *log = DefaultEnv::GetLog();
1520 ServerResponse *rsp = (ServerResponse *)msg.GetBuffer();
1521 char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
1522 memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
1523 log->Error( XRootDTransportMsg, "Server responded with an error [%d]: %s",
1524 rsp->body.error.errnum, errmsg );
1525 delete [] errmsg;
1526 }
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 1113 of file XrdClXRootDTransport.cc.

1114 {
1115 ClientRequest *req = (ClientRequest*)msg;
1116 switch( req->header.requestid )
1117 {
1118 //------------------------------------------------------------------------
1119 // kXR_protocol
1120 //------------------------------------------------------------------------
1121 case kXR_protocol:
1122 req->protocol.clientpv = htonl( req->protocol.clientpv );
1123 break;
1124
1125 //------------------------------------------------------------------------
1126 // kXR_login
1127 //------------------------------------------------------------------------
1128 case kXR_login:
1129 req->login.pid = htonl( req->login.pid );
1130 break;
1131
1132 //------------------------------------------------------------------------
1133 // kXR_locate
1134 //------------------------------------------------------------------------
1135 case kXR_locate:
1136 req->locate.options = htons( req->locate.options );
1137 break;
1138
1139 //------------------------------------------------------------------------
1140 // kXR_query
1141 //------------------------------------------------------------------------
1142 case kXR_query:
1143 req->query.infotype = htons( req->query.infotype );
1144 break;
1145
1146 //------------------------------------------------------------------------
1147 // kXR_truncate
1148 //------------------------------------------------------------------------
1149 case kXR_truncate:
1150 req->truncate.offset = htonll( req->truncate.offset );
1151 break;
1152
1153 //------------------------------------------------------------------------
1154 // kXR_mkdir
1155 //------------------------------------------------------------------------
1156 case kXR_mkdir:
1157 req->mkdir.mode = htons( req->mkdir.mode );
1158 break;
1159
1160 //------------------------------------------------------------------------
1161 // kXR_chmod
1162 //------------------------------------------------------------------------
1163 case kXR_chmod:
1164 req->chmod.mode = htons( req->chmod.mode );
1165 break;
1166
1167 //------------------------------------------------------------------------
1168 // kXR_open
1169 //------------------------------------------------------------------------
1170 case kXR_open:
1171 req->open.mode = htons( req->open.mode );
1172 req->open.options = htons( req->open.options );
1173 break;
1174
1175 //------------------------------------------------------------------------
1176 // kXR_read
1177 //------------------------------------------------------------------------
1178 case kXR_read:
1179 req->read.offset = htonll( req->read.offset );
1180 req->read.rlen = htonl( req->read.rlen );
1181 break;
1182
1183 //------------------------------------------------------------------------
1184 // kXR_write
1185 //------------------------------------------------------------------------
1186 case kXR_write:
1187 req->write.offset = htonll( req->write.offset );
1188 break;
1189
1190 //------------------------------------------------------------------------
1191 // kXR_mv
1192 //------------------------------------------------------------------------
1193 case kXR_mv:
1194 req->mv.arg1len = htons( req->mv.arg1len );
1195 break;
1196
1197 //------------------------------------------------------------------------
1198 // kXR_readv
1199 //------------------------------------------------------------------------
1200 case kXR_readv:
1201 {
1202 uint16_t numChunks = (req->readv.dlen)/16;
1203 readahead_list *dataChunk = (readahead_list*)( msg + 24 );
1204 for( size_t i = 0; i < numChunks; ++i )
1205 {
1206 dataChunk[i].rlen = htonl( dataChunk[i].rlen );
1207 dataChunk[i].offset = htonll( dataChunk[i].offset );
1208 }
1209 break;
1210 }
1211
1212 //------------------------------------------------------------------------
1213 // kXR_writev
1214 //------------------------------------------------------------------------
1215 case kXR_writev:
1216 {
1217 uint16_t numChunks = (req->writev.dlen)/16;
1218 XrdProto::write_list *wrtList =
1219 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
1220 for( size_t i = 0; i < numChunks; ++i )
1221 {
1222 wrtList[i].wlen = htonl( wrtList[i].wlen );
1223 wrtList[i].offset = htonll( wrtList[i].offset );
1224 }
1225
1226 break;
1227 }
1228
1229 case kXR_pgread:
1230 {
1231 req->pgread.offset = htonll( req->pgread.offset );
1232 req->pgread.rlen = htonl( req->pgread.rlen );
1233 break;
1234 }
1235
1236 case kXR_pgwrite:
1237 {
1238 req->pgwrite.offset = htonll( req->pgwrite.offset );
1239 break;
1240 }
1241
1242 //------------------------------------------------------------------------
1243 // kXR_prepare
1244 //------------------------------------------------------------------------
1245 case kXR_prepare:
1246 {
1247 req->prepare.optionX = htons( req->prepare.optionX );
1248 req->prepare.port = htons( req->prepare.port );
1249 break;
1250 }
1251
1252 case kXR_chkpoint:
1253 {
1254 if( req->chkpoint.opcode == kXR_ckpXeq )
1255 MarshallRequest( msg + 24 );
1256 break;
1257 }
1258 };
1259
1260 req->header.requestid = htons( req->header.requestid );
1261 req->header.dlen = htonl( req->header.dlen );
1262 return XRootDStatus();
1263 }
struct ClientTruncateRequest truncate
Definition XProtocol.hh:875
struct ClientPgReadRequest pgread
Definition XProtocol.hh:861
struct ClientMkdirRequest mkdir
Definition XProtocol.hh:858
struct ClientPgWriteRequest pgwrite
Definition XProtocol.hh:862
struct ClientReadVRequest readv
Definition XProtocol.hh:868
struct ClientOpenRequest open
Definition XProtocol.hh:860
struct ClientRequestHdr header
Definition XProtocol.hh:846
struct ClientWriteVRequest writev
Definition XProtocol.hh:877
struct ClientLoginRequest login
Definition XProtocol.hh:857
@ kXR_login
Definition XProtocol.hh:119
struct ClientChmodRequest chmod
Definition XProtocol.hh:850
struct ClientQueryRequest query
Definition XProtocol.hh:866
struct ClientReadRequest read
Definition XProtocol.hh:867
struct ClientMvRequest mv
Definition XProtocol.hh:859
struct ClientChkPointRequest chkpoint
Definition XProtocol.hh:849
struct ClientPrepareRequest prepare
Definition XProtocol.hh:864
struct ClientWriteRequest write
Definition XProtocol.hh:876
struct ClientProtocolRequest protocol
Definition XProtocol.hh:865
struct ClientLocateRequest locate
Definition XProtocol.hh:856
static XRootDStatus MarshallRequest(Message *msg)
Marshal the outgoing message.

References ClientMvRequest::arg1len, ClientRequest::chkpoint, ClientRequest::chmod, ClientProtocolRequest::clientpv, ClientReadVRequest::dlen, ClientRequestHdr::dlen, ClientWriteVRequest::dlen, ClientRequest::header, ClientQueryRequest::infotype, kXR_chkpoint, kXR_chmod, kXR_ckpXeq, 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, 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, ClientRequest::truncate, XrdProto::write_list::wlen, ClientRequest::write, and ClientRequest::writev.

+ Here is the call graph for this function:

◆ MarshallRequest() [2/2]

static 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::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 1640 of file XrdClXRootDTransport.cc.

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

1753 {
1754 // Called when a message has been sent. For messages that return on a
1755 // different pathid (and hence may use a different poller) it is possible
1756 // that the server has already replied and the reply will trigger
1757 // MessageReceived() before this method has been called. However for open
1758 // and close this is never the case and this method is used for tracking
1759 // only those.
1760 XRootDChannelInfo *info = 0;
1761 channelData.Get( info );
1762 if( !info ) return;
1763 XrdSysMutexHelper scopedLock( info->mutex );
1764 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1765 uint16_t reqid = ntohs( req->header.requestid );
1766
1767
1768 //--------------------------------------------------------------------------
1769 // We need to track opens to know if we can close streams due to idleness
1770 //--------------------------------------------------------------------------
1771 uint16_t sid;
1772 memcpy( &sid, req->header.streamid, 2 );
1773
1774 if( reqid == kXR_open )
1775 info->sentOpens.insert( sid );
1776 else if( reqid == kXR_close )
1777 info->sentCloses.insert( sid );
1778 }
kXR_char streamid[2]
Definition XProtocol.hh:156

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 872 of file XrdClXRootDTransport.cc.

873 {
874 return PathID( 0, 0 );
875 }

◆ 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 880 of file XrdClXRootDTransport.cc.

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

References 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 1531 of file XrdClXRootDTransport.cc.

1532 {
1533 XRootDChannelInfo *info = 0;
1534 channelData.Get( info );
1535
1536 if (!info) {
1537 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1538 return 0;
1539 }
1540
1541 XrdSysMutexHelper scopedLock( info->mutex );
1542
1543 uint16_t nbConnected = 0;
1544 for( size_t i = 1; i < info->stream.size(); ++i )
1545 if( info->stream[i].status == XRootDStreamInfo::Connected )
1546 ++nbConnected;
1547
1548 return nbConnected;
1549 }

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 1844 of file XrdClXRootDTransport.cc.

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

1591 {
1592 XRootDChannelInfo *info = 0;
1593 channelData.Get( info );
1594
1595 if (!info)
1596 return XRootDStatus(stFatal, errInternal);
1597
1598 XrdSysMutexHelper scopedLock( info->mutex );
1599
1600 switch( query )
1601 {
1602 //------------------------------------------------------------------------
1603 // Protocol name
1604 //------------------------------------------------------------------------
1606 result.Set( (const char*)"XRootD", false );
1607 return Status();
1608
1609 //------------------------------------------------------------------------
1610 // Authentication
1611 //------------------------------------------------------------------------
1613 result.Set( new std::string( info->authProtocolName ), false );
1614 return Status();
1615
1616 //------------------------------------------------------------------------
1617 // Server flags
1618 //------------------------------------------------------------------------
1620 result.Set( new int( info->serverFlags ), false );
1621 return Status();
1622
1623 //------------------------------------------------------------------------
1624 // Protocol version
1625 //------------------------------------------------------------------------
1627 result.Set( new int( info->protocolVersion ), false );
1628 return Status();
1629
1631 result.Set( new bool( info->encrypted ), false );
1632 return Status();
1633 };
1634 return Status( stError, errQueryNotSupported );
1635 }
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::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()

static 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::Close(), XrdCl::FileSystem::DirList(), XrdCl::FileStateHandler::Fcntl(), XrdCl::FileSystem::Locate(), XrdCl::FileSystem::MkDir(), XrdCl::FileSystem::Mv(), XrdCl::FileStateHandler::Open(), XrdCl::FileStateHandler::PgReadImpl(), XrdCl::FileStateHandler::PgWriteImpl(), XrdCl::FileSystem::Ping(), XrdCl::FileSystem::Prepare(), 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 1052 of file XrdClXRootDTransport.cc.

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

1445 {
1446 ServerResponseV2 *rsp = (ServerResponseV2*)msg.GetBuffer();
1447 uint16_t reqType = rsp->status.bdy.requestid + kXR_1stRequest;
1448
1449 switch( reqType )
1450 {
1451 case kXR_pgwrite:
1452 {
1453 //--------------------------------------------------------------------------
1454 // If there's no additional data there's nothing to unmarshal
1455 //--------------------------------------------------------------------------
1456 if( rsp->status.bdy.dlen == 0 ) return XRootDStatus();
1457 //--------------------------------------------------------------------------
1458 // If there's not enough data to form correction-segment report an error
1459 //--------------------------------------------------------------------------
1460 if( size_t( rsp->status.bdy.dlen ) < sizeof( ServerResponseBody_pgWrCSE ) )
1461 return XRootDStatus( stError, errInvalidMessage, 0,
1462 "kXR_status: invalid message size." );
1463
1464 //--------------------------------------------------------------------------
1465 // Calculate the crc32c for the additional data
1466 //--------------------------------------------------------------------------
1467 ServerResponseBody_pgWrCSE *cse = (ServerResponseBody_pgWrCSE*)msg.GetBuffer( sizeof( ServerResponseV2 ) );
1468 cse->cseCRC = ntohl( cse->cseCRC );
1469 size_t length = rsp->status.bdy.dlen - sizeof( uint32_t );
1470 void* buffer = msg.GetBuffer( sizeof( ServerResponseV2 ) + sizeof( uint32_t ) );
1471 uint32_t crcval = XrdOucCRC::Calc32C( buffer, length );
1472
1473 //--------------------------------------------------------------------------
1474 // Do the integrity checks
1475 //--------------------------------------------------------------------------
1476 if( crcval != cse->cseCRC )
1477 {
1478 return XRootDStatus( stError, errDataError, 0, "kXR_status response header "
1479 "corrupted (crc32c integrity check failed)." );
1480 }
1481
1482 cse->dlFirst = ntohs( cse->dlFirst );
1483 cse->dlLast = ntohs( cse->dlLast );
1484
1485 size_t pgcnt = ( rsp->status.bdy.dlen - sizeof( ServerResponseBody_pgWrCSE ) ) /
1486 sizeof( kXR_int64 );
1487 kXR_int64 *pgoffs = (kXR_int64*)msg.GetBuffer( sizeof( ServerResponseV2 ) +
1488 sizeof( ServerResponseBody_pgWrCSE ) );
1489
1490 for( size_t i = 0; i < pgcnt; ++i )
1491 pgoffs[i] = ntohll( pgoffs[i] );
1492
1493 return XRootDStatus();
1494 break;
1495 }
1496
1497 default:
1498 break;
1499 }
1500
1501 return XRootDStatus( stError, errNotSupported );
1502 }
ServerResponseStatus status
@ kXR_1stRequest
Definition XProtocol.hh:111
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 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 1290 of file XrdClXRootDTransport.cc.

1291 {
1292 ServerResponse *m = (ServerResponse *)msg->GetBuffer();
1293
1294 //--------------------------------------------------------------------------
1295 // kXR_ok
1296 //--------------------------------------------------------------------------
1297 if( m->hdr.status == kXR_ok )
1298 {
1299 switch( reqType )
1300 {
1301 //----------------------------------------------------------------------
1302 // kXR_protocol
1303 //----------------------------------------------------------------------
1304 case kXR_protocol:
1305 if( m->hdr.dlen < 8 )
1306 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_protocol: body too short." );
1307 m->body.protocol.pval = ntohl( m->body.protocol.pval );
1308 m->body.protocol.flags = ntohl( m->body.protocol.flags );
1309 break;
1310 }
1311 }
1312 //--------------------------------------------------------------------------
1313 // kXR_error
1314 //--------------------------------------------------------------------------
1315 else if( m->hdr.status == kXR_error )
1316 {
1317 if( m->hdr.dlen < 4 )
1318 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_error: body too short." );
1319 m->body.error.errnum = ntohl( m->body.error.errnum );
1320 }
1321
1322 //--------------------------------------------------------------------------
1323 // kXR_wait
1324 //--------------------------------------------------------------------------
1325 else if( m->hdr.status == kXR_wait )
1326 {
1327 if( m->hdr.dlen < 4 )
1328 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_wait: body too short." );
1329 m->body.wait.seconds = htonl( m->body.wait.seconds );
1330 }
1331
1332 //--------------------------------------------------------------------------
1333 // kXR_redirect
1334 //--------------------------------------------------------------------------
1335 else if( m->hdr.status == kXR_redirect )
1336 {
1337 if( m->hdr.dlen < 4 )
1338 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_redirect: body too short." );
1339 m->body.redirect.port = htonl( m->body.redirect.port );
1340 }
1341
1342 //--------------------------------------------------------------------------
1343 // kXR_waitresp
1344 //--------------------------------------------------------------------------
1345 else if( m->hdr.status == kXR_waitresp )
1346 {
1347 if( m->hdr.dlen < 4 )
1348 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_waitresp: body too short." );
1349 m->body.waitresp.seconds = htonl( m->body.waitresp.seconds );
1350 }
1351
1352 //--------------------------------------------------------------------------
1353 // kXR_attn
1354 //--------------------------------------------------------------------------
1355 else if( m->hdr.status == kXR_attn )
1356 {
1357 if( m->hdr.dlen < 4 )
1358 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_attn: body too short." );
1359 m->body.attn.actnum = htonl( m->body.attn.actnum );
1360 }
1361
1362 return XRootDStatus();
1363 }
@ kXR_redirect
Definition XProtocol.hh:904
@ kXR_error
Definition XProtocol.hh:903

References 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 1507 of file XrdClXRootDTransport.cc.

1508 {
1509 ServerResponseHeader *header = (ServerResponseHeader *)msg.GetBuffer();
1510 header->status = ntohs( header->status );
1511 header->dlen = ntohl( header->dlen );
1512 }

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 1269 of file XrdClXRootDTransport.cc.

1270 {
1271 if( !msg->IsMarshalled() ) return XRootDStatus( stOK, suAlreadyDone );
1272 // We rely on the marshaling process to be symmetric!
1273 // First we unmarshall the request ID and the length because
1274 // MarshallRequest() relies on these, and then we need to unmarshall these
1275 // two again, because they get marshalled in MarshallRequest().
1276 // All this is pretty damn ugly and should be rewritten.
1277 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1278 req->header.requestid = htons( req->header.requestid );
1279 req->header.dlen = htonl( req->header.dlen );
1280 XRootDStatus st = MarshallRequest( msg );
1281 req->header.requestid = htons( req->header.requestid );
1282 req->header.dlen = htonl( req->header.dlen );
1283 msg->SetIsMarshalled( false );
1284 return st;
1285 }
const uint16_t suAlreadyDone

References 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 1368 of file XrdClXRootDTransport.cc.

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

References 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 1835 of file XrdClXRootDTransport.cc.

1836 {
1837 XrdSysRWLockHelper scope( pSecUnloadHandler->lock, false ); // obtain write lock
1838 pSecUnloadHandler->unloaded = true;
1839 }

Friends And Related Symbol Documentation

◆ 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: