XRootD
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. More...
 
 ~XRootDTransport ()
 Destructor. More...
 
virtual void DecFileInstCnt (AnyObject &channelData)
 Decrement file object instance count bound to this channel. More...
 
virtual void Disconnect (AnyObject &channelData, uint16_t subStreamId)
 The stream has been disconnected, do the cleanups. More...
 
virtual void FinalizeChannel (AnyObject &channelData)
 Finalize channel. More...
 
virtual URL GetBindPreference (const URL &url, AnyObject &channelData)
 Get bind preference for the next data stream. More...
 
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. More...
 
virtual Status GetSignature (Message *toSign, Message *&sign, XRootDChannelInfo *info)
 Get signature for given message. More...
 
virtual XRootDStatus HandShake (HandShakeData *handShakeData, AnyObject &channelData)
 HandShake. More...
 
virtual bool HandShakeDone (HandShakeData *handShakeData, AnyObject &channelData)
 
virtual void InitializeChannel (const URL &url, AnyObject &channelData)
 Initialize channel. More...
 
virtual Status IsStreamBroken (time_t inactiveTime, AnyObject &channelData)
 
virtual bool IsStreamTTLElapsed (time_t time, AnyObject &channelData)
 Check if the stream should be disconnected. More...
 
virtual uint32_t MessageReceived (Message &msg, uint16_t subStream, AnyObject &channelData)
 Check if the message invokes a stream action. More...
 
virtual void MessageSent (Message *msg, uint16_t subStream, uint32_t bytesSent, AnyObject &channelData)
 Notify the transport about a message having been sent. More...
 
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. More...
 
virtual uint16_t SubStreamNumber (AnyObject &channelData)
 Return a number of substreams per stream that should be created. More...
 
virtual void WaitBeforeExit ()
 Wait until the program can safely exit. More...
 
- 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. More...
 
static void LogErrorResponse (const Message &msg)
 Log server error response. More...
 
static XRootDStatus MarshallRequest (char *msg)
 Marshal the outgoing message. More...
 
static XRootDStatus MarshallRequest (Message *msg)
 Marshal the outgoing message. More...
 
static uint16_t NbConnectedStrm (AnyObject &channelData)
 Number of currently connected data streams. More...
 
static void SetDescription (Message *msg)
 Get the description of a message. More...
 
static XRootDStatus UnMarchalStatusMore (Message &msg)
 Unmarshall the correction-segment of the status response for pgwrite. More...
 
static XRootDStatus UnMarshallBody (Message *msg, uint16_t reqType)
 Unmarshall the body of the incoming message. More...
 
static void UnMarshallHeader (Message &msg)
 Unmarshall the header incoming message. More...
 
static XRootDStatus UnMarshallRequest (Message *msg)
 
static XRootDStatus UnMarshalStatusBody (Message &msg, uint16_t reqType)
 Unmarshall the body of the status response. More...
 

Friends

struct PluginUnloadHandler
 

Additional Inherited Members

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

Detailed Description

XRootD transport handler.

Definition at line 47 of file XrdClXRootDTransport.hh.

Constructor & Destructor Documentation

◆ XRootDTransport()

XrdCl::XRootDTransport::XRootDTransport ( )

Constructor.

Definition at line 290 of file XrdClXRootDTransport.cc.

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

◆ ~XRootDTransport()

XrdCl::XRootDTransport::~XRootDTransport ( )

Destructor.

Definition at line 298 of file XrdClXRootDTransport.cc.

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

Member Function Documentation

◆ DecFileInstCnt()

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

Decrement file object instance count bound to this channel.

Implements XrdCl::TransportHandler.

Definition at line 1824 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ Disconnect()

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

The stream has been disconnected, do the cleanups.

Implements XrdCl::TransportHandler.

Definition at line 1554 of file XrdClXRootDTransport.cc.

1556  {
1557  XRootDChannelInfo *info = 0;
1558  channelData.Get( info );
1559 
1560  if (!info) {
1561  DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1562  return;
1563  }
1564 
1565  XrdSysMutexHelper scopedLock( info->mutex );
1566 
1567  CleanUpProtection( info );
1568 
1569  if( !info->stream.empty() )
1570  {
1571  XRootDStreamInfo &sInfo = info->stream[subStreamId];
1572  sInfo.status = XRootDStreamInfo::Disconnected;
1573  }
1574 
1575  if( subStreamId == 0 )
1576  {
1577  info->sidManager->ReleaseAllTimedOut();
1578  info->sentOpens.clear();
1579  info->sentCloses.clear();
1580  info->openFiles = 0;
1581  info->waitBarrier = 0;
1582  }
1583  }
static Log * GetLog()
Get default log.
void Error(uint64_t topic, const char *format,...)
Report an error.
Definition: XrdClLog.cc:231
const uint64_t XRootDTransportMsg

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

+ Here is the call graph for this function:

◆ FinalizeChannel()

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

Finalize channel.

Implements XrdCl::TransportHandler.

Definition at line 470 of file XrdClXRootDTransport.cc.

471  {
472  }

◆ GenerateDescription()

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

Get the description of a message.

Definition at line 3003 of file XrdClXRootDTransport.cc.

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

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

◆ GetBindPreference()

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

Get bind preference for the next data stream.

Implements XrdCl::TransportHandler.

Definition at line 1923 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ GetBody()

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

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

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

Implements XrdCl::TransportHandler.

Definition at line 346 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ GetHeader()

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

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

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

Implements XrdCl::TransportHandler.

Definition at line 306 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ GetMore()

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

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

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

Implements XrdCl::TransportHandler.

Definition at line 382 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ GetSignature() [1/2]

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

Get signature for given message.

Implements XrdCl::TransportHandler.

