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

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

References PluginUnloadHandler.

+ Here is the call graph for this function:

◆ ~XRootDTransport()

XrdCl::XRootDTransport::~XRootDTransport ( )

Destructor.

Definition at line 299 of file XrdClXRootDTransport.cc.

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

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

1815 {
1816 XRootDChannelInfo *info = 0;
1817 channelData.Get( info );
1818 if( info->finstcnt.load( std::memory_order_relaxed ) > 0 )
1819 info->finstcnt.fetch_sub( 1, std::memory_order_relaxed );
1820 }

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

1546 {
1547 XRootDChannelInfo *info = 0;
1548 channelData.Get( info );
1549
1550 if (!info) {
1551 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1552 return;
1553 }
1554
1555 XrdSysMutexHelper scopedLock( info->mutex );
1556
1557 CleanUpProtection( info );
1558
1559 if( !info->stream.empty() )
1560 {
1561 XRootDStreamInfo &sInfo = info->stream[subStreamId];
1562 sInfo.status = XRootDStreamInfo::Disconnected;
1563 }
1564
1565 if( subStreamId == 0 )
1566 {
1567 info->sidManager->ReleaseAllTimedOut();
1568 info->sentOpens.clear();
1569 info->sentCloses.clear();
1570 info->openFiles = 0;
1571 info->waitBarrier = 0;
1572 }
1573 }
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 460 of file XrdClXRootDTransport.cc.

461 {
462 }

◆ GenerateDescription()

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

Get the description of a message.

Definition at line 2977 of file XrdClXRootDTransport.cc.

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

1915 {
1916 XRootDChannelInfo *info = 0;
1917 channelData.Get( info );
1918
1919 if(!info || !info->bindSelector)
1920 return url;
1921
1922 return URL( info->bindSelector->Get() );
1923 }

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

348 {
349 //--------------------------------------------------------------------------
350 // Retrieve the body
351 //--------------------------------------------------------------------------
352 size_t leftToBeRead = 0;
353 uint32_t bodySize = 0;
354 ServerResponseHeader* rsphdr = (ServerResponseHeader*)message.GetBuffer();
355 bodySize = rsphdr->dlen;
356
357 if( message.GetSize() < bodySize + 8 )
358 message.ReAllocate( bodySize + 8 );
359
360 leftToBeRead = bodySize-(message.GetCursor()-8);
361 while( leftToBeRead )
362 {
363 int bytesRead = 0;
364 XRootDStatus status = socket->Read( message.GetBufferAtCursor(), leftToBeRead, bytesRead );
365
366 if( !status.IsOK() || status.code == suRetry )
367 return status;
368
369 leftToBeRead -= bytesRead;
370 message.AdvanceCursor( bytesRead );
371 }
372
373 return XRootDStatus( stOK, suDone );
374 }
const uint16_t suRetry
const uint16_t stOK
Everything went OK.
const uint16_t suDone

References XrdCl::Buffer::AdvanceCursor(), XrdCl::Status::code, ServerResponseHeader::dlen, XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::GetBufferAtCursor(), XrdCl::Buffer::GetCursor(), XrdCl::Buffer::GetSize(), XrdCl::Status::IsOK(), XrdCl::Socket::Read(), XrdCl::Buffer::ReAllocate(), 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 307 of file XrdClXRootDTransport.cc.

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

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

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

1775 {
1776 XRootDChannelInfo *info = 0;
1777 channelData.Get( info );
1778 return GetSignature( toSign, sign, info );
1779 }
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 1784 of file XrdClXRootDTransport.cc.

1787 {
1788 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
1789 if( pSecUnloadHandler->unloaded ) return Status( stError, errInvalidOp );
1790
1791 ClientRequest *thereq = reinterpret_cast<ClientRequest*>( toSign->GetBuffer() );
1792 if( !info ) return Status( stError, errInternal );
1793 if( info->protection )
1794 {
1795 SecurityRequest *newreq = 0;
1796 // check if we have to secure the request in the first place
1797 if( !( NEED2SECURE ( info->protection )( *thereq ) ) ) return Status();
1798 // secure (sign/encrypt) the request
1799 int rc = info->protection->Secure( newreq, *thereq, 0 );
1800 // there was an error
1801 if( rc < 0 )
1802 return Status( stError, errInternal, -rc );
1803
1804 sign = new Message();
1805 sign->Grab( reinterpret_cast<char*>( newreq ), rc );
1806 }
1807
1808 return Status();
1809 }
#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 467 of file XrdClXRootDTransport.cc.

