51#include "XrdVersion.hh"
88 std::pair< std::set<std::string>::iterator,
bool > ret =
protocols.insert( protocol );
147 strmqueues.resize( size - 1, 0 );
155 strmqueues.resize( size - 1, 0);
163 uint16_t
Select(
const std::vector<bool> &connected )
166 size_t minval = std::numeric_limits<size_t>::max();
168 for(
size_t i = 0; i < connected.size() && i < strmqueues.size(); ++i )
170 if( !connected[i] )
continue;
172 if( strmqueues[i] < minval )
175 minval = strmqueues[i];
189 --strmqueues[substrm - 1];
194 std::vector<size_t> strmqueues;
200 bindprefs( std::move( bindprefs ) ), next( 0 )
204 inline const std::string&
Get()
206 std::string &ret = bindprefs[next];
208 if( next >= bindprefs.size() )
214 std::vector<std::string> bindprefs;
302 delete pSecUnloadHandler; pSecUnloadHandler = 0;
321 size_t leftToBeRead = 8 - message.
GetCursor();
322 while( leftToBeRead )
326 leftToBeRead, bytesRead );
330 leftToBeRead -= bytesRead;
335 uint32_t bodySize = *(uint32_t*)(message.
GetBuffer(4));
338 "body", (
void*)&message, bodySize );
353 size_t leftToBeRead = 0;
354 uint32_t bodySize = 0;
356 bodySize = rsphdr->
dlen;
358 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
360 "Response body too large." );
362 if( message.
GetSize() < bodySize + 8 )
365 leftToBeRead = bodySize-(message.
GetCursor()-8);
366 while( leftToBeRead )
374 leftToBeRead -= bytesRead;
397 uint32_t bodySize = rsphdr->
dlen;
398 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
400 "kXR_status: response body too large." );
403 "kXR_status: invalid message size." );
406 uint32_t moreSize =
static_cast<uint32_t
>( rspst->
bdy.
dlen );
407 if( moreSize > std::numeric_limits<uint32_t>::max() - 8 - bodySize )
409 "kXR_status: response body too large." );
410 bodySize += moreSize;
412 if( message.
GetSize() < bodySize + 8 )
415 size_t leftToBeRead = bodySize-(message.
GetCursor()-8);
416 while( leftToBeRead )
424 leftToBeRead -= bytesRead;
456 channelData.
Set( info );
460 env->
GetInt(
"SubStreamsPerChannel", streams );
461 if( streams < 1 ) streams = 1;
462 info->
stream.resize( streams );
483 channelData.
Get( info );
494 "[%s] Internal error: not enough substreams",
502 return HandShakeMain( handShakeData, channelData );
504 return HandShakeParallel( handShakeData, channelData );
514 channelData.
Get( info );
518 "[%s] Internal error: no channel info",
531 handShakeData->
out = GenerateInitialHSProtocol( handShakeData, info,
542 XRootDStatus st = ProcessServerHS( handShakeData, info );
556 XRootDStatus st = ProcessProtocolResp( handShakeData, info );
566 handShakeData->
out = GenerateProtocol( handShakeData, info,
572 handShakeData->
out = GenerateLogIn( handShakeData, info );
583 XRootDStatus st = ProcessLogInResp( handShakeData, info );
591 if( st.IsOK() && st.code ==
suDone )
600 handShakeData->
out = GenerateEndSession( handShakeData, info );
610 st = DoAuthentication( handShakeData, info );
623 XRootDStatus st = DoAuthentication( handShakeData, info );
631 if( st.IsOK() && st.code ==
suDone )
638 handShakeData->
out = GenerateEndSession( handShakeData, info );
656 XRootDStatus st = ProcessEndSessionResp( handShakeData, info );
658 if( st.IsOK() && st.code ==
suDone )
662 else if( !st.IsOK() )
679 XRootDChannelInfo *info = 0;
680 channelData.Get( info );
684 "[%s] Internal error: no channel info",
685 handShakeData->streamName.c_str());
689 XRootDStreamInfo &sInfo = info->
stream[handShakeData->subStreamId];
697 handShakeData->out = GenerateInitialHSProtocol( handShakeData, info,
709 XRootDStatus st = ProcessServerHS( handShakeData, info );
723 XRootDStatus st = ProcessProtocolResp( handShakeData, info );
731 handShakeData->out = GenerateBind( handShakeData, info );
741 XRootDStatus st = ProcessBindResp( handShakeData, info );
762 channelData.
Get( info );
766 "[%s] Internal error: no channel info",
782 channelData.
Get( info );
789 "Internal error: no channel info, behaving as if TTL has elapsed");
800 env->
GetInt(
"DataServerTTL", ttl );
805 env->
GetInt(
"LoadBalancerTTL", ttl );
812 uint16_t allocatedSIDs = info->
sidManager->GetNumberOfAllocatedSIDs();
814 "TTL: %d, allocated SIDs: %d, open files: %d, bound file objects: %d",
815 info->
streamName.c_str(), (
long long) inactiveTime, ttl, allocatedSIDs,
818 if( info->
openFiles != 0 && info->
finstcnt.load( std::memory_order_relaxed ) != 0 )
821 if( !allocatedSIDs && inactiveTime > ttl )
835 channelData.
Get( info );
841 "Internal error: no channel info, behaving as if stream is broken");
846 env->
GetInt(
"StreamTimeout", streamTimeout );
850 const time_t now = time(0);
852 info->
sidManager->IsAnySIDOldAs( now - streamTimeout );
855 "stream timeout: %d, any SID: %d, wait barrier: %s",
856 info->
streamName.c_str(), (
long long) inactiveTime, streamTimeout,
859 if( inactiveTime < streamTimeout )
862 if( now < info->waitBarrier )
887 channelData.
Get( info );
891 "Internal error: no channel info, cannot multiplex");
908 uint16_t upStream = 0;
909 uint16_t downStream = 0;
914 downStream = hint->
down;
919 std::vector<bool> connected;
920 connected.reserve( info->
stream.size() - 1 );
921 size_t nbConnected = 0;
922 for(
size_t i = 1; i < info->
stream.size(); ++i )
925 connected.push_back(
true );
929 connected.push_back(
false );
931 if( nbConnected == 0 )
937 if( upStream >= info->
stream.size() )
940 "[%s] Up link stream %d does not exist, using 0",
945 if( downStream >= info->
stream.size() )
948 "[%s] Down link stream %d does not exist, using 0",
972 memset( newBuf, 0, 8 );
1046 return PathID( upStream, downStream );
1057 channelData.
Get( info );
1076 uint16_t ret = info->
stream.size();
1080 env->
GetInt(
"TlsNoData", nodata );
1093 if( ( usrTlsStrm0 && usrNoTlsData && srvNoTlsData ) ||
1094 ( srvTlsStrm0 && srvNoTlsData && usrNoTlsData ) )
1100 if( ret == 1 ) ++ret;
1103 if( ret > info->
stream.size() )
1105 info->
stream.resize( ret );
1205 uint16_t numChunks = (req->
readv.
dlen)/16;
1207 for(
size_t i = 0; i < numChunks; ++i )
1209 dataChunk[i].
rlen = htonl( dataChunk[i].rlen );
1210 dataChunk[i].
offset = htonll( dataChunk[i].offset );
1220 for(
size_t i = 0; i < numChunks; ++i )
1222 dataChunk[i].
srcOffs = htonll( dataChunk[i].srcOffs );
1223 dataChunk[i].
srcLen = htonll( dataChunk[i].srcLen );
1224 dataChunk[i].
dstOffs = htonll( dataChunk[i].dstOffs );
1237 for(
size_t i = 0; i < numChunks; ++i )
1239 wrtList[i].
wlen = htonl( wrtList[i].wlen );
1240 wrtList[i].
offset = htonll( wrtList[i].offset );
1324 m->
body.protocol.pval = ntohl( m->
body.protocol.pval );
1325 m->
body.protocol.flags = ntohl( m->
body.protocol.flags );
1336 m->
body.error.errnum = ntohl( m->
body.error.errnum );
1346 m->
body.wait.seconds = htonl( m->
body.wait.seconds );
1356 m->
body.redirect.port = htonl( m->
body.redirect.port );
1366 m->
body.waitresp.seconds = htonl( m->
body.waitresp.seconds );
1376 m->
body.attn.actnum = htonl( m->
body.attn.actnum );
1412 "kXR_status: invalid message size." );
1440 "corrupted (crc32c integrity check failed)." );
1447 "(stream ID mismatch)." );
1455 "(request ID mismatch)." );
1479 "kXR_status: invalid message size." );
1493 if( crcval != cse->
cseCRC )
1496 "corrupted (crc32c integrity check failed)." );
1507 for(
size_t i = 0; i < pgcnt; ++i )
1508 pgoffs[i] = ntohll( pgoffs[i] );
1528 header->
dlen = ntohl( header->
dlen );
1538 char *errmsg =
new char[rsp->
hdr.
dlen-3]; errmsg[rsp->
hdr.
dlen-4] = 0;
1539 memcpy( errmsg, rsp->
body.error.errmsg, rsp->
hdr.
dlen-4 );
1541 rsp->
body.error.errnum, errmsg );
1551 channelData.
Get( info );
1560 uint16_t nbConnected = 0;
1561 for(
size_t i = 1; i < info->
stream.size(); ++i )
1572 uint16_t subStreamId )
1575 channelData.
Get( info );
1584 if( !info->
stream.empty() )
1590 if( subStreamId == 0 )
1592 CleanUpProtection( info );
1609 channelData.
Get( info );
1622 result.
Set( (
const char*)
"XRootD",
false );
1661 channelData.
Get( info );
1684 "response that we're no longer interested in (timed out)",
1696 uint16_t sid; memcpy( &sid, rsp->
hdr.
streamid, 2 );
1697 std::set<uint16_t>::iterator sidIt = info->
sentOpens.find( sid );
1709 uint32_t seconds = 0;
1711 seconds = ntohl( rsp->
body.wait.seconds ) + 5;
1715 seconds = ntohl( rsp->
body.waitresp.seconds );
1717 log->
Dump(
XRootDMsg,
"[%s] Got kXR_waitresp response of %u seconds, "
1718 "setting up wait barrier.",
1723 time_t barrier = time(0) + seconds;
1731 uint16_t sid; memcpy( &sid, rsp->
hdr.
streamid, 2 );
1732 std::set<uint16_t>::iterator sidIt = info->
sentOpens.find( sid );
1741 info->
finstcnt.fetch_add( 1, std::memory_order_relaxed );
1777 channelData.
Get( info );
1803 channelData.
Get( info );
1831 sign->
Grab(
reinterpret_cast<char*
>( newreq ), rc );
1843 channelData.
Get( info );
1844 if( info->
finstcnt.load( std::memory_order_relaxed ) > 0 )
1845 info->
finstcnt.fetch_sub( 1, std::memory_order_relaxed );
1854 pSecUnloadHandler->unloaded =
true;
1864 channelData.
Get( info );
1868 env->
GetInt(
"NoTlsOK", notlsok );
1942 channelData.
Get( info );
1960 "[%s] Sending out the initial hand shake + kXR_protocol",
1970 init->
fifth = htonl(2012);
1973 InitProtocolReq( proto, info, expect );
1981 Message *XRootDTransport::GenerateProtocol( HandShakeData *hsData,
1982 XRootDChannelInfo *info,
1987 "[%s] Sending out the kXR_protocol",
1988 hsData->streamName.c_str() );
1995 InitProtocolReq( proto, info, expect );
2003 void XRootDTransport::InitProtocolReq( ClientProtocolRequest *request,
2017 env->
GetInt(
"NoTlsOK", notlsok );
2020 env->
GetInt(
"TlsNoData", tlsnodata );
2022 if (info->encrypted ||
InitTLS())
2025 if (info->encrypted && !(notlsok || tlsnodata))
2028 request->
expect = expect;
2046 Message *msg = hsData->in;
2047 ServerResponseHeader *respHdr = (ServerResponseHeader *)msg->GetBuffer();
2048 ServerInitHandShake *hs = (ServerInitHandShake *)msg->GetBuffer(4);
2053 hsData->streamName.c_str() );
2058 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2059 const uint32_t pv = ntohl(hs->
protover);
2064 if( hsData->subStreamId == 0 )
2066 info->protocolVersion = pv;
2067 info->serverFlags = sInfo.serverFlags;
2071 "[%s] Got the server hand shake response (%s, protocol "
2073 hsData->streamName.c_str(),
2074 ServerFlagsToStr( sInfo.serverFlags ).c_str(),
2075 info->protocolVersion );
2092 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2098 hsData->streamName.c_str() );
2103 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2104 if( rsp->
body.protocol.pval >= 0x297 )
2105 sInfo.serverFlags = rsp->
body.protocol.flags;
2107 if( hsData->subStreamId > 0 )
2110 info->serverFlags = sInfo.serverFlags;
2114 env->
GetInt(
"NoTlsOK", notlsok );
2129 "[%s] Falling back to unencrypted transmission, server does "
2130 "not support TLS encryption.",
2131 hsData->streamName.c_str() );
2132 info->encrypted =
false;
2135 if( rsp->
body.protocol.pval >= 0x297 )
2136 info->serverFlags = rsp->
body.protocol.flags;
2140 info->protRespBuff.assign(
sizeof( ServerResponseBody_Protocol ), 0 );
2141 info->protRespSize = 0;
2142 ServerResponseBody_Protocol *protRespBody =
2143 reinterpret_cast<ServerResponseBody_Protocol*
>( info->protRespBuff.data() );
2144 protRespBody->
flags = rsp->
body.protocol.flags;
2145 protRespBody->
pval = rsp->
body.protocol.pval;
2147 char* bodybuff =
reinterpret_cast<char*
>( &rsp->
body.protocol.secreq );
2148 size_t bodysize = rsp->
hdr.
dlen - 8;
2149 XRootDStatus st = ProcessProtocolBody( bodybuff, bodysize, info );
2155 "[%s] kXR_protocol successful (%s, protocol version %x)",
2156 hsData->streamName.c_str(),
2157 ServerFlagsToStr( info->serverFlags ).c_str(),
2158 info->protocolVersion );
2160 if( !( info->serverFlags &
kXR_haveTLS ) && info->encrypted )
2167 "Server was not configured to support encryption." );
2175 env->
GetInt(
"WantTlsOnNoPgrw", tlsOnNoPgrw );
2176 if( !( info->serverFlags &
kXR_suppgrw ) && tlsOnNoPgrw )
2182 if( info->encrypted )
2185 "[%s] Server does not support PgRead/PgWrite and"
2186 " WantTlsOnNoPgrw is on; enforcing encryption for data.",
2187 hsData->streamName.c_str() );
2197 info->encrypted =
true;
2205 XRootDStatus XRootDTransport::ProcessProtocolBody(
char *bodybuff,
2218 if( bodysize < bifreq->bifILen )
2220 "protocol response." );
2221 std::string bindprefs_str( bodybuff, bifreq->
bifILen );
2222 std::vector<std::string> bindprefs;
2224 info->bindSelector.reset(
new BindPrefSelector( std::move( bindprefs ) ) );
2233 sizeof( ServerResponseSVec_Protocol );
2234 if( bodysize >= secHdrLen && secreq->
theTag ==
'S' )
2240 size_t secsize = secHdrLen + secreq->
secvsz *
2241 sizeof( ServerResponseSVec_Protocol );
2242 if( bodysize < secsize )
2244 "protocol response." );
2247 if( info->protRespBuff.size() < respsize )
2248 info->protRespBuff.resize( respsize, 0 );
2250 info->protRespSize = respsize;
2265 "[%s] Sending out the bind request",
2266 hsData->streamName.c_str() );
2269 Message *msg =
new Message(
sizeof( ClientBindRequest ) );
2270 ClientBindRequest *bindReq = (ClientBindRequest *)msg->GetBuffer();
2273 memcpy( bindReq->
sessid, info->sessionId, 16 );
2291 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2296 hsData->streamName.c_str() );
2300 info->stream[hsData->subStreamId].pathId = rsp->
body.bind.substreamid;
2302 hsData->streamName.c_str() );
2322 char *cgiBuffer =
new char[1024 + info->logintoken.size()];
2323 std::string appName;
2324 std::string monInfo;
2325 env->GetString(
"AppName", appName );
2326 env->GetString(
"MonInfo", monInfo );
2327 if( info->logintoken.empty() )
2329 snprintf( cgiBuffer, 1024,
2330 "xrd.cc=%s&xrd.tz=%d&xrd.appname=%s&xrd.info=%s&"
2331 "xrd.hostname=%s&xrd.rn=%s", countryCode.c_str(), timeZone,
2332 appName.c_str(), monInfo.c_str(), hostName, XrdVERSION );
2336 snprintf( cgiBuffer, 1024,
2337 "xrd.cc=%s&xrd.tz=%d&xrd.appname=%s&xrd.info=%s&"
2338 "xrd.hostname=%s&xrd.rn=%s&%s", countryCode.c_str(), timeZone,
2339 appName.c_str(), monInfo.c_str(), hostName, XrdVERSION, info->logintoken.c_str() );
2341 uint16_t cgiLen = strlen( cgiBuffer );
2347 Message *msg =
new Message(
sizeof(ClientLoginRequest) + cgiLen );
2348 ClientLoginRequest *loginReq = (ClientLoginRequest *)msg->GetBuffer();
2351 loginReq->
pid = ::getpid();
2353 loginReq->
dlen = cgiLen;
2359 int multiProtocol = 0;
2360 env->GetInt(
"MultiProtocol", multiProtocol );
2368 bool dualStack =
false;
2369 bool privateIPv6 =
false;
2370 bool privateIPv4 =
false;
2394 if( !dualStack && hsData->serverAddr )
2407 std::string buffer( 8, 0 );
2408 if( hsData->url->GetUserName().length() )
2409 buffer = hsData->url->GetUserName();
2412 char *name =
new char[1024];
2419 buffer.resize( 8, 0 );
2420 std::copy( buffer.begin(), buffer.end(), (
char*)loginReq->
username );
2422 msg->Append( cgiBuffer, cgiLen, 24 );
2425 "username: %s, cgi: %s, dual-stack: %s, private IPv4: %s, "
2426 "private IPv6: %s", hsData->streamName.c_str(),
2427 loginReq->
username, cgiBuffer, dualStack ?
"true" :
"false",
2428 privateIPv4 ?
"true" :
"false",
2429 privateIPv6 ?
"true" :
"false" );
2431 delete [] cgiBuffer;
2448 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2453 hsData->streamName.c_str() );
2457 if( !info->firstLogIn )
2458 memcpy( info->oldSessionId, info->sessionId, 16 );
2460 if( rsp->
hdr.
dlen == 0 && info->protocolVersion <= 0x289 )
2467 memset( info->sessionId, 0, 16 );
2469 "[%s] Logged in, accepting empty login response.",
2470 hsData->streamName.c_str() );
2477 memcpy( info->sessionId, rsp->
body.login.sessid, 16 );
2482 hsData->streamName.c_str(), sessId.c_str() );
2489 size_t len = rsp->
hdr.
dlen-16;
2490 info->authBuffer =
new char[len+1];
2491 info->authBuffer[len] = 0;
2492 memcpy( info->authBuffer, rsp->
body.login.sec, len );
2494 hsData->streamName.c_str(), info->authBuffer );
2512 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2514 std::string protocolName;
2522 hsData->streamName.c_str() );
2528 info->authEnv->Put(
"sockname", hsData->clientName.c_str() );
2529 info->authEnv->Put(
"username", hsData->url->GetUserName().c_str() );
2530 info->authEnv->Put(
"password", hsData->url->GetPassword().c_str() );
2533 URL::ParamsMap::const_iterator it;
2534 for( it = urlParams.begin(); it != urlParams.end(); ++it )
2536 if( it->first.compare( 0, 4,
"xrd." ) == 0 ||
2537 it->first.compare( 0, 6,
"xrdcl." ) == 0 )
2538 info->authEnv->Put( it->first.c_str(), it->second.c_str() );
2544 size_t authBuffLen = strlen( info->authBuffer );
2545 char *pars = (
char *)malloc( authBuffLen + 1 );
2546 memcpy( pars, info->authBuffer, authBuffLen );
2549 delete [] info->authBuffer;
2550 info->authBuffer = 0;
2555 XRootDStatus st = GetCredentials( credentials, hsData, info );
2558 CleanUpAuthentication( info );
2561 protocolName = info->authProtocol->Entity.prot;
2569 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2570 protocolName = info->authProtocol->Entity.prot;
2578 "[%s] Sending more authentication data for %s",
2579 hsData->streamName.c_str(), protocolName.c_str() );
2582 char *secTokenData = (
char*)malloc( len );
2583 memcpy( secTokenData, rsp->
body.authmore.data, len );
2585 XrdOucErrInfo ei(
"", info->authEnv);
2586 credentials = info->authProtocol->getCredentials( secToken, &ei );
2595 "[%s] Auth protocol handler for %s refuses to give "
2596 "us more credentials %s",
2597 hsData->streamName.c_str(), protocolName.c_str(),
2599 CleanUpAuthentication( info );
2609 info->authProtocolName = info->authProtocol->Entity.prot;
2614 if( !info->protRespBuff.empty() )
2616 ServerResponseBody_Protocol *protRespBody =
2617 reinterpret_cast<ServerResponseBody_Protocol*
>( info->protRespBuff.data() );
2618 int rc =
XrdSecGetProtection( info->protection, *info->authProtocol, *protRespBody, info->protRespSize );
2622 "[%s] XrdSecProtect loaded.", hsData->streamName.c_str() );
2627 "[%s] XrdSecProtect: no protection needed.",
2628 hsData->streamName.c_str() );
2633 "[%s] Failed to load XrdSecProtect: %s",
2634 hsData->streamName.c_str(),
XrdSysE2T( -rc ) );
2635 CleanUpAuthentication( info );
2641 if( !info->protection )
2642 CleanUpAuthentication( info );
2644 pSecUnloadHandler->Register( info->authProtocolName );
2647 "[%s] Authenticated with %s.", hsData->streamName.c_str(),
2648 protocolName.c_str() );
2663 char *errmsg =
new char[rsp->
hdr.
dlen-3]; errmsg[rsp->
hdr.
dlen-4] = 0;
2664 memcpy( errmsg, rsp->
body.error.errmsg, rsp->
hdr.
dlen-4 );
2666 "[%s] Authentication with %s failed: %s",
2667 hsData->streamName.c_str(), protocolName.c_str(),
2671 info->authProtocol->Delete();
2672 info->authProtocol = 0;
2677 XRootDStatus st = GetCredentials( credentials, hsData, info );
2680 CleanUpAuthentication( info );
2683 protocolName = info->authProtocol->Entity.prot;
2690 info->authProtocolName = info->authProtocol->Entity.prot;
2691 CleanUpAuthentication( info );
2694 "[%s] Authentication with %s failed: unexpected answer",
2695 hsData->streamName.c_str(), protocolName.c_str() );
2703 Message *msg =
new Message(
sizeof(ClientAuthRequest)+credentials->
size );
2705 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
2706 char *reqBuffer = msg->GetBuffer(
sizeof(ClientAuthRequest));
2711 protocolName.length() > 4 ? 4 : protocolName.length() );
2713 memcpy( reqBuffer, credentials->
buffer, credentials->
size );
2738 XrdOucErrInfo ei(
"", info->authEnv);
2748 char *secuidc = (ei.getEnv()) ? ei.getEnv()->Get(
"xrdcl.secuid") : 0;
2749 char *secgidc = (ei.getEnv()) ? ei.getEnv()->Get(
"xrdcl.secgid") : 0;
2754 if(secuidc) secuid = atoi(secuidc);
2755 if(secgidc) secgid = atoi(secgidc);
2758 ScopedFsUidSetter uidSetter(secuid, secgid, hsData->streamName);
2759 if(!uidSetter.IsOk()) {
2760 log->Error(
XRootDTransportMsg,
"[%s] Error while setting (fsuid, fsgid) to (%d, %d)",
2761 hsData->streamName.c_str(), secuid, secgid );
2765 if(secuid >= 0 || secgid >= 0) {
2766 log->Error(
XRootDTransportMsg,
"[%s] xrdcl.secuid and xrdcl.secgid only supported on Linux.",
2767 hsData->streamName.c_str() );
2769 " only supported on Linux" );
2777 XrdNetAddr &srvAddrInfo = *
const_cast<XrdNetAddr *
>(hsData->serverAddr);
2778 srvAddrInfo.
SetTLS( info->encrypted );
2784 info->authProtocol = (*authHandler)( hsData->url->GetHostName().c_str(),
2788 if( !info->authProtocol )
2791 hsData->streamName.c_str() );
2795 std::string protocolName = info->authProtocol->Entity.prot;
2797 hsData->streamName.c_str(), protocolName.c_str() );
2802 credentials = info->authProtocol->getCredentials( 0, &ei );
2806 "[%s] Cannot get credentials for protocol %s: %s",
2807 hsData->streamName.c_str(), protocolName.c_str(),
2809 info->authProtocol->Delete();
2821 if( info->authProtocol )
2822 info->authProtocol->Delete();
2823 delete info->authParams;
2824 delete info->authEnv;
2825 info->authProtocol = 0;
2826 info->authParams = 0;
2837 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
2840 if( info->protection )
2842 info->protection->Delete();
2843 info->protection = 0;
2845 CleanUpAuthentication( info );
2848 info->protRespBuff.clear();
2849 info->protRespSize = 0;
2860 char errorBuff[1024];
2865 auto ret = authHandler.load( std::memory_order_relaxed );
2866 if( ret )
return ret;
2872 static XrdSysMutex mtx;
2873 XrdSysMutexHelper lck( mtx );
2875 ret = authHandler.load( std::memory_order_relaxed );
2876 if( ret )
return ret;
2880 authHandler.store( ret, std::memory_order_relaxed );
2885 "Unable to get the security framework: %s", errorBuff );
2902 Message *msg =
new Message(
sizeof(ClientEndsessRequest) );
2903 ClientEndsessRequest *endsessReq = (ClientEndsessRequest *)msg->GetBuffer();
2906 memcpy( endsessReq->
sessid, info->oldSessionId, 16 );
2910 " %s", hsData->streamName.c_str(), sessId.c_str() );
2921 uint32_t size = sign->GetSize();
2922 sign->ReAllocate( size + msg->GetSize() );
2923 char* buffer = sign->GetBuffer( size );
2924 memcpy( buffer, msg->GetBuffer(), msg->GetSize() );
2925 msg->Grab( sign->GetBuffer(), sign->GetSize() );
2943 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2957 std::string errorMsg( rsp->
body.error.errmsg, rsp->
hdr.
dlen - 4 );
2959 "kXR_endsess: %s", hsData->streamName.c_str(),
2967 std::string msg( rsp->
body.wait.infomsg, rsp->
hdr.
dlen - 4 );
2969 "kXR_endsess: %s", hsData->streamName.c_str(),
2971 hsData->out = GenerateEndSession( hsData, info );
2982 std::string XRootDTransport::ServerFlagsToStr( uint32_t flags )
2984 std::string repr =
"type: ";
3008 repr.erase( repr.length()-1, 1 );
3019 char *GetDataAsString(
char *msg )
3022 char *fn =
new char[req->
dlen+1];
3023 memcpy( fn, msg + 24, req->
dlen );
3050 char *fn = GetDataAsString( msg );
3051 o <<
"file: " << fn <<
", ";
3053 o <<
"mode: 0" << std::setbase(8) << sreq->
mode <<
", ";
3054 o << std::setbase(10);
3061 o <<
"kXR_compress ";
3073 o <<
"kXR_open_apnd ";
3075 o <<
"kXR_open_read ";
3077 o <<
"kXR_open_updt ";
3079 o <<
"kXR_open_wrto ";
3083 o <<
"kXR_prefname ";
3085 o <<
"kXR_refresh ";
3087 o <<
"kXR_4dirlist ";
3089 o <<
"kXR_replica ";
3095 o <<
"kXR_retstat ";
3107 o <<
"fhtemplt: " << FileHandleToStr( sreq->
fhtemplt );
3119 o <<
"kXR_clone ( ";
3120 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3121 o << std::setbase(10);
3125 o <<
"(src_handle: ";
3126 o << FileHandleToStr( dataChunk[i].srcFH );
3128 o << std::setbase(10);
3129 o <<
"src_offset: " << dataChunk[i].
srcOffs;
3130 o <<
", src_length: " << dataChunk[i].
srcLen;
3131 o <<
", dst_offset: " << dataChunk[i].
dstOffs <<
"); ";
3145 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3159 char *fn = GetDataAsString( msg );;
3160 o <<
"path: " << fn <<
", ";
3165 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3187 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3188 o << std::setbase(10);
3190 o <<
"offset: " << sreq->
offset <<
", ";
3191 o <<
"size: " << sreq->
rlen <<
")";
3201 o <<
"kXR_pgread (";
3202 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3203 o << std::setbase(10);
3205 o <<
"offset: " << sreq->
offset <<
", ";
3206 o <<
"size: " << sreq->
rlen <<
")";
3217 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3218 o << std::setbase(10);
3220 o <<
"offset: " << sreq->
offset <<
", ";
3221 o <<
"size: " << sreq->
dlen <<
")";
3231 o <<
"kXR_pgwrite (";
3232 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3233 o << std::setbase(10);
3235 o <<
"offset: " << sreq->
offset <<
", ";
3236 o <<
"size: " << sreq->
dlen <<
")";
3263 o <<
" unknown subcode: " << sreq->
subcode;
3266 o <<
" (handle: " << FileHandleToStr( sreq->
fhandle );
3267 o << std::setbase(10);
3269 o <<
", numattr: " << nattr;
3277 o <<
", total size: " << req->
dlen <<
")";
3288 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3299 o <<
"kXR_truncate (";
3301 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3304 char *fn = GetDataAsString( msg );
3305 o <<
"file: " << fn;
3308 o << std::setbase(10);
3310 o <<
"offset: " << sreq->
offset;
3320 unsigned char *fhandle = 0;
3325 fhandle = dataChunk[0].
fhandle;
3327 o << FileHandleToStr( fhandle );
3331 o << std::setbase(10);
3336 size += dataChunk[i].
rlen;
3337 o <<
"(offset: " << dataChunk[i].
offset;
3338 o <<
", size: " << dataChunk[i].
rlen <<
"); ";
3341 o <<
"total size: " << size <<
")";
3350 unsigned char *fhandle = 0;
3351 o <<
"kXR_writev (";
3356 uint32_t numChunks = 0;
3360 size += wrtList[i].
wlen;
3365 o << FileHandleToStr( fhandle );
3369 o << std::setbase(10);
3370 o <<
"chunks: " << numChunks <<
", ";
3371 o <<
"total size: " << size <<
")";
3381 char *fn = GetDataAsString( msg );;
3382 o <<
"kXR_locate (";
3383 o <<
"path: " << fn <<
", ";
3391 o <<
"kXR_refresh ";
3393 o <<
"kXR_prefname ";
3399 o <<
"kXR_compress ";
3415 o <<
"destination: ";
3437 case kXR_QPrep: o <<
"kXR_QPrep";
break;
3440 case kXR_Qvisa: o <<
"kXR_Qvisa";
break;
3442 default: o << sreq->
infotype;
break;
3448 o <<
"handle: " << FileHandleToStr( sreq->
fhandle );
3452 o <<
"arg length: " << sreq->
dlen <<
")";
3462 char *fn = GetDataAsString( msg );;
3463 o <<
"path: " << fn <<
")";
3475 char *fn = GetDataAsString( msg );
3476 o <<
"path: " << fn <<
", ";
3478 o <<
"mode: 0" << std::setbase(8) << sreq->
mode <<
", ";
3479 o << std::setbase(10);
3486 o <<
"kXR_mkdirpath";
3498 char *fn = GetDataAsString( msg );
3499 o <<
"path: " << fn <<
")";
3511 char *fn = GetDataAsString( msg );
3512 o <<
"path: " << fn <<
", ";
3514 o <<
"mode: 0" << std::setbase(8) << sreq->
mode <<
")";
3533 o <<
"kXR_protocol (";
3534 o <<
"clientpv: 0x" << std::setbase(16) << sreq->
clientpv <<
")";
3543 o <<
"kXR_dirlist (";
3544 char *fn = GetDataAsString( msg );;
3545 o <<
"path: " << fn <<
")";
3556 char *fn = GetDataAsString( msg );;
3557 o <<
"data: " << fn <<
")";
3568 o <<
"kXR_prepare (";
3585 o <<
", priority: " << (int) sreq->
prty <<
", ";
3587 char *fn = GetDataAsString( msg );
3589 for( cursor = fn; *cursor; ++cursor )
3590 if( *cursor ==
'\n' ) *cursor =
' ';
3592 o <<
"paths: " << fn <<
")";
3600 o <<
"kXR_chkpoint (";
3602 if( sreq->
opcode == kXR_ckpBegin ) o <<
"kXR_ckpBegin)";
3603 else if( sreq->
opcode == kXR_ckpCommit ) o <<
"kXR_ckpCommit)";
3604 else if( sreq->
opcode == kXR_ckpQuery ) o <<
"kXR_ckpQuery)";
3605 else if( sreq->
opcode == kXR_ckpRollback ) o <<
"kXR_ckpRollback)";
3606 else if( sreq->
opcode == kXR_ckpXeq )
3608 o <<
"kXR_ckpXeq) ";
3622 o <<
"kXR_unknown (length: " << req->
dlen <<
")";
3631 std::string XRootDTransport::FileHandleToStr(
const unsigned char handle[4] )
3633 std::ostringstream o;
3635 for( uint8_t i = 0; i < 4; ++i )
3637 o << std::setbase(16) << std::setfill(
'0') << std::setw(2);
3638 o << (int)handle[i];
struct ClientTruncateRequest truncate
#define kXR_ShortProtRespLen
ServerResponseStatus status
struct ClientPgReadRequest pgread
union ServerResponse::@040373375333017131300127053271011057331004327334 body
struct ClientMkdirRequest mkdir
struct ClientAuthRequest auth
struct ClientPgWriteRequest pgwrite
struct ClientReadVRequest readv
struct ClientOpenRequest open
struct ServerResponseBody_Status bdy
struct ClientRequestHdr header
struct ClientWriteVRequest writev
struct ClientLoginRequest login
struct ClientChmodRequest chmod
struct ClientQueryRequest query
struct ClientReadRequest read
struct ClientMvRequest mv
struct ClientChkPointRequest chkpoint
struct ServerResponseHeader hdr
#define kXR_PROTOCOLVERSION
struct ClientPrepareRequest prepare
struct ClientWriteRequest write
#define kXR_PROTTLSVERSION
struct ClientProtocolRequest protocol
struct ClientLocateRequest locate
struct ClientCloneRequest clone
XrdSecProtocol *(*) XrdSecGetProt_t(const char *hostname, XrdNetAddrInfo &endPoint, XrdSecParameters §oken, XrdOucErrInfo *einfo)
Typedef to simplify the encoding of methods returning XrdSecProtocol.
XrdSecBuffer XrdSecParameters
XrdSecBuffer XrdSecCredentials
XrdSecGetProt_t XrdSecLoadSecFactory(char *eBuff, int eBlen, const char *seclib)
int XrdSecGetProtection(XrdSecProtect *&protP, XrdSecProtocol &aprot, ServerResponseBody_Protocol &resp, unsigned int resplen)
#define NEED2SECURE(protP)
This class implements the XRootD protocol security protection.
const char * XrdSysE2T(int errcode)
void Set(Type object, bool own=true)
void Get(Type &object)
Retrieve the object being held.
void AdvanceCursor(uint32_t delta)
Advance the cursor.
void Grab(char *buffer, uint32_t size)
Grab a buffer allocated outside.
char * GetBufferAtCursor()
Get the buffer pointer at the append cursor.
void ReAllocate(uint32_t size)
Reallocate the buffer to a new location of a given size.
void Allocate(uint32_t size)
Allocate the buffer.
const char * GetBuffer(uint32_t offset=0) const
Get the message buffer.
uint32_t GetCursor() const
Get append cursor.
uint32_t GetSize() const
Get the size of the message.
static TransportManager * GetTransportManager()
Get transport manager.
static Log * GetLog()
Get default log.
static Env * GetEnv()
Get default client environment.
bool PutInt(const std::string &key, int value)
bool GetInt(const std::string &key, int &value)
void Error(uint64_t topic, const char *format,...)
Report an error.
LogLevel GetLevel() const
Get the log level.
void Dump(uint64_t topic, const char *format,...)
Print a dump message.
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
The message representation used throughout the system.
void SetIsMarshalled(bool isMarshalled)
Set the marshalling status.
bool IsMarshalled() const
Check if the message is marshalled.
Message(uint32_t size=0)
Constructor.
static SIDMgrPool & Instance()
std::shared_ptr< SIDManager > GetSIDMgr(const URL &url)
virtual XRootDStatus Read(char *buffer, size_t size, int &bytesRead)
static void ClearErrorQueue()
Clear the error queue for the calling thread.
Perform the handshake and the authentication for each physical stream.
@ RequestClose
Send a close request.
virtual void WaitBeforeExit()=0
Wait before exit.
Manage transport handler objects.
TransportHandler * GetHandler(const std::string &protocol)
Get a transport handler object for a given protocol.
std::string GetChannelId() const
std::map< std::string, std::string > ParamsMap
URL()
Default constructor.
bool IsSecure() const
Does the protocol indicate encryption.
bool IsTPC() const
Is the URL used in TPC context.
std::string GetLoginToken() const
Get the login token if present in the opaque info.
static std::string TimeToString(time_t timestamp)
Convert timestamp to a string.
static std::string FQDNToCC(const std::string &fqdn)
Convert the fully qualified host name to country code.
static std::string Char2Hex(uint8_t *array, uint16_t size)
Print a char array as hex.
static void splitString(Container &result, const std::string &input, const std::string &delimiter)
Split a string.
const std::string & GetErrorMessage() const
Get error message.
XRootDStatus(uint16_t st=0, uint16_t code=0, uint32_t errN=0, const std::string &message="")
Constructor.
static uint16_t NbConnectedStrm(AnyObject &channelData)
Number of currently connected data streams.
virtual bool IsStreamTTLElapsed(time_t time, AnyObject &channelData)
Check if the stream should be disconnected.
virtual void Disconnect(AnyObject &channelData, uint16_t subStreamId)
The stream has been disconnected, do the cleanups.
XRootDTransport()
Constructor.
virtual uint32_t MessageReceived(Message &msg, uint16_t subStream, AnyObject &channelData)
Check if the message invokes a stream action.
virtual void WaitBeforeExit()
Wait until the program can safely exit.
static XRootDStatus UnMarshallBody(Message *msg, uint16_t reqType)
Unmarshall the body of the incoming message.
virtual XRootDStatus GetBody(Message &message, Socket *socket)
virtual XRootDStatus GetHeader(Message &message, Socket *socket)
~XRootDTransport()
Destructor.
virtual uint16_t SubStreamNumber(AnyObject &channelData)
Return a number of substreams per stream that should be created.
virtual void FinalizeChannel(AnyObject &channelData)
Finalize channel.
virtual bool HandShakeDone(HandShakeData *handShakeData, AnyObject &channelData)
virtual Status GetSignature(Message *toSign, Message *&sign, AnyObject &channelData)
Get signature for given message.
virtual void MessageSent(Message *msg, uint16_t subStream, uint32_t bytesSent, AnyObject &channelData)
Notify the transport about a message having been sent.
virtual XRootDStatus HandShake(HandShakeData *handShakeData, AnyObject &channelData)
HandShake.
virtual XRootDStatus GetMore(Message &message, Socket *socket)
static void GenerateDescription(char *msg, std::ostringstream &o)
Get the description of a message.
static XRootDStatus UnMarshallRequest(Message *msg)
static XRootDStatus UnMarchalStatusMore(Message &msg)
Unmarshall the correction-segment of the status response for pgwrite.
static void LogErrorResponse(const Message &msg)
Log server error response.
virtual void DecFileInstCnt(AnyObject &channelData)
Decrement file object instance count bound to this channel.
virtual PathID Multiplex(Message *msg, AnyObject &channelData, PathID *hint=0)
virtual void InitializeChannel(const URL &url, AnyObject &channelData)
Initialize channel.
virtual Status Query(uint16_t query, AnyObject &result, AnyObject &channelData)
Query the channel.
static void UnMarshallHeader(Message &msg)
Unmarshall the header incoming message.
static XRootDStatus UnMarshalStatusBody(Message &msg, uint16_t reqType)
Unmarshall the body of the status response.
friend struct PluginUnloadHandler
static XRootDStatus MarshallRequest(Message *msg)
Marshal the outgoing message.
virtual URL GetBindPreference(const URL &url, AnyObject &channelData)
Get bind preference for the next data stream.
virtual PathID MultiplexSubStream(Message *msg, AnyObject &channelData, PathID *hint=0)
virtual bool NeedEncryption(HandShakeData *handShakeData, AnyObject &channelData)
virtual Status IsStreamBroken(time_t inactiveTime, AnyObject &channelData)
static char * MyHostName(const char *eName="*unknown*", const char **eText=0)
static NetProt NetConfig(NetType netquery=qryINET, const char **eText=0)
static uint32_t Calc32C(const void *data, size_t count, uint32_t prevcs=0)
XrdOucEnv(const char *vardata=0, int vardlen=0, const XrdSecEntity *secent=0)
static int UserName(uid_t uID, char *uName, int uNsz)
virtual int Secure(SecurityRequest *&newreq, ClientRequest &thereq, const char *thedata)
const uint16_t errQueryNotSupported
const int DefaultLoadBalancerTTL
const uint64_t XRootDTransportMsg
const uint16_t errTlsError
const uint16_t stFatal
Fatal error, it's still an error.
const uint16_t stError
An error occurred that could potentially be retried.
const uint16_t errLoginFailed
const int DefaultWantTlsOnNoPgrw
const uint16_t errSocketTimeout
const uint16_t errDataError
data is corrupted
const uint16_t errInternal
Internal error.
const uint16_t stOK
Everything went OK.
const int DefaultSubStreamsPerChannel
const uint16_t errInvalidOp
const int DefaultDataServerTTL
const uint16_t errHandShakeFailed
const int DefaultStreamTimeout
const uint16_t suAlreadyDone
const uint16_t errNotSupported
const uint16_t suContinue
const int DefaultTlsNoData
const uint16_t errAuthFailed
const uint16_t errInvalidMessage
struct ServerResponseBifs_Protocol bifReqs
struct ServerResponseReqs_Protocol secReqs
BindPrefSelector(std::vector< std::string > &&bindprefs)
const std::string & Get()
Data structure that carries the handshake information.
std::string streamName
Name of the stream.
uint16_t subStreamId
Sub-stream id.
Message * out
Message to be sent out.
PathID(uint16_t u=0, uint16_t d=0)
static void UnloadHandler(const std::string &trProt)
void Register(const std::string &protocol)
static void UnloadHandler()
std::set< std::string > protocols
Procedure execution status.
Status(uint16_t st=stOK, uint16_t cod=errNone, uint32_t errN=0)
Constructor.
uint16_t code
Error type, or additional hints on what to do.
bool IsOK() const
We're fine.
void AdjustQueues(uint16_t size)
StreamSelector(uint16_t size)
void MsgReceived(uint16_t substrm)
uint16_t Select(const std::vector< bool > &connected)
static const uint16_t Name
Transport name, returns const char *.
static const uint16_t Auth
Transport name, returns std::string *.
Information holder for xrootd channels.
std::vector< XRootDStreamInfo > StreamInfoVector
std::set< uint16_t > sentCloses
std::unique_ptr< StreamSelector > strmSelector
XrdSecParameters * authParams
XrdSecProtocol * authProtocol
XrdSecProtect * protection
std::unique_ptr< BindPrefSelector > bindSelector
std::string authProtocolName
std::atomic< uint32_t > finstcnt
unsigned int protRespSize
std::vector< char > protRespBuff
XRootDChannelInfo(const URL &url)
std::set< uint16_t > sentOpens
std::shared_ptr< SIDManager > sidManager
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
Information holder for XRootDStreams.
char * buffer
Pointer to the buffer.
int size
Size of the buffer or length of data in the buffer.