Definition at line 1784 of file XrdClXRootDTransport.cc.

1785  {
1786  XRootDChannelInfo *info = 0;
1787  channelData.Get( info );
1788  return GetSignature( toSign, sign, info );
1789  }
virtual Status GetSignature(Message *toSign, Message *&sign, AnyObject &channelData)
Get signature for given message.

References XrdCl::AnyObject::Get().

+ Here is the call graph for this function:

◆ GetSignature() [2/2]

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

Get signature for given message.

Definition at line 1794 of file XrdClXRootDTransport.cc.

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

References XrdCl::errInternal, XrdCl::errInvalidOp, XrdCl::Buffer::GetBuffer(), XrdCl::Buffer::Grab(), XrdCl::PluginUnloadHandler::lock, NEED2SECURE, XrdCl::XRootDChannelInfo::protection, XrdSecProtect::Secure(), XrdCl::stError, and XrdCl::PluginUnloadHandler::unloaded.

+ Here is the call graph for this function:

◆ HandShake()

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

HandShake.

Implements XrdCl::TransportHandler.

Definition at line 477 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ HandShakeDone()

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

Implements XrdCl::TransportHandler.

Definition at line 756 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ InitializeChannel()

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

Initialize channel.

Implements XrdCl::TransportHandler.

Definition at line 449 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ IsStreamBroken()

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

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

Implements XrdCl::TransportHandler.

Definition at line 829 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ IsStreamTTLElapsed()

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

Check if the stream should be disconnected.

Implements XrdCl::TransportHandler.

Definition at line 776 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ LogErrorResponse()

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

Log server error response.

Definition at line 1517 of file XrdClXRootDTransport.cc.

1518  {
1519  Log *log = DefaultEnv::GetLog();
1520  ServerResponse *rsp = (ServerResponse *)msg.GetBuffer();
1521  char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
1522  memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
1523  log->Error( XRootDTransportMsg, "Server responded with an error [%d]: %s",
1524  rsp->body.error.errnum, errmsg );
1525  delete [] errmsg;
1526  }
union ServerResponse::@0 body
ServerResponseHeader hdr
Definition: XProtocol.hh:1288

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

+ Here is the call graph for this function:

◆ MarshallRequest() [1/2]

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

Marshal the outgoing message.

Definition at line 1113 of file XrdClXRootDTransport.cc.

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

References ClientMvRequest::arg1len, ClientRequest::chkpoint, ClientRequest::chmod, ClientProtocolRequest::clientpv, ClientRequestHdr::dlen, ClientReadVRequest::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, readahead_list::offset, ClientTruncateRequest::offset, ClientWriteRequest::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(), and XrdCl::Message::SetIsMarshalled().

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

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

◆ MessageReceived()

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

Check if the message invokes a stream action.

Implements XrdCl::TransportHandler.

Definition at line 1640 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ MessageSent()

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

Notify the transport about a message having been sent.

Implements XrdCl::TransportHandler.

Definition at line 1749 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ Multiplex()

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

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

Implements XrdCl::TransportHandler.

Definition at line 872 of file XrdClXRootDTransport.cc.

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

◆ MultiplexSubStream()

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

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

Implements XrdCl::TransportHandler.

Definition at line 880 of file XrdClXRootDTransport.cc.

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

+ Here is the call graph for this function:

◆ NbConnectedStrm()

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

Number of currently connected data streams.

Definition at line 1531 of file XrdClXRootDTransport.cc.

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

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

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

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

◆ NeedControlConnection()

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

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

Definition at line 167 of file XrdClXRootDTransport.hh.

168  {
169  return true;
170  }

◆ NeedEncryption()

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

Implements XrdCl::TransportHandler.

Definition at line 1844 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ Query()

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

Query the channel.

Implements XrdCl::TransportHandler.

Definition at line 1588 of file XrdClXRootDTransport.cc.

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

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

◆ SubStreamNumber()

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

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

Implements XrdCl::TransportHandler.

Definition at line 1052 of file XrdClXRootDTransport.cc.

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

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

+ Here is the call graph for this function:

◆ UnMarchalStatusMore()

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

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

Definition at line 1444 of file XrdClXRootDTransport.cc.

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

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

Referenced by GetMore().

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

◆ UnMarshallBody()

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

Unmarshall the body of the incoming message.

Definition at line 1290 of file XrdClXRootDTransport.cc.

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

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

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

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

◆ UnMarshallHeader()

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

Unmarshall the header incoming message.

Definition at line 1507 of file XrdClXRootDTransport.cc.

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

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

Referenced by GetHeader().

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

◆ UnMarshallRequest()

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

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

Definition at line 1269 of file XrdClXRootDTransport.cc.

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

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

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

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

◆ UnMarshalStatusBody()

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

Unmarshall the body of the status response.

Definition at line 1368 of file XrdClXRootDTransport.cc.

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

References ServerResponseStatus::bdy, XrdOucCRC::Calc32C(), ServerResponseBody_Status::crc32c, ServerResponseHeader::dlen, ServerResponseBody_Status::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, ServerResponseHeader::streamid, and ServerResponseBody_Status::streamID.

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

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

◆ WaitBeforeExit()

void XrdCl::XRootDTransport::WaitBeforeExit ( )
virtual

Wait until the program can safely exit.

Implements XrdCl::TransportHandler.

Definition at line 1835 of file XrdClXRootDTransport.cc.

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

References XrdCl::PluginUnloadHandler::lock, and XrdCl::PluginUnloadHandler::unloaded.

Friends And Related Function Documentation

◆ PluginUnloadHandler

friend struct PluginUnloadHandler
friend

Definition at line 432 of file XrdClXRootDTransport.hh.


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