469 {
470 XRootDChannelInfo *info = 0;
471 channelData.Get( info );
472
473 if (!info)
474 return XRootDStatus(stFatal, errInternal);
475
476 XrdSysMutexHelper scopedLock( info->mutex );
477
478 if( info->stream.size() <= handShakeData->subStreamId )
479 {
480 Log *log = DefaultEnv::GetLog();
481 log->Error( XRootDTransportMsg,
482 "[%s] Internal error: not enough substreams",
483 handShakeData->streamName.c_str() );
484 return XRootDStatus( stFatal, errInternal );
485 }
486
487 if( handShakeData->subStreamId == 0 )
488 {
489 info->streamName = handShakeData->streamName;
490 return HandShakeMain( handShakeData, channelData );
491 }
492 return HandShakeParallel( handShakeData, channelData );
493 }
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 746 of file XrdClXRootDTransport.cc.

748 {
749 XRootDChannelInfo *info = 0;
750 channelData.Get( info );
751
752 if (!info) {
754 "[%s] Internal error: no channel info",
755 handShakeData->streamName.c_str());
756 return false;
757 }
758
759 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
760 return ( sInfo.status == XRootDStreamInfo::Connected );
761 }

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

441 {
442 XRootDChannelInfo *info = new XRootDChannelInfo( url );
443 XrdSysMutexHelper scopedLock( info->mutex );
444 channelData.Set( info );
445
446 Env *env = DefaultEnv::GetEnv();
447 int streams = DefaultSubStreamsPerChannel;
448 env->GetInt( "SubStreamsPerChannel", streams );
449 if( streams < 1 ) streams = 1;
450 info->stream.resize( streams );
451 info->strmSelector.reset( new StreamSelector( streams ) );
452 info->encrypted = url.IsSecure();
453 info->istpc = url.IsTPC();
454 info->logintoken = url.GetLoginToken();
455 }
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 819 of file XrdClXRootDTransport.cc.

821 {
822 XRootDChannelInfo *info = 0;
823 channelData.Get( info );
824 Env *env = DefaultEnv::GetEnv();
825 Log *log = DefaultEnv::GetLog();
826
827 if (!info) {
828 log->Error(XRootDTransportMsg,
829 "Internal error: no channel info, behaving as if stream is broken");
830 return true;
831 }
832
833 int streamTimeout = DefaultStreamTimeout;
834 env->GetInt( "StreamTimeout", streamTimeout );
835
836 XrdSysMutexHelper scopedLock( info->mutex );
837
838 const time_t now = time(0);
839 const bool anySID =
840 info->sidManager->IsAnySIDOldAs( now - streamTimeout );
841
842 log->Dump( XRootDTransportMsg, "[%s] Stream inactive since %lld seconds, "
843 "stream timeout: %d, any SID: %d, wait barrier: %s",
844 info->streamName.c_str(), (long long) inactiveTime, streamTimeout,
845 anySID, Utils::TimeToString(info->waitBarrier).c_str() );
846
847 if( inactiveTime < streamTimeout )
848 return Status();
849
850 if( now < info->waitBarrier )
851 return Status();
852
853 if( !anySID )
854 return Status();
855
856 return Status( stError, errSocketTimeout );
857 }
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 766 of file XrdClXRootDTransport.cc.

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

1508 {
1509 Log *log = DefaultEnv::GetLog();
1510 ServerResponse *rsp = (ServerResponse *)msg.GetBuffer();
1511 char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
1512 memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
1513 log->Error( XRootDTransportMsg, "Server responded with an error [%d]: %s",
1514 rsp->body.error.errnum, errmsg );
1515 delete [] errmsg;
1516 }
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 1103 of file XrdClXRootDTransport.cc.

1104 {
1105 ClientRequest *req = (ClientRequest*)msg;
1106 switch( req->header.requestid )
1107 {
1108 //------------------------------------------------------------------------
1109 // kXR_protocol
1110 //------------------------------------------------------------------------
1111 case kXR_protocol:
1112 req->protocol.clientpv = htonl( req->protocol.clientpv );
1113 break;
1114
1115 //------------------------------------------------------------------------
1116 // kXR_login
1117 //------------------------------------------------------------------------
1118 case kXR_login:
1119 req->login.pid = htonl( req->login.pid );
1120 break;
1121
1122 //------------------------------------------------------------------------
1123 // kXR_locate
1124 //------------------------------------------------------------------------
1125 case kXR_locate:
1126 req->locate.options = htons( req->locate.options );
1127 break;
1128
1129 //------------------------------------------------------------------------
1130 // kXR_query
1131 //------------------------------------------------------------------------
1132 case kXR_query:
1133 req->query.infotype = htons( req->query.infotype );
1134 break;
1135
1136 //------------------------------------------------------------------------
1137 // kXR_truncate
1138 //------------------------------------------------------------------------
1139 case kXR_truncate:
1140 req->truncate.offset = htonll( req->truncate.offset );
1141 break;
1142
1143 //------------------------------------------------------------------------
1144 // kXR_mkdir
1145 //------------------------------------------------------------------------
1146 case kXR_mkdir:
1147 req->mkdir.mode = htons( req->mkdir.mode );
1148 break;
1149
1150 //------------------------------------------------------------------------
1151 // kXR_chmod
1152 //------------------------------------------------------------------------
1153 case kXR_chmod:
1154 req->chmod.mode = htons( req->chmod.mode );
1155 break;
1156
1157 //------------------------------------------------------------------------
1158 // kXR_open
1159 //------------------------------------------------------------------------
1160 case kXR_open:
1161 req->open.mode = htons( req->open.mode );
1162 req->open.options = htons( req->open.options );
1163 break;
1164
1165 //------------------------------------------------------------------------
1166 // kXR_read
1167 //------------------------------------------------------------------------
1168 case kXR_read:
1169 req->read.offset = htonll( req->read.offset );
1170 req->read.rlen = htonl( req->read.rlen );
1171 break;
1172
1173 //------------------------------------------------------------------------
1174 // kXR_write
1175 //------------------------------------------------------------------------
1176 case kXR_write:
1177 req->write.offset = htonll( req->write.offset );
1178 break;
1179
1180 //------------------------------------------------------------------------
1181 // kXR_mv
1182 //------------------------------------------------------------------------
1183 case kXR_mv:
1184 req->mv.arg1len = htons( req->mv.arg1len );
1185 break;
1186
1187 //------------------------------------------------------------------------
1188 // kXR_readv
1189 //------------------------------------------------------------------------
1190 case kXR_readv:
1191 {
1192 uint16_t numChunks = (req->readv.dlen)/16;
1193 readahead_list *dataChunk = (readahead_list*)( msg + 24 );
1194 for( size_t i = 0; i < numChunks; ++i )
1195 {
1196 dataChunk[i].rlen = htonl( dataChunk[i].rlen );
1197 dataChunk[i].offset = htonll( dataChunk[i].offset );
1198 }
1199 break;
1200 }
1201
1202 //------------------------------------------------------------------------
1203 // kXR_writev
1204 //------------------------------------------------------------------------
1205 case kXR_writev:
1206 {
1207 uint16_t numChunks = (req->writev.dlen)/16;
1208 XrdProto::write_list *wrtList =
1209 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
1210 for( size_t i = 0; i < numChunks; ++i )
1211 {
1212 wrtList[i].wlen = htonl( wrtList[i].wlen );
1213 wrtList[i].offset = htonll( wrtList[i].offset );
1214 }
1215
1216 break;
1217 }
1218
1219 case kXR_pgread:
1220 {
1221 req->pgread.offset = htonll( req->pgread.offset );
1222 req->pgread.rlen = htonl( req->pgread.rlen );
1223 break;
1224 }
1225
1226 case kXR_pgwrite:
1227 {
1228 req->pgwrite.offset = htonll( req->pgwrite.offset );
1229 break;
1230 }
1231
1232 //------------------------------------------------------------------------
1233 // kXR_prepare
1234 //------------------------------------------------------------------------
1235 case kXR_prepare:
1236 {
1237 req->prepare.optionX = htons( req->prepare.optionX );
1238 req->prepare.port = htons( req->prepare.port );
1239 break;
1240 }
1241
1242 case kXR_chkpoint:
1243 {
1244 if( req->chkpoint.opcode == kXR_ckpXeq )
1245 MarshallRequest( msg + 24 );
1246 break;
1247 }
1248 };
1249
1250 req->header.requestid = htons( req->header.requestid );
1251 req->header.dlen = htonl( req->header.dlen );
1252 return XRootDStatus();
1253 }
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 1630 of file XrdClXRootDTransport.cc.

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

1743 {
1744 // Called when a message has been sent. For messages that return on a
1745 // different pathid (and hence may use a different poller) it is possible
1746 // that the server has already replied and the reply will trigger
1747 // MessageReceived() before this method has been called. However for open
1748 // and close this is never the case and this method is used for tracking
1749 // only those.
1750 XRootDChannelInfo *info = 0;
1751 channelData.Get( info );
1752 if( !info ) return;
1753 XrdSysMutexHelper scopedLock( info->mutex );
1754 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1755 uint16_t reqid = ntohs( req->header.requestid );
1756
1757
1758 //--------------------------------------------------------------------------
1759 // We need to track opens to know if we can close streams due to idleness
1760 //--------------------------------------------------------------------------
1761 uint16_t sid;
1762 memcpy( &sid, req->header.streamid, 2 );
1763
1764 if( reqid == kXR_open )
1765 info->sentOpens.insert( sid );
1766 else if( reqid == kXR_close )
1767 info->sentCloses.insert( sid );
1768 }
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 862 of file XrdClXRootDTransport.cc.

863 {
864 return PathID( 0, 0 );
865 }

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

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

1522 {
1523 XRootDChannelInfo *info = 0;
1524 channelData.Get( info );
1525
1526 if (!info) {
1527 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1528 return 0;
1529 }
1530
1531 XrdSysMutexHelper scopedLock( info->mutex );
1532
1533 uint16_t nbConnected = 0;
1534 for( size_t i = 1; i < info->stream.size(); ++i )
1535 if( info->stream[i].status == XRootDStreamInfo::Connected )
1536 ++nbConnected;
1537
1538 return nbConnected;
1539 }

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

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

1581 {
1582 XRootDChannelInfo *info = 0;
1583 channelData.Get( info );
1584
1585 if (!info)
1586 return XRootDStatus(stFatal, errInternal);
1587
1588 XrdSysMutexHelper scopedLock( info->mutex );
1589
1590 switch( query )
1591 {
1592 //------------------------------------------------------------------------
1593 // Protocol name
1594 //------------------------------------------------------------------------
1596 result.Set( (const char*)"XRootD", false );
1597 return Status();
1598
1599 //------------------------------------------------------------------------
1600 // Authentication
1601 //------------------------------------------------------------------------
1603 result.Set( new std::string( info->authProtocolName ), false );
1604 return Status();
1605
1606 //------------------------------------------------------------------------
1607 // Server flags
1608 //------------------------------------------------------------------------
1610 result.Set( new int( info->serverFlags ), false );
1611 return Status();
1612
1613 //------------------------------------------------------------------------
1614 // Protocol version
1615 //------------------------------------------------------------------------
1617 result.Set( new int( info->protocolVersion ), false );
1618 return Status();
1619
1621 result.Set( new bool( info->encrypted ), false );
1622 return Status();
1623 };
1624 return Status( stError, errQueryNotSupported );
1625 }
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 1042 of file XrdClXRootDTransport.cc.

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

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

1281 {
1282 ServerResponse *m = (ServerResponse *)msg->GetBuffer();
1283
1284 //--------------------------------------------------------------------------
1285 // kXR_ok
1286 //--------------------------------------------------------------------------
1287 if( m->hdr.status == kXR_ok )
1288 {
1289 switch( reqType )
1290 {
1291 //----------------------------------------------------------------------
1292 // kXR_protocol
1293 //----------------------------------------------------------------------
1294 case kXR_protocol:
1295 if( m->hdr.dlen < 8 )
1296 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_protocol: body too short." );
1297 m->body.protocol.pval = ntohl( m->body.protocol.pval );
1298 m->body.protocol.flags = ntohl( m->body.protocol.flags );
1299 break;
1300 }
1301 }
1302 //--------------------------------------------------------------------------
1303 // kXR_error
1304 //--------------------------------------------------------------------------
1305 else if( m->hdr.status == kXR_error )
1306 {
1307 if( m->hdr.dlen < 4 )
1308 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_error: body too short." );
1309 m->body.error.errnum = ntohl( m->body.error.errnum );
1310 }
1311
1312 //--------------------------------------------------------------------------
1313 // kXR_wait
1314 //--------------------------------------------------------------------------
1315 else if( m->hdr.status == kXR_wait )
1316 {
1317 if( m->hdr.dlen < 4 )
1318 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_wait: body too short." );
1319 m->body.wait.seconds = htonl( m->body.wait.seconds );
1320 }
1321
1322 //--------------------------------------------------------------------------
1323 // kXR_redirect
1324 //--------------------------------------------------------------------------
1325 else if( m->hdr.status == kXR_redirect )
1326 {
1327 if( m->hdr.dlen < 4 )
1328 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_redirect: body too short." );
1329 m->body.redirect.port = htonl( m->body.redirect.port );
1330 }
1331
1332 //--------------------------------------------------------------------------
1333 // kXR_waitresp
1334 //--------------------------------------------------------------------------
1335 else if( m->hdr.status == kXR_waitresp )
1336 {
1337 if( m->hdr.dlen < 4 )
1338 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_waitresp: body too short." );
1339 m->body.waitresp.seconds = htonl( m->body.waitresp.seconds );
1340 }
1341
1342 //--------------------------------------------------------------------------
1343 // kXR_attn
1344 //--------------------------------------------------------------------------
1345 else if( m->hdr.status == kXR_attn )
1346 {
1347 if( m->hdr.dlen < 4 )
1348 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_attn: body too short." );
1349 m->body.attn.actnum = htonl( m->body.attn.actnum );
1350 }
1351
1352 return XRootDStatus();
1353 }
@ 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 1497 of file XrdClXRootDTransport.cc.

1498 {
1499 ServerResponseHeader *header = (ServerResponseHeader *)msg.GetBuffer();
1500 header->status = ntohs( header->status );
1501 header->dlen = ntohl( header->dlen );
1502 }

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

1260 {
1261 if( !msg->IsMarshalled() ) return XRootDStatus( stOK, suAlreadyDone );
1262 // We rely on the marshaling process to be symmetric!
1263 // First we unmarshall the request ID and the length because
1264 // MarshallRequest() relies on these, and then we need to unmarshall these
1265 // two again, because they get marshalled in MarshallRequest().
1266 // All this is pretty damn ugly and should be rewritten.
1267 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1268 req->header.requestid = htons( req->header.requestid );
1269 req->header.dlen = htonl( req->header.dlen );
1270 XRootDStatus st = MarshallRequest( msg );
1271 req->header.requestid = htons( req->header.requestid );
1272 req->header.dlen = htonl( req->header.dlen );
1273 msg->SetIsMarshalled( false );
1274 return st;
1275 }
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 1358 of file XrdClXRootDTransport.cc.

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

1826 {
1827 XrdSysRWLockHelper scope( pSecUnloadHandler->lock, false ); // obtain write lock
1828 pSecUnloadHandler->unloaded = true;
1829 }

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: