XRootD
Loading...
Searching...
No Matches
XrdClXRootDTransport.cc
Go to the documentation of this file.
1//------------------------------------------------------------------------------
2// Copyright (c) 2011-2014 by European Organization for Nuclear Research (CERN)
3// Author: Lukasz Janyst <ljanyst@cern.ch>
4//------------------------------------------------------------------------------
5// This file is part of the XRootD software suite.
6//
7// XRootD is free software: you can redistribute it and/or modify
8// it under the terms of the GNU Lesser General Public License as published by
9// the Free Software Foundation, either version 3 of the License, or
10// (at your option) any later version.
11//
12// XRootD is distributed in the hope that it will be useful,
13// but WITHOUT ANY WARRANTY; without even the implied warranty of
14// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
15// GNU General Public License for more details.
16//
17// You should have received a copy of the GNU Lesser General Public License
18// along with XRootD. If not, see <http://www.gnu.org/licenses/>.
19//
20// In applying this licence, CERN does not waive the privileges and immunities
21// granted to it by virtue of its status as an Intergovernmental Organization
22// or submit itself to any jurisdiction.
23//------------------------------------------------------------------------------
24
27#include "XrdCl/XrdClLog.hh"
28#include "XrdCl/XrdClSocket.hh"
29#include "XrdCl/XrdClMessage.hh"
32#include "XrdCl/XrdClUtils.hh"
34#include "XrdCl/XrdClTls.hh"
35#include "XrdNet/XrdNetAddr.hh"
36#include "XrdNet/XrdNetUtils.hh"
39#include "XrdOuc/XrdOucUtils.hh"
40#include "XrdOuc/XrdOucCRC.hh"
42#include "XrdSys/XrdSysTimer.hh"
47#include "XrdSys/XrdSysE2T.hh"
48#include "XrdCl/XrdClTls.hh"
49#include "XrdCl/XrdClSocket.hh"
51#include "XrdVersion.hh"
52
53#include <arpa/inet.h>
54#include <sys/types.h>
55#include <unistd.h>
56#include <dlfcn.h>
57#include <sstream>
58#include <iomanip>
59#include <set>
60#include <limits>
61
62#include <atomic>
63
65
66namespace XrdCl
67{
69 {
71
72 static void UnloadHandler()
73 {
74 UnloadHandler( "root" );
75 UnloadHandler( "xroot" );
76 }
77
78 static void UnloadHandler( const std::string &trProt )
79 {
81 TransportHandler *trHandler = trManager->GetHandler( trProt );
82 trHandler->WaitBeforeExit();
83 }
84
85 void Register( const std::string &protocol )
86 {
87 XrdSysRWLockHelper scope( lock, false ); // obtain write lock
88 std::pair< std::set<std::string>::iterator, bool > ret = protocols.insert( protocol );
89 // if that's the first time we are using the protocol, the sec lib
90 // was just loaded so now's the time to register the atexit handler
91 if( ret.second )
92 {
93 atexit( UnloadHandler );
94 }
95 }
96
99 std::set<std::string> protocols;
100 };
101
102 //----------------------------------------------------------------------------
104 //----------------------------------------------------------------------------
106 {
107 //--------------------------------------------------------------------------
108 // Define the stream status for the link negotiation purposes
109 //--------------------------------------------------------------------------
122
123 //--------------------------------------------------------------------------
124 // Constructor
125 //--------------------------------------------------------------------------
129
131 uint8_t pathId;
132 };
133
134 //----------------------------------------------------------------------------
136 //----------------------------------------------------------------------------
138 {
139 StreamSelector( uint16_t size )
140 {
141 //----------------------------------------------------------------------
142 // Subtract one because we shouldn't take into account the control
143 // stream.
144 //----------------------------------------------------------------------
145 strmqueues.resize( size - 1, 0 );
146 }
147
148 //------------------------------------------------------------------------
149 // @param size : number of streams
150 //------------------------------------------------------------------------
151 void AdjustQueues( uint16_t size )
152 {
153 strmqueues.resize( size - 1, 0);
154 }
155
156 //------------------------------------------------------------------------
157 // @param connected : bitarray stating if given sub-stream is connected
158 //
159 // @return : substream number
160 //------------------------------------------------------------------------
161 uint16_t Select( const std::vector<bool> &connected )
162 {
163 uint16_t ret = 0;
164 size_t minval = std::numeric_limits<size_t>::max();
165
166 for( size_t i = 0; i < connected.size() && i < strmqueues.size(); ++i )
167 {
168 if( !connected[i] ) continue;
169
170 if( strmqueues[i] < minval )
171 {
172 ret = i;
173 minval = strmqueues[i];
174 }
175 }
176
177 ++strmqueues[ret];
178 return ret + 1;
179 }
180
181 //--------------------------------------------------------------------------
182 // Update queue for given substream
183 //--------------------------------------------------------------------------
184 void MsgReceived( uint16_t substrm )
185 {
186 if( substrm > 0 )
187 --strmqueues[substrm - 1];
188 }
189
190 private:
191
192 std::vector<size_t> strmqueues;
193 };
194
196 {
197 BindPrefSelector( std::vector<std::string> && bindprefs ) :
198 bindprefs( std::move( bindprefs ) ), next( 0 )
199 {
200 }
201
202 inline const std::string& Get()
203 {
204 std::string &ret = bindprefs[next];
205 ++next;
206 if( next >= bindprefs.size() )
207 next = 0;
208 return ret;
209 }
210
211 private:
212 std::vector<std::string> bindprefs;
213 size_t next;
214 };
215
216 //----------------------------------------------------------------------------
218 //----------------------------------------------------------------------------
220 {
221 //--------------------------------------------------------------------------
222 // Constructor
223 //--------------------------------------------------------------------------
224 XRootDChannelInfo( const URL &url ):
225 serverFlags(0),
227 firstLogIn(true),
228 authBuffer(0),
229 authProtocol(0),
230 authParams(0),
231 authEnv(0),
232 finstcnt(0),
233 openFiles(0),
234 waitBarrier(0),
235 protection(0),
236 protRespSize(0),
237 encrypted(false),
238 istpc(false)
239 {
241 memset( sessionId, 0, 16 );
242 memset( oldSessionId, 0, 16 );
243 }
244
245 //--------------------------------------------------------------------------
246 // Destructor
247 //--------------------------------------------------------------------------
249 {
250 delete [] authBuffer;
251 }
252
253 typedef std::vector<XRootDStreamInfo> StreamInfoVector;
254
255 //--------------------------------------------------------------------------
256 // Data
257 //--------------------------------------------------------------------------
258 uint32_t serverFlags;
260 uint8_t sessionId[16];
261 uint8_t oldSessionId[16];
263 std::shared_ptr<SIDManager> sidManager;
269 std::string streamName;
270 std::string authProtocolName;
271 std::set<uint16_t> sentOpens;
272 std::set<uint16_t> sentCloses;
273 std::atomic<uint32_t> finstcnt; // file instance count
274 uint32_t openFiles;
277 std::vector<char> protRespBuff;
278 unsigned int protRespSize;
279 std::unique_ptr<StreamSelector> strmSelector;
281 bool istpc;
282 std::unique_ptr<BindPrefSelector> bindSelector;
283 std::string logintoken;
285 };
286
287 //----------------------------------------------------------------------------
288 // Constructor
289 //----------------------------------------------------------------------------
291 pSecUnloadHandler( new PluginUnloadHandler() )
292 {
293 }
294
295 //----------------------------------------------------------------------------
296 // Destructor
297 //----------------------------------------------------------------------------
299 {
300 delete pSecUnloadHandler; pSecUnloadHandler = 0;
301 }
302
303 //----------------------------------------------------------------------------
304 // Read message header from socket
305 //----------------------------------------------------------------------------
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 }
341 }
342
343 //----------------------------------------------------------------------------
344 // Read message body from socket
345 //----------------------------------------------------------------------------
347 {
348 //--------------------------------------------------------------------------
349 // Retrieve the body
350 //--------------------------------------------------------------------------
351 size_t leftToBeRead = 0;
352 uint32_t bodySize = 0;
354 bodySize = rsphdr->dlen;
355
356 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
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 }
378
379 //----------------------------------------------------------------------------
380 // Read more of the message body from socket
381 //----------------------------------------------------------------------------
383 {
385 if( rsphdr->status != kXR_status )
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 )
398 "kXR_status: response body too large." );
399 if( bodySize+8 < sizeof( ServerResponseStatus ) )
401 "kXR_status: invalid message size." );
402
404 uint32_t moreSize = static_cast<uint32_t>( rspst->bdy.dlen );
405 if( moreSize > std::numeric_limits<uint32_t>::max() - 8 - bodySize )
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();
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 }
445
446 //----------------------------------------------------------------------------
447 // Initialize channel
448 //----------------------------------------------------------------------------
450 AnyObject &channelData )
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 }
466
467 //----------------------------------------------------------------------------
468 // Finalize channel
469 //----------------------------------------------------------------------------
473
474 //----------------------------------------------------------------------------
475 // HandShake
476 //----------------------------------------------------------------------------
478 AnyObject &channelData )
479 {
480 XRootDChannelInfo *info = 0;
481 channelData.Get( info );
482
483 if (!info)
485
486 XrdSysMutexHelper scopedLock( info->mutex );
487
488 if( info->stream.size() <= handShakeData->subStreamId )
489 {
490 Log *log = DefaultEnv::GetLog();
492 "[%s] Internal error: not enough substreams",
493 handShakeData->streamName.c_str() );
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 }
504
505 //----------------------------------------------------------------------------
506 // Hand shake the main stream
507 //----------------------------------------------------------------------------
508 XRootDStatus XRootDTransport::HandShakeMain( HandShakeData *handShakeData,
509 AnyObject &channelData )
510 {
511 XRootDChannelInfo *info = 0;
512 channelData.Get( info );
513
514 if (!info) {
516 "[%s] Internal error: no channel info",
517 handShakeData->streamName.c_str());
519 }
520
521 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
522
523 //--------------------------------------------------------------------------
524 // First step - we need to create and initial handshake and send it out
525 //--------------------------------------------------------------------------
526 if( sInfo.status == XRootDStreamInfo::Disconnected ||
527 sInfo.status == XRootDStreamInfo::Broken )
528 {
529 handShakeData->out = GenerateInitialHSProtocol( handShakeData, info,
531 sInfo.status = XRootDStreamInfo::HandShakeSent;
532 return XRootDStatus( stOK, suContinue );
533 }
534
535 //--------------------------------------------------------------------------
536 // Second step - we got the reply message to the initial handshake
537 //--------------------------------------------------------------------------
538 if( sInfo.status == XRootDStreamInfo::HandShakeSent )
539 {
540 XRootDStatus st = ProcessServerHS( handShakeData, info );
541 if( st.IsOK() )
543 else
544 sInfo.status = XRootDStreamInfo::Broken;
545 return st;
546 }
547
548 //--------------------------------------------------------------------------
549 // Third step - we got the response to the protocol request, we need
550 // to process it and send out a login request
551 //--------------------------------------------------------------------------
552 if( sInfo.status == XRootDStreamInfo::HandShakeReceived )
553 {
554 XRootDStatus st = ProcessProtocolResp( handShakeData, info );
555
556 if( !st.IsOK() )
557 {
558 sInfo.status = XRootDStreamInfo::Broken;
559 return st;
560 }
561
562 if( st.code == suRetry )
563 {
564 handShakeData->out = GenerateProtocol( handShakeData, info,
567 return XRootDStatus( stOK, suRetry );
568 }
569
570 handShakeData->out = GenerateLogIn( handShakeData, info );
571 sInfo.status = XRootDStreamInfo::LoginSent;
572 return XRootDStatus( stOK, suContinue );
573 }
574
575 //--------------------------------------------------------------------------
576 // Fourth step - handle the log in response and proceed with the
577 // authentication if required by the server
578 //--------------------------------------------------------------------------
579 if( sInfo.status == XRootDStreamInfo::LoginSent )
580 {
581 XRootDStatus st = ProcessLogInResp( handShakeData, info );
582
583 if( !st.IsOK() )
584 {
585 sInfo.status = XRootDStreamInfo::Broken;
586 return st;
587 }
588
589 if( st.IsOK() && st.code == suDone )
590 {
591 //----------------------------------------------------------------------
592 // If it's not our first log in we need to end the previous session
593 // to make sure that the server noticed our disconnection and closed
594 // all the writable handles that we owned
595 //----------------------------------------------------------------------
596 if( !info->firstLogIn )
597 {
598 handShakeData->out = GenerateEndSession( handShakeData, info );
600 return XRootDStatus( stOK, suContinue );
601 }
602
603 sInfo.status = XRootDStreamInfo::Connected;
604 info->firstLogIn = false;
605 return st;
606 }
607
608 st = DoAuthentication( handShakeData, info );
609 if( !st.IsOK() )
610 sInfo.status = XRootDStreamInfo::Broken;
611 else
612 sInfo.status = XRootDStreamInfo::AuthSent;
613 return st;
614 }
615
616 //--------------------------------------------------------------------------
617 // Fifth step and later - proceed with the authentication
618 //--------------------------------------------------------------------------
619 if( sInfo.status == XRootDStreamInfo::AuthSent )
620 {
621 XRootDStatus st = DoAuthentication( handShakeData, info );
622
623 if( !st.IsOK() )
624 {
625 sInfo.status = XRootDStreamInfo::Broken;
626 return st;
627 }
628
629 if( st.IsOK() && st.code == suDone )
630 {
631 //----------------------------------------------------------------------
632 // If it's not our first log in we need to end the previous session
633 //----------------------------------------------------------------------
634 if( !info->firstLogIn )
635 {
636 handShakeData->out = GenerateEndSession( handShakeData, info );
638 return XRootDStatus( stOK, suContinue );
639 }
640
641 sInfo.status = XRootDStreamInfo::Connected;
642 info->firstLogIn = false;
643 return st;
644 }
645
646 return st;
647 }
648
649 //--------------------------------------------------------------------------
650 // The last step - kXR_endsess returned
651 //--------------------------------------------------------------------------
652 if( sInfo.status == XRootDStreamInfo::EndSessionSent )
653 {
654 XRootDStatus st = ProcessEndSessionResp( handShakeData, info );
655
656 if( st.IsOK() && st.code == suDone )
657 {
658 sInfo.status = XRootDStreamInfo::Connected;
659 }
660 else if( !st.IsOK() )
661 {
662 sInfo.status = XRootDStreamInfo::Broken;
663 }
664
665 return st;
666 }
667
668 return XRootDStatus( stOK, suDone );
669 }
670
671 //----------------------------------------------------------------------------
672 // Hand shake parallel stream
673 //----------------------------------------------------------------------------
674 XRootDStatus XRootDTransport::HandShakeParallel( HandShakeData *handShakeData,
675 AnyObject &channelData )
676 {
677 XRootDChannelInfo *info = 0;
678 channelData.Get( info );
679
680 if (!info) {
682 "[%s] Internal error: no channel info",
683 handShakeData->streamName.c_str());
684 return XRootDStatus(stFatal, errInternal);
685 }
686
687 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
688
689 //--------------------------------------------------------------------------
690 // First step - we need to create and initial handshake and send it out
691 //--------------------------------------------------------------------------
692 if( sInfo.status == XRootDStreamInfo::Disconnected ||
693 sInfo.status == XRootDStreamInfo::Broken )
694 {
695 handShakeData->out = GenerateInitialHSProtocol( handShakeData, info,
697 sInfo.status = XRootDStreamInfo::HandShakeSent;
698 return XRootDStatus( stOK, suContinue );
699 }
700
701 //--------------------------------------------------------------------------
702 // Second step - we got the reply message to the initial handshake,
703 // if successful we need to send bind
704 //--------------------------------------------------------------------------
705 if( sInfo.status == XRootDStreamInfo::HandShakeSent )
706 {
707 XRootDStatus st = ProcessServerHS( handShakeData, info );
708 if( st.IsOK() )
710 else
711 sInfo.status = XRootDStreamInfo::Broken;
712 return st;
713 }
714
715 //--------------------------------------------------------------------------
716 // Second step bis - we got the response to the protocol request, we need
717 // to process it and send out a bind request
718 //--------------------------------------------------------------------------
719 if( sInfo.status == XRootDStreamInfo::HandShakeReceived )
720 {
721 XRootDStatus st = ProcessProtocolResp( handShakeData, info );
722
723 if( !st.IsOK() )
724 {
725 sInfo.status = XRootDStreamInfo::Broken;
726 return st;
727 }
728
729 handShakeData->out = GenerateBind( handShakeData, info );
730 sInfo.status = XRootDStreamInfo::BindSent;
731 return XRootDStatus( stOK, suContinue );
732 }
733
734 //--------------------------------------------------------------------------
735 // Third step - we got the response to the kXR_bind
736 //--------------------------------------------------------------------------
737 if( sInfo.status == XRootDStreamInfo::BindSent )
738 {
739 XRootDStatus st = ProcessBindResp( handShakeData, info );
740
741 if( !st.IsOK() )
742 {
743 sInfo.status = XRootDStreamInfo::Broken;
744 return st;
745 }
746 sInfo.status = XRootDStreamInfo::Connected;
747 return XRootDStatus();
748 }
749 return XRootDStatus();
750 }
751
752 //------------------------------------------------------------------------
753 // @return true if handshake has been done and stream is connected,
754 // false otherwise
755 //------------------------------------------------------------------------
757 AnyObject &channelData )
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 }
772
773 //----------------------------------------------------------------------------
774 // Check if the stream should be disconnected
775 //----------------------------------------------------------------------------
776 bool XRootDTransport::IsStreamTTLElapsed( time_t inactiveTime,
777 AnyObject &channelData )
778 {
779 XRootDChannelInfo *info = 0;
780 channelData.Get( info );
781
782 Env *env = DefaultEnv::GetEnv();
783 Log *log = DefaultEnv::GetLog();
784
785 if (!info) {
787 "Internal error: no channel info, behaving as if TTL has elapsed");
788 return true;
789 }
790
791 //--------------------------------------------------------------------------
792 // Check the TTL settings for the current server
793 //--------------------------------------------------------------------------
794 int ttl;
795 if( info->serverFlags & kXR_isServer )
796 {
798 env->GetInt( "DataServerTTL", ttl );
799 }
800 else
801 {
803 env->GetInt( "LoadBalancerTTL", ttl );
804 }
805
806 //--------------------------------------------------------------------------
807 // See whether we can give a go-ahead for the disconnection
808 //--------------------------------------------------------------------------
809 XrdSysMutexHelper scopedLock( info->mutex );
810 uint16_t allocatedSIDs = info->sidManager->GetNumberOfAllocatedSIDs();
811 log->Dump( XRootDTransportMsg, "[%s] Stream inactive since %lld seconds, "
812 "TTL: %d, allocated SIDs: %d, open files: %d, bound file objects: %d",
813 info->streamName.c_str(), (long long) inactiveTime, ttl, allocatedSIDs,
814 info->openFiles, info->finstcnt.load( std::memory_order_relaxed ) );
815
816 if( info->openFiles != 0 && info->finstcnt.load( std::memory_order_relaxed ) != 0 )
817 return false;
818
819 if( !allocatedSIDs && inactiveTime > ttl )
820 return true;
821
822 return false;
823 }
824
825 //----------------------------------------------------------------------------
826 // Check the stream is broken - ie. TCP connection got broken and
827 // went undetected by the TCP stack
828 //----------------------------------------------------------------------------
830 AnyObject &channelData )
831 {
832 XRootDChannelInfo *info = 0;
833 channelData.Get( info );
834 Env *env = DefaultEnv::GetEnv();
835 Log *log = DefaultEnv::GetLog();
836
837 if (!info) {
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
867 }
868
869 //----------------------------------------------------------------------------
870 // Multiplex
871 //----------------------------------------------------------------------------
873 {
874 return PathID( 0, 0 );
875 }
876
877 //----------------------------------------------------------------------------
878 // Multiplex
879 //----------------------------------------------------------------------------
881 AnyObject &channelData,
882 PathID *hint )
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 {
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 {
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 );
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 );
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 ) );
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 {
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 }
1046
1047 //----------------------------------------------------------------------------
1048 // Return a number of substreams per stream that should be created
1049 // This depends on the environment and whether we are connected to
1050 // a data server or not
1051 //----------------------------------------------------------------------------
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 }
1109
1110 //----------------------------------------------------------------------------
1111 // Marshall
1112 //----------------------------------------------------------------------------
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 }
1264
1265 //----------------------------------------------------------------------------
1266 // Unmarshall the request - sometimes the requests need to be rewritten,
1267 // so we need to unmarshall them
1268 //----------------------------------------------------------------------------
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 }
1286
1287 //----------------------------------------------------------------------------
1288 // Unmarshall the body of the incoming message
1289 //----------------------------------------------------------------------------
1291 {
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 }
1364
1365 //------------------------------------------------------------------------
1367 //------------------------------------------------------------------------
1369 {
1370 //--------------------------------------------------------------------------
1371 // Calculate the crc32c before the unmarshaling the body!
1372 //--------------------------------------------------------------------------
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 {
1405 pgrdbdy->offset = ntohll( pgrdbdy->offset );
1406 break;
1407 }
1408
1409 case kXR_pgwrite:
1410 {
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 }
1443
1445 {
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 ) )
1462 "kXR_status: invalid message size." );
1463
1464 //--------------------------------------------------------------------------
1465 // Calculate the crc32c for the additional data
1466 //--------------------------------------------------------------------------
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
1502 }
1503
1504 //----------------------------------------------------------------------------
1505 // Unmarshall the header of the incoming message
1506 //----------------------------------------------------------------------------
1508 {
1510 header->status = ntohs( header->status );
1511 header->dlen = ntohl( header->dlen );
1512 }
1513
1514 //----------------------------------------------------------------------------
1515 // Log server error response
1516 //----------------------------------------------------------------------------
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 }
1527
1528 //------------------------------------------------------------------------
1529 // Number of currently connected data streams
1530 //------------------------------------------------------------------------
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 }
1550
1551 //----------------------------------------------------------------------------
1552 // The stream has been disconnected, do the cleanups
1553 //----------------------------------------------------------------------------
1555 uint16_t subStreamId )
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];
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 }
1584
1585 //------------------------------------------------------------------------
1586 // Query the channel
1587 //------------------------------------------------------------------------
1589 AnyObject &result,
1590 AnyObject &channelData )
1591 {
1592 XRootDChannelInfo *info = 0;
1593 channelData.Get( info );
1594
1595 if (!info)
1597
1598 XrdSysMutexHelper scopedLock( info->mutex );
1599
1600 switch( query )
1601 {
1602 //------------------------------------------------------------------------
1603 // Protocol name
1604 //------------------------------------------------------------------------
1606 result.Set( (const char*)"XRootD", false );
1607 return Status();
1608
1609 //------------------------------------------------------------------------
1610 // Authentication
1611 //------------------------------------------------------------------------
1613 result.Set( new std::string( info->authProtocolName ), false );
1614 return Status();
1615
1616 //------------------------------------------------------------------------
1617 // Server flags
1618 //------------------------------------------------------------------------
1620 result.Set( new int( info->serverFlags ), false );
1621 return Status();
1622
1623 //------------------------------------------------------------------------
1624 // Protocol version
1625 //------------------------------------------------------------------------
1627 result.Set( new int( info->protocolVersion ), false );
1628 return Status();
1629
1631 result.Set( new bool( info->encrypted ), false );
1632 return Status();
1633 };
1635 }
1636
1637 //----------------------------------------------------------------------------
1638 // Check whether the transport can hijack the message
1639 //----------------------------------------------------------------------------
1641 uint16_t subStream,
1642 AnyObject &channelData )
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 //--------------------------------------------------------------------------
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 }
1745
1746 //----------------------------------------------------------------------------
1747 // Notify the transport about a message having been sent
1748 //----------------------------------------------------------------------------
1750 uint16_t subStream,
1751 uint32_t bytesSent,
1752 AnyObject &channelData )
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 }
1779
1780
1781 //----------------------------------------------------------------------------
1782 // Get signature for given message
1783 //----------------------------------------------------------------------------
1785 {
1786 XRootDChannelInfo *info = 0;
1787 channelData.Get( info );
1788 return GetSignature( toSign, sign, info );
1789 }
1790
1791 //------------------------------------------------------------------------
1793 //------------------------------------------------------------------------
1795 Message *&sign,
1796 XRootDChannelInfo *info )
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 }
1820
1821 //------------------------------------------------------------------------
1823 //------------------------------------------------------------------------
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 }
1831
1832 //----------------------------------------------------------------------------
1833 // Wait before exit
1834 //----------------------------------------------------------------------------
1836 {
1837 XrdSysRWLockHelper scope( pSecUnloadHandler->lock, false ); // obtain write lock
1838 pSecUnloadHandler->unloaded = true;
1839 }
1840
1841 //----------------------------------------------------------------------------
1842 // @return : true if encryption should be turned on, false otherwise
1843 //----------------------------------------------------------------------------
1845 AnyObject &channelData )
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 //--------------------------------------------------------------------
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 }
1919
1920 //------------------------------------------------------------------------
1921 // Get bind preference for the next data stream
1922 //------------------------------------------------------------------------
1924 AnyObject &channelData )
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 }
1934
1935 //----------------------------------------------------------------------------
1936 // Generate the message to be sent as an initial handshake
1937 // (handshake+kXR_protocol)
1938 //----------------------------------------------------------------------------
1939 Message *XRootDTransport::GenerateInitialHSProtocol( HandShakeData *hsData,
1940 XRootDChannelInfo *info,
1941 kXR_char expect )
1942 {
1943 Log *log = DefaultEnv::GetLog();
1945 "[%s] Sending out the initial hand shake + kXR_protocol",
1946 hsData->streamName.c_str() );
1947
1948 Message *msg = new Message();
1949
1950 msg->Allocate( 20+sizeof(ClientProtocolRequest) );
1951 msg->Zero();
1952
1954 init->fourth = htonl(4);
1955 init->fifth = htonl(2012);
1956
1958 InitProtocolReq( proto, info, expect );
1959
1960 return msg;
1961 }
1962
1963 //------------------------------------------------------------------------
1964 // Generate the protocol message
1965 //------------------------------------------------------------------------
1966 Message *XRootDTransport::GenerateProtocol( HandShakeData *hsData,
1967 XRootDChannelInfo *info,
1968 kXR_char expect )
1969 {
1970 Log *log = DefaultEnv::GetLog();
1971 log->Debug( XRootDTransportMsg,
1972 "[%s] Sending out the kXR_protocol",
1973 hsData->streamName.c_str() );
1974
1975 Message *msg = new Message();
1976 msg->Allocate( sizeof(ClientProtocolRequest) );
1977 msg->Zero();
1978
1979 ClientProtocolRequest *proto = (ClientProtocolRequest *)msg->GetBuffer();
1980 InitProtocolReq( proto, info, expect );
1981
1982 return msg;
1983 }
1984
1985 //------------------------------------------------------------------------
1986 // Initialize protocol request
1987 //------------------------------------------------------------------------
1988 void XRootDTransport::InitProtocolReq( ClientProtocolRequest *request,
1989 XRootDChannelInfo *info,
1990 kXR_char expect )
1991 {
1992 request->requestid = htons(kXR_protocol);
1993 request->clientpv = htonl(kXR_PROTOCOLVERSION);
1996
1997 int notlsok = DefaultNoTlsOK;
1998 int tlsnodata = DefaultTlsNoData;
1999
2000 XrdCl::Env *env = XrdCl::DefaultEnv::GetEnv();
2001
2002 env->GetInt( "NoTlsOK", notlsok );
2003
2005 env->GetInt( "TlsNoData", tlsnodata );
2006
2007 if (info->encrypted || InitTLS())
2009
2010 if (info->encrypted && !(notlsok || tlsnodata))
2012
2013 request->expect = expect;
2014
2015 //--------------------------------------------------------------------------
2016 // If we are in the curse of establishing a connection in the context of
2017 // TPC update the expect! (this will be never followed be a bind)
2018 //--------------------------------------------------------------------------
2019 if( info->istpc )
2021 }
2022
2023 //----------------------------------------------------------------------------
2024 // Process the server initial handshake response
2025 //----------------------------------------------------------------------------
2026 XRootDStatus XRootDTransport::ProcessServerHS( HandShakeData *hsData,
2027 XRootDChannelInfo *info )
2028 {
2029 Log *log = DefaultEnv::GetLog();
2030
2031 Message *msg = hsData->in;
2032 ServerResponseHeader *respHdr = (ServerResponseHeader *)msg->GetBuffer();
2033 ServerInitHandShake *hs = (ServerInitHandShake *)msg->GetBuffer(4);
2034
2035 if( respHdr->status != kXR_ok )
2036 {
2037 log->Error( XRootDTransportMsg, "[%s] Invalid hand shake response",
2038 hsData->streamName.c_str() );
2039
2040 return XRootDStatus( stFatal, errHandShakeFailed, 0, "Invalid hand shake response." );
2041 }
2042
2043 info->protocolVersion = ntohl(hs->protover);
2044 info->serverFlags = ntohl(hs->msgval) == kXR_DataServer ?
2047
2048 log->Debug( XRootDTransportMsg,
2049 "[%s] Got the server hand shake response (%s, protocol "
2050 "version %x)",
2051 hsData->streamName.c_str(),
2052 ServerFlagsToStr( info->serverFlags ).c_str(),
2053 info->protocolVersion );
2054
2055 return XRootDStatus( stOK, suContinue );
2056 }
2057
2058 //----------------------------------------------------------------------------
2059 // Process the protocol response
2060 //----------------------------------------------------------------------------
2061 XRootDStatus XRootDTransport::ProcessProtocolResp( HandShakeData *hsData,
2062 XRootDChannelInfo *info )
2063 {
2064 Log *log = DefaultEnv::GetLog();
2065
2066 XRootDStatus st = UnMarshallBody( hsData->in, kXR_protocol );
2067 if( !st.IsOK() )
2068 return st;
2069
2070 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2071
2072
2073 if( rsp->hdr.status != kXR_ok )
2074 {
2075 log->Error( XRootDTransportMsg, "[%s] kXR_protocol request failed",
2076 hsData->streamName.c_str() );
2077
2078 return XRootDStatus( stFatal, errHandShakeFailed, 0, "kXR_protocol request failed" );
2079 }
2080
2081 XrdCl::Env *env = XrdCl::DefaultEnv::GetEnv();
2082 int notlsok = DefaultNoTlsOK;
2083 env->GetInt( "NoTlsOK", notlsok );
2084
2085 if( rsp->body.protocol.pval < kXR_PROTTLSVERSION && info->encrypted )
2086 {
2087 //------------------------------------------------------------------------
2088 // User requested an encrypted connection but the server is to old to
2089 // support it!
2090 //------------------------------------------------------------------------
2091 if( !notlsok ) return XRootDStatus( stFatal, errTlsError, ENOTSUP, "TLS not supported" );
2092
2093 //------------------------------------------------------------------------
2094 // We are falling back to unencrypted data transmission, as configured
2095 // in XRD_NOTLSOK environment variable
2096 //------------------------------------------------------------------------
2097 log->Info( XRootDTransportMsg,
2098 "[%s] Falling back to unencrypted transmission, server does "
2099 "not support TLS encryption.",
2100 hsData->streamName.c_str() );
2101 info->encrypted = false;
2102 }
2103
2104 if( rsp->body.protocol.pval >= 0x297 )
2105 info->serverFlags = rsp->body.protocol.flags;
2106
2107 if( rsp->hdr.dlen > 8 )
2108 {
2109 info->protRespBuff.assign( sizeof( ServerResponseBody_Protocol ), 0 );
2110 info->protRespSize = 0;
2111 ServerResponseBody_Protocol *protRespBody =
2112 reinterpret_cast<ServerResponseBody_Protocol*>( info->protRespBuff.data() );
2113 protRespBody->flags = rsp->body.protocol.flags;
2114 protRespBody->pval = rsp->body.protocol.pval;
2115
2116 char* bodybuff = reinterpret_cast<char*>( &rsp->body.protocol.secreq );
2117 size_t bodysize = rsp->hdr.dlen - 8;
2118 XRootDStatus st = ProcessProtocolBody( bodybuff, bodysize, info );
2119 if( !st.IsOK() )
2120 return st;
2121 }
2122
2123 log->Debug( XRootDTransportMsg,
2124 "[%s] kXR_protocol successful (%s, protocol version %x)",
2125 hsData->streamName.c_str(),
2126 ServerFlagsToStr( info->serverFlags ).c_str(),
2127 info->protocolVersion );
2128
2129 if( !( info->serverFlags & kXR_haveTLS ) && info->encrypted )
2130 {
2131 //------------------------------------------------------------------------
2132 // User requested an encrypted connection but the server was not configured
2133 // to support encryption!
2134 //------------------------------------------------------------------------
2135 return XRootDStatus( stFatal, errTlsError, ECONNREFUSED,
2136 "Server was not configured to support encryption." );
2137 }
2138
2139 //--------------------------------------------------------------------------
2140 // Now see if we have to enforce encryption in case the server does not
2141 // support PgRead/PgWrite
2142 //--------------------------------------------------------------------------
2143 int tlsOnNoPgrw = DefaultWantTlsOnNoPgrw;
2144 env->GetInt( "WantTlsOnNoPgrw", tlsOnNoPgrw );
2145 if( !( info->serverFlags & kXR_suppgrw ) && tlsOnNoPgrw )
2146 {
2147 //------------------------------------------------------------------------
2148 // If user requested encryption just make sure it is not switched off for
2149 // data
2150 //------------------------------------------------------------------------
2151 if( info->encrypted )
2152 {
2153 log->Debug( XRootDTransportMsg,
2154 "[%s] Server does not support PgRead/PgWrite and"
2155 " WantTlsOnNoPgrw is on; enforcing encryption for data.",
2156 hsData->streamName.c_str() );
2157 env->PutInt( "TlsNoData", DefaultTlsNoData );
2158 }
2159 //------------------------------------------------------------------------
2160 // Otherwise, if server is not enforcing data encryption, we will need to
2161 // redo the protocol request with kXR_wantTLS set.
2162 //------------------------------------------------------------------------
2163 else if( !( info->serverFlags & kXR_tlsData ) &&
2164 ( info->serverFlags & kXR_haveTLS ) )
2165 {
2166 info->encrypted = true;
2167 return XRootDStatus( stOK, suRetry );
2168 }
2169 }
2170
2171 return XRootDStatus( stOK, suContinue );
2172 }
2173
2174 XRootDStatus XRootDTransport::ProcessProtocolBody( char *bodybuff,
2175 size_t bodysize,
2176 XRootDChannelInfo *info )
2177 {
2178 //--------------------------------------------------------------------------
2179 // Parse bind preferences
2180 //--------------------------------------------------------------------------
2181 XrdProto::bifReqs *bifreq = reinterpret_cast<XrdProto::bifReqs*>( bodybuff );
2182 if( bodysize >= sizeof( XrdProto::bifReqs ) && bifreq->theTag == 'B' )
2183 {
2184 bodybuff += sizeof( XrdProto::bifReqs );
2185 bodysize -= sizeof( XrdProto::bifReqs );
2186
2187 if( bodysize < bifreq->bifILen )
2188 return XRootDStatus( stError, errDataError, 0, "Received incomplete "
2189 "protocol response." );
2190 std::string bindprefs_str( bodybuff, bifreq->bifILen );
2191 std::vector<std::string> bindprefs;
2192 Utils::splitString( bindprefs, bindprefs_str, "," );
2193 info->bindSelector.reset( new BindPrefSelector( std::move( bindprefs ) ) );
2194 bodybuff += bifreq->bifILen;
2195 bodysize -= bifreq->bifILen;
2196 }
2197 //--------------------------------------------------------------------------
2198 // Parse security requirements
2199 //--------------------------------------------------------------------------
2200 XrdProto::secReqs *secreq = reinterpret_cast<XrdProto::secReqs*>( bodybuff );
2201 static const size_t secHdrLen = sizeof( XrdProto::secReqs ) -
2202 sizeof( ServerResponseSVec_Protocol );
2203 if( bodysize >= secHdrLen && secreq->theTag == 'S' )
2204 {
2205 //------------------------------------------------------------------------
2206 // Copy only the header and the secvsz entries of the security vector,
2207 // the server may send fewer bytes than declared or trailing garbage
2208 //------------------------------------------------------------------------
2209 size_t secsize = secHdrLen + secreq->secvsz *
2210 sizeof( ServerResponseSVec_Protocol );
2211 if( bodysize < secsize )
2212 return XRootDStatus( stError, errDataError, 0, "Received incomplete "
2213 "protocol response." );
2214
2215 size_t respsize = kXR_ShortProtRespLen + secsize;
2216 if( info->protRespBuff.size() < respsize )
2217 info->protRespBuff.resize( respsize, 0 );
2218 memcpy( info->protRespBuff.data() + kXR_ShortProtRespLen, secreq, secsize );
2219 info->protRespSize = respsize;
2220 }
2221
2222 return XRootDStatus();
2223 }
2224
2225 //----------------------------------------------------------------------------
2226 // Generate the bind message
2227 //----------------------------------------------------------------------------
2228 Message *XRootDTransport::GenerateBind( HandShakeData *hsData,
2229 XRootDChannelInfo *info )
2230 {
2231 Log *log = DefaultEnv::GetLog();
2232
2233 log->Debug( XRootDTransportMsg,
2234 "[%s] Sending out the bind request",
2235 hsData->streamName.c_str() );
2236
2237
2238 Message *msg = new Message( sizeof( ClientBindRequest ) );
2239 ClientBindRequest *bindReq = (ClientBindRequest *)msg->GetBuffer();
2240
2241 bindReq->requestid = kXR_bind;
2242 memcpy( bindReq->sessid, info->sessionId, 16 );
2243 bindReq->dlen = 0;
2244 MarshallRequest( msg );
2245 return msg;
2246 }
2247
2248 //----------------------------------------------------------------------------
2249 // Generate the bind message
2250 //----------------------------------------------------------------------------
2251 XRootDStatus XRootDTransport::ProcessBindResp( HandShakeData *hsData,
2252 XRootDChannelInfo *info )
2253 {
2254 Log *log = DefaultEnv::GetLog();
2255
2256 XRootDStatus st = UnMarshallBody( hsData->in, kXR_bind );
2257 if( !st.IsOK() )
2258 return st;
2259
2260 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2261
2262 if( rsp->hdr.status != kXR_ok )
2263 {
2264 log->Error( XRootDTransportMsg, "[%s] kXR_bind request failed",
2265 hsData->streamName.c_str() );
2266 return XRootDStatus( stFatal, errHandShakeFailed, 0, "kXR_bind request failed" );
2267 }
2268
2269 info->stream[hsData->subStreamId].pathId = rsp->body.bind.substreamid;
2270 log->Debug( XRootDTransportMsg, "[%s] kXR_bind successful",
2271 hsData->streamName.c_str() );
2272
2273 return XRootDStatus();
2274 }
2275
2276 //----------------------------------------------------------------------------
2277 // Generate the login message
2278 //----------------------------------------------------------------------------
2279 Message *XRootDTransport::GenerateLogIn( HandShakeData *hsData,
2280 XRootDChannelInfo *info )
2281 {
2282 Log *log = DefaultEnv::GetLog();
2283 Env *env = DefaultEnv::GetEnv();
2284
2285 //--------------------------------------------------------------------------
2286 // Compute the login cgi
2287 //--------------------------------------------------------------------------
2288 int timeZone = XrdSysTimer::TimeZone();
2289 char *hostName = XrdNetUtils::MyHostName();
2290 std::string countryCode = Utils::FQDNToCC( hostName );
2291 char *cgiBuffer = new char[1024 + info->logintoken.size()];
2292 std::string appName;
2293 std::string monInfo;
2294 env->GetString( "AppName", appName );
2295 env->GetString( "MonInfo", monInfo );
2296 if( info->logintoken.empty() )
2297 {
2298 snprintf( cgiBuffer, 1024,
2299 "xrd.cc=%s&xrd.tz=%d&xrd.appname=%s&xrd.info=%s&"
2300 "xrd.hostname=%s&xrd.rn=%s", countryCode.c_str(), timeZone,
2301 appName.c_str(), monInfo.c_str(), hostName, XrdVERSION );
2302 }
2303 else
2304 {
2305 snprintf( cgiBuffer, 1024,
2306 "xrd.cc=%s&xrd.tz=%d&xrd.appname=%s&xrd.info=%s&"
2307 "xrd.hostname=%s&xrd.rn=%s&%s", countryCode.c_str(), timeZone,
2308 appName.c_str(), monInfo.c_str(), hostName, XrdVERSION, info->logintoken.c_str() );
2309 }
2310 uint16_t cgiLen = strlen( cgiBuffer );
2311 free( hostName );
2312
2313 //--------------------------------------------------------------------------
2314 // Generate the message
2315 //--------------------------------------------------------------------------
2316 Message *msg = new Message( sizeof(ClientLoginRequest) + cgiLen );
2317 ClientLoginRequest *loginReq = (ClientLoginRequest *)msg->GetBuffer();
2318
2319 loginReq->requestid = kXR_login;
2320 loginReq->pid = ::getpid();
2321 loginReq->capver[0] = (kXR_char) kXR_asyncap | (kXR_char) kXR_ver005;
2322 loginReq->dlen = cgiLen;
2324#ifdef WITH_XRDEC
2325 loginReq->ability2 = kXR_ecredir;
2326#endif
2327
2328 int multiProtocol = 0;
2329 env->GetInt( "MultiProtocol", multiProtocol );
2330 if(multiProtocol)
2331 loginReq->ability |= kXR_multipr;
2332
2333 //--------------------------------------------------------------------------
2334 // Check the IP stacks
2335 //--------------------------------------------------------------------------
2337 bool dualStack = false;
2338 bool privateIPv6 = false;
2339 bool privateIPv4 = false;
2340
2341 if( (stacks & XrdNetUtils::hasIP64) == XrdNetUtils::hasIP64 )
2342 {
2343 dualStack = true;
2344 loginReq->ability |= kXR_hasipv64;
2345 }
2346
2347 if( (stacks & XrdNetUtils::hasIPv6) && !(stacks & XrdNetUtils::hasPub6) )
2348 {
2349 privateIPv6 = true;
2350 loginReq->ability |= kXR_onlyprv6;
2351 }
2352
2353 if( (stacks & XrdNetUtils::hasIPv4) && !(stacks & XrdNetUtils::hasPub4) )
2354 {
2355 privateIPv4 = true;
2356 loginReq->ability |= kXR_onlyprv4;
2357 }
2358
2359 // The following code snippet tries to overcome the problem that this host
2360 // may still be dual-stacked but we don't know it because one of the
2361 // interfaces was not registered in DNS.
2362 //
2363 if( !dualStack && hsData->serverAddr )
2364 {if ( ( ( stacks & XrdNetUtils::hasIPv4 )
2365 && hsData->serverAddr->isIPType(XrdNetAddrInfo::IPv6))
2366 || ( ( stacks & XrdNetUtils::hasIPv6 )
2367 && hsData->serverAddr->isIPType(XrdNetAddrInfo::IPv4)))
2368 {dualStack = true;
2369 loginReq->ability |= kXR_hasipv64;
2370 }
2371 }
2372
2373 //--------------------------------------------------------------------------
2374 // Check the username
2375 //--------------------------------------------------------------------------
2376 std::string buffer( 8, 0 );
2377 if( hsData->url->GetUserName().length() )
2378 buffer = hsData->url->GetUserName();
2379 else
2380 {
2381 char *name = new char[1024];
2382 if( !XrdOucUtils::UserName( geteuid(), name, 1024 ) )
2383 buffer = name;
2384 else
2385 buffer = "_anon_";
2386 delete [] name;
2387 }
2388 buffer.resize( 8, 0 );
2389 std::copy( buffer.begin(), buffer.end(), (char*)loginReq->username );
2390
2391 msg->Append( cgiBuffer, cgiLen, 24 );
2392
2393 log->Debug( XRootDTransportMsg, "[%s] Sending out kXR_login request, "
2394 "username: %s, cgi: %s, dual-stack: %s, private IPv4: %s, "
2395 "private IPv6: %s", hsData->streamName.c_str(),
2396 loginReq->username, cgiBuffer, dualStack ? "true" : "false",
2397 privateIPv4 ? "true" : "false",
2398 privateIPv6 ? "true" : "false" );
2399
2400 delete [] cgiBuffer;
2401 MarshallRequest( msg );
2402 return msg;
2403 }
2404
2405 //----------------------------------------------------------------------------
2406 // Process the protocol response
2407 //----------------------------------------------------------------------------
2408 XRootDStatus XRootDTransport::ProcessLogInResp( HandShakeData *hsData,
2409 XRootDChannelInfo *info )
2410 {
2411 Log *log = DefaultEnv::GetLog();
2412
2413 XRootDStatus st = UnMarshallBody( hsData->in, kXR_login );
2414 if( !st.IsOK() )
2415 return st;
2416
2417 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2418
2419 if( rsp->hdr.status != kXR_ok )
2420 {
2421 log->Error( XRootDTransportMsg, "[%s] Got invalid login response",
2422 hsData->streamName.c_str() );
2423 return XRootDStatus( stFatal, errLoginFailed, 0, "Got invalid login response." );
2424 }
2425
2426 if( !info->firstLogIn )
2427 memcpy( info->oldSessionId, info->sessionId, 16 );
2428
2429 if( rsp->hdr.dlen == 0 && info->protocolVersion <= 0x289 )
2430 {
2431 //--------------------------------------------------------------------------
2432 // This if statement is there only to support dCache inaccurate
2433 // implementation of XRoot protocol, that in some cases returns
2434 // an empty login response for protocol version <= 2.8.9.
2435 //--------------------------------------------------------------------------
2436 memset( info->sessionId, 0, 16 );
2437 log->Warning( XRootDTransportMsg,
2438 "[%s] Logged in, accepting empty login response.",
2439 hsData->streamName.c_str() );
2440 return XRootDStatus();
2441 }
2442
2443 if( rsp->hdr.dlen < 16 )
2444 return XRootDStatus( stError, errDataError, 0, "Login response too short." );
2445
2446 memcpy( info->sessionId, rsp->body.login.sessid, 16 );
2447
2448 std::string sessId = Utils::Char2Hex( rsp->body.login.sessid, 16 );
2449
2450 log->Debug( XRootDTransportMsg, "[%s] Logged in, session: %s",
2451 hsData->streamName.c_str(), sessId.c_str() );
2452
2453 //--------------------------------------------------------------------------
2454 // We have an authentication info to process
2455 //--------------------------------------------------------------------------
2456 if( rsp->hdr.dlen > 16 )
2457 {
2458 size_t len = rsp->hdr.dlen-16;
2459 info->authBuffer = new char[len+1];
2460 info->authBuffer[len] = 0;
2461 memcpy( info->authBuffer, rsp->body.login.sec, len );
2462 log->Debug( XRootDTransportMsg, "[%s] Authentication is required: %s",
2463 hsData->streamName.c_str(), info->authBuffer );
2464
2465 return XRootDStatus( stOK, suContinue );
2466 }
2467
2468 return XRootDStatus();
2469 }
2470
2471 //----------------------------------------------------------------------------
2472 // Do the authentication
2473 //----------------------------------------------------------------------------
2474 XRootDStatus XRootDTransport::DoAuthentication( HandShakeData *hsData,
2475 XRootDChannelInfo *info )
2476 {
2477 //--------------------------------------------------------------------------
2478 // Prepare
2479 //--------------------------------------------------------------------------
2480 Log *log = DefaultEnv::GetLog();
2481 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2482 XrdSecCredentials *credentials = 0;
2483 std::string protocolName;
2484
2485 //--------------------------------------------------------------------------
2486 // We're doing this for the first time
2487 //--------------------------------------------------------------------------
2488 if( sInfo.status == XRootDStreamInfo::LoginSent )
2489 {
2490 log->Debug( XRootDTransportMsg, "[%s] Sending authentication data",
2491 hsData->streamName.c_str() );
2492
2493 //------------------------------------------------------------------------
2494 // Set up the authentication environment
2495 //------------------------------------------------------------------------
2496 info->authEnv = new XrdOucEnv();
2497 info->authEnv->Put( "sockname", hsData->clientName.c_str() );
2498 info->authEnv->Put( "username", hsData->url->GetUserName().c_str() );
2499 info->authEnv->Put( "password", hsData->url->GetPassword().c_str() );
2500
2501 const URL::ParamsMap &urlParams = hsData->url->GetParams();
2502 URL::ParamsMap::const_iterator it;
2503 for( it = urlParams.begin(); it != urlParams.end(); ++it )
2504 {
2505 if( it->first.compare( 0, 4, "xrd." ) == 0 ||
2506 it->first.compare( 0, 6, "xrdcl." ) == 0 )
2507 info->authEnv->Put( it->first.c_str(), it->second.c_str() );
2508 }
2509
2510 //------------------------------------------------------------------------
2511 // Initialize some other structs
2512 //------------------------------------------------------------------------
2513 size_t authBuffLen = strlen( info->authBuffer );
2514 char *pars = (char *)malloc( authBuffLen + 1 );
2515 memcpy( pars, info->authBuffer, authBuffLen );
2516 info->authParams = new XrdSecParameters( pars, authBuffLen );
2517 sInfo.status = XRootDStreamInfo::AuthSent;
2518 delete [] info->authBuffer;
2519 info->authBuffer = 0;
2520
2521 //------------------------------------------------------------------------
2522 // Find a protocol that gives us valid credentials
2523 //------------------------------------------------------------------------
2524 XRootDStatus st = GetCredentials( credentials, hsData, info );
2525 if( !st.IsOK() )
2526 {
2527 CleanUpAuthentication( info );
2528 return st;
2529 }
2530 protocolName = info->authProtocol->Entity.prot;
2531 }
2532
2533 //--------------------------------------------------------------------------
2534 // We've been here already
2535 //--------------------------------------------------------------------------
2536 else
2537 {
2538 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2539 protocolName = info->authProtocol->Entity.prot;
2540
2541 //------------------------------------------------------------------------
2542 // We're required to send out more authentication data
2543 //------------------------------------------------------------------------
2544 if( rsp->hdr.status == kXR_authmore )
2545 {
2546 log->Debug( XRootDTransportMsg,
2547 "[%s] Sending more authentication data for %s",
2548 hsData->streamName.c_str(), protocolName.c_str() );
2549
2550 uint32_t len = rsp->hdr.dlen;
2551 char *secTokenData = (char*)malloc( len );
2552 memcpy( secTokenData, rsp->body.authmore.data, len );
2553 XrdSecParameters *secToken = new XrdSecParameters( secTokenData, len );
2554 XrdOucErrInfo ei( "", info->authEnv);
2555 credentials = info->authProtocol->getCredentials( secToken, &ei );
2556 delete secToken;
2557
2558 //----------------------------------------------------------------------
2559 // The protocol handler refuses to give us the data
2560 //----------------------------------------------------------------------
2561 if( !credentials )
2562 {
2563 log->Error( XRootDTransportMsg,
2564 "[%s] Auth protocol handler for %s refuses to give "
2565 "us more credentials %s",
2566 hsData->streamName.c_str(), protocolName.c_str(),
2567 ei.getErrText() );
2568 CleanUpAuthentication( info );
2569 return XRootDStatus( stFatal, errAuthFailed, 0, ei.getErrText() );
2570 }
2571 }
2572
2573 //------------------------------------------------------------------------
2574 // We have succeeded
2575 //------------------------------------------------------------------------
2576 else if( rsp->hdr.status == kXR_ok )
2577 {
2578 info->authProtocolName = info->authProtocol->Entity.prot;
2579
2580 //----------------------------------------------------------------------
2581 // Do we need protection?
2582 //----------------------------------------------------------------------
2583 if( !info->protRespBuff.empty() )
2584 {
2585 ServerResponseBody_Protocol *protRespBody =
2586 reinterpret_cast<ServerResponseBody_Protocol*>( info->protRespBuff.data() );
2587 int rc = XrdSecGetProtection( info->protection, *info->authProtocol, *protRespBody, info->protRespSize );
2588 if( rc > 0 )
2589 {
2590 log->Debug( XRootDTransportMsg,
2591 "[%s] XrdSecProtect loaded.", hsData->streamName.c_str() );
2592 }
2593 else if( rc == 0 )
2594 {
2595 log->Debug( XRootDTransportMsg,
2596 "[%s] XrdSecProtect: no protection needed.",
2597 hsData->streamName.c_str() );
2598 }
2599 else
2600 {
2601 log->Debug( XRootDTransportMsg,
2602 "[%s] Failed to load XrdSecProtect: %s",
2603 hsData->streamName.c_str(), XrdSysE2T( -rc ) );
2604 CleanUpAuthentication( info );
2605
2606 return XRootDStatus( stError, errAuthFailed, -rc, XrdSysE2T( -rc ) );
2607 }
2608 }
2609
2610 if( !info->protection )
2611 CleanUpAuthentication( info );
2612 else
2613 pSecUnloadHandler->Register( info->authProtocolName );
2614
2615 log->Debug( XRootDTransportMsg,
2616 "[%s] Authenticated with %s.", hsData->streamName.c_str(),
2617 protocolName.c_str() );
2618
2619 //--------------------------------------------------------------------
2620 // Clear the SSL error queue of the calling thread, as there might be
2621 // some leftover from the authentication!
2622 //--------------------------------------------------------------------
2624
2625 return XRootDStatus();
2626 }
2627 //------------------------------------------------------------------------
2628 // Failure
2629 //------------------------------------------------------------------------
2630 else if( rsp->hdr.status == kXR_error )
2631 {
2632 char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
2633 memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
2634 log->Error( XRootDTransportMsg,
2635 "[%s] Authentication with %s failed: %s",
2636 hsData->streamName.c_str(), protocolName.c_str(),
2637 errmsg );
2638 delete [] errmsg;
2639
2640 info->authProtocol->Delete();
2641 info->authProtocol = 0;
2642
2643 //----------------------------------------------------------------------
2644 // Find another protocol that gives us valid credentials
2645 //----------------------------------------------------------------------
2646 XRootDStatus st = GetCredentials( credentials, hsData, info );
2647 if( !st.IsOK() )
2648 {
2649 CleanUpAuthentication( info );
2650 return st;
2651 }
2652 protocolName = info->authProtocol->Entity.prot;
2653 }
2654 //------------------------------------------------------------------------
2655 // God knows what
2656 //------------------------------------------------------------------------
2657 else
2658 {
2659 info->authProtocolName = info->authProtocol->Entity.prot;
2660 CleanUpAuthentication( info );
2661
2662 log->Error( XRootDTransportMsg,
2663 "[%s] Authentication with %s failed: unexpected answer",
2664 hsData->streamName.c_str(), protocolName.c_str() );
2665 return XRootDStatus( stFatal, errAuthFailed, 0, "Authentication failed: unexpected answer." );
2666 }
2667 }
2668
2669 //--------------------------------------------------------------------------
2670 // Generate the client request
2671 //--------------------------------------------------------------------------
2672 Message *msg = new Message( sizeof(ClientAuthRequest)+credentials->size );
2673 msg->Zero();
2674 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
2675 char *reqBuffer = msg->GetBuffer(sizeof(ClientAuthRequest));
2676
2677 req->header.requestid = kXR_auth;
2678 req->auth.dlen = credentials->size;
2679 memcpy( req->auth.credtype, protocolName.c_str(),
2680 protocolName.length() > 4 ? 4 : protocolName.length() );
2681
2682 memcpy( reqBuffer, credentials->buffer, credentials->size );
2683 hsData->out = msg;
2684 MarshallRequest( msg );
2685 delete credentials;
2686
2687 //------------------------------------------------------------------------
2688 // Clear the SSL error queue of the calling thread, as there might be
2689 // some leftover from the authentication!
2690 //------------------------------------------------------------------------
2692
2693 return XRootDStatus( stOK, suContinue );
2694 }
2695
2696 //------------------------------------------------------------------------
2697 // Get the initial credentials using one of the protocols
2698 //------------------------------------------------------------------------
2699 XRootDStatus XRootDTransport::GetCredentials( XrdSecCredentials *&credentials,
2700 HandShakeData *hsData,
2701 XRootDChannelInfo *info )
2702 {
2703 //--------------------------------------------------------------------------
2704 // Set up the auth handler
2705 //--------------------------------------------------------------------------
2706 Log *log = DefaultEnv::GetLog();
2707 XrdOucErrInfo ei( "", info->authEnv);
2708 XrdSecGetProt_t authHandler = GetAuthHandler();
2709 if( !authHandler )
2710 return XRootDStatus( stFatal, errAuthFailed, 0, "Could not load authentication handler." );
2711
2712 //--------------------------------------------------------------------------
2713 // Retrieve secuid and secgid, if available. These will override the fsuid
2714 // and fsgid of the current thread reading the credentials to prevent
2715 // security holes in case this process is running with elevated permissions.
2716 //--------------------------------------------------------------------------
2717 char *secuidc = (ei.getEnv()) ? ei.getEnv()->Get("xrdcl.secuid") : 0;
2718 char *secgidc = (ei.getEnv()) ? ei.getEnv()->Get("xrdcl.secgid") : 0;
2719
2720 int secuid = -1;
2721 int secgid = -1;
2722
2723 if(secuidc) secuid = atoi(secuidc);
2724 if(secgidc) secgid = atoi(secgidc);
2725
2726#ifdef __linux__
2727 ScopedFsUidSetter uidSetter(secuid, secgid, hsData->streamName);
2728 if(!uidSetter.IsOk()) {
2729 log->Error( XRootDTransportMsg, "[%s] Error while setting (fsuid, fsgid) to (%d, %d)",
2730 hsData->streamName.c_str(), secuid, secgid );
2731 return XRootDStatus( stFatal, errAuthFailed, 0, "Error while setting (fsuid, fsgid)." );
2732 }
2733#else
2734 if(secuid >= 0 || secgid >= 0) {
2735 log->Error( XRootDTransportMsg, "[%s] xrdcl.secuid and xrdcl.secgid only supported on Linux.",
2736 hsData->streamName.c_str() );
2737 return XRootDStatus( stFatal, errAuthFailed, 0, "xrdcl.secuid and xrdcl.secgid"
2738 " only supported on Linux" );
2739 }
2740#endif
2741
2742 //--------------------------------------------------------------------------
2743 // Loop over the possible protocols to find one that gives us valid
2744 // credentials
2745 //--------------------------------------------------------------------------
2746 XrdNetAddr &srvAddrInfo = *const_cast<XrdNetAddr *>(hsData->serverAddr);
2747 srvAddrInfo.SetTLS( info->encrypted );
2748 while(1)
2749 {
2750 //------------------------------------------------------------------------
2751 // Get the protocol
2752 //------------------------------------------------------------------------
2753 info->authProtocol = (*authHandler)( hsData->url->GetHostName().c_str(),
2754 srvAddrInfo,
2755 *info->authParams,
2756 &ei );
2757 if( !info->authProtocol )
2758 {
2759 log->Error( XRootDTransportMsg, "[%s] No protocols left to try",
2760 hsData->streamName.c_str() );
2761 return XRootDStatus( stFatal, errAuthFailed, 0, "No protocols left to try" );
2762 }
2763
2764 std::string protocolName = info->authProtocol->Entity.prot;
2765 log->Debug( XRootDTransportMsg, "[%s] Trying to authenticate using %s",
2766 hsData->streamName.c_str(), protocolName.c_str() );
2767
2768 //------------------------------------------------------------------------
2769 // Get the credentials from the current protocol
2770 //------------------------------------------------------------------------
2771 credentials = info->authProtocol->getCredentials( 0, &ei );
2772 if( !credentials )
2773 {
2774 log->Debug( XRootDTransportMsg,
2775 "[%s] Cannot get credentials for protocol %s: %s",
2776 hsData->streamName.c_str(), protocolName.c_str(),
2777 ei.getErrText() );
2778 info->authProtocol->Delete();
2779 continue;
2780 }
2781 return XRootDStatus( stOK, suContinue );
2782 }
2783 }
2784
2785 //------------------------------------------------------------------------
2786 // Clean up the data structures created for the authentication process
2787 //------------------------------------------------------------------------
2788 Status XRootDTransport::CleanUpAuthentication( XRootDChannelInfo *info )
2789 {
2790 if( info->authProtocol )
2791 info->authProtocol->Delete();
2792 delete info->authParams;
2793 delete info->authEnv;
2794 info->authProtocol = 0;
2795 info->authParams = 0;
2796 info->authEnv = 0;
2798 return Status();
2799 }
2800
2801 //------------------------------------------------------------------------
2802 // Clean up the data structures created for the protection purposes
2803 //------------------------------------------------------------------------
2804 Status XRootDTransport::CleanUpProtection( XRootDChannelInfo *info )
2805 {
2806 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
2807 if( pSecUnloadHandler->unloaded ) return Status( stError, errInvalidOp );
2808
2809 if( info->protection )
2810 {
2811 info->protection->Delete();
2812 info->protection = 0;
2813
2814 CleanUpAuthentication( info );
2815 }
2816
2817 info->protRespBuff.clear();
2818 info->protRespSize = 0;
2819
2820 return Status();
2821 }
2822
2823 //----------------------------------------------------------------------------
2824 // Get the authentication function handle
2825 //----------------------------------------------------------------------------
2826 XrdSecGetProt_t XRootDTransport::GetAuthHandler()
2827 {
2828 Log *log = DefaultEnv::GetLog();
2829 char errorBuff[1024];
2830
2831 // the static constructor is invoked only once and it is guaranteed that this
2832 // is thread safe
2833 static std::atomic<XrdSecGetProt_t> authHandler( XrdSecLoadSecFactory( errorBuff, 1024 ) );
2834 auto ret = authHandler.load( std::memory_order_relaxed );
2835 if( ret ) return ret;
2836
2837 // if we are here it means we failed to load the security library for the
2838 // first time and we hope the environment changed
2839
2840 // obtain a lock
2841 static XrdSysMutex mtx;
2842 XrdSysMutexHelper lck( mtx );
2843 // check if in the meanwhile some else didn't load the library
2844 ret = authHandler.load( std::memory_order_relaxed );
2845 if( ret ) return ret;
2846
2847 // load the library
2848 ret = XrdSecLoadSecFactory( errorBuff, 1024 );
2849 authHandler.store( ret, std::memory_order_relaxed );
2850 // if we failed report an error
2851 if( !ret )
2852 {
2853 log->Error( XRootDTransportMsg,
2854 "Unable to get the security framework: %s", errorBuff );
2855 return 0;
2856 }
2857 return ret;
2858 }
2859
2860 //----------------------------------------------------------------------------
2861 // Generate the end session message
2862 //----------------------------------------------------------------------------
2863 Message *XRootDTransport::GenerateEndSession( HandShakeData *hsData,
2864 XRootDChannelInfo *info )
2865 {
2866 Log *log = DefaultEnv::GetLog();
2867
2868 //--------------------------------------------------------------------------
2869 // Generate the message
2870 //--------------------------------------------------------------------------
2871 Message *msg = new Message( sizeof(ClientEndsessRequest) );
2872 ClientEndsessRequest *endsessReq = (ClientEndsessRequest *)msg->GetBuffer();
2873
2874 endsessReq->requestid = kXR_endsess;
2875 memcpy( endsessReq->sessid, info->oldSessionId, 16 );
2876 std::string sessId = Utils::Char2Hex( endsessReq->sessid, 16 );
2877
2878 log->Debug( XRootDTransportMsg, "[%s] Sending out kXR_endsess for session:"
2879 " %s", hsData->streamName.c_str(), sessId.c_str() );
2880
2881 MarshallRequest( msg );
2882
2883 Message *sign = 0;
2884 GetSignature( msg, sign, info );
2885 if( sign )
2886 {
2887 //------------------------------------------------------------------------
2888 // Now place both the signature and the request in a single buffer
2889 //------------------------------------------------------------------------
2890 uint32_t size = sign->GetSize();
2891 sign->ReAllocate( size + msg->GetSize() );
2892 char* buffer = sign->GetBuffer( size );
2893 memcpy( buffer, msg->GetBuffer(), msg->GetSize() );
2894 msg->Grab( sign->GetBuffer(), sign->GetSize() );
2895 }
2896
2897 return msg;
2898 }
2899
2900 //----------------------------------------------------------------------------
2901 // Process the protocol response
2902 //----------------------------------------------------------------------------
2903 Status XRootDTransport::ProcessEndSessionResp( HandShakeData *hsData,
2904 XRootDChannelInfo *info )
2905 {
2906 Log *log = DefaultEnv::GetLog();
2907
2908 Status st = UnMarshallBody( hsData->in, kXR_endsess );
2909 if( !st.IsOK() )
2910 return st;
2911
2912 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2913
2914 // If we're good, we're good!
2915 if( rsp->hdr.status == kXR_ok )
2916 return Status();
2917
2918 // we ignore not found errors as such an error means the connection
2919 // has been already terminated
2920 if( rsp->hdr.status == kXR_error && rsp->body.error.errnum == kXR_NotFound )
2921 return Status();
2922
2923 // other errors
2924 if( rsp->hdr.status == kXR_error )
2925 {
2926 std::string errorMsg( rsp->body.error.errmsg, rsp->hdr.dlen - 4 );
2927 log->Error( XRootDTransportMsg, "[%s] Got error response to "
2928 "kXR_endsess: %s", hsData->streamName.c_str(),
2929 errorMsg.c_str() );
2930 return Status( stFatal, errHandShakeFailed );
2931 }
2932
2933 // Wait Response.
2934 if( rsp->hdr.status == kXR_wait )
2935 {
2936 std::string msg( rsp->body.wait.infomsg, rsp->hdr.dlen - 4 );
2937 log->Info( XRootDTransportMsg, "[%s] Got wait response to "
2938 "kXR_endsess: %s", hsData->streamName.c_str(),
2939 msg.c_str() );
2940 hsData->out = GenerateEndSession( hsData, info );
2941 return Status( stOK, suRetry );
2942 }
2943
2944 // Any other response is protocol violation
2945 return Status( stError, errDataError );
2946 }
2947
2948 //----------------------------------------------------------------------------
2949 // Get a string representation of the server flags
2950 //----------------------------------------------------------------------------
2951 std::string XRootDTransport::ServerFlagsToStr( uint32_t flags )
2952 {
2953 std::string repr = "type: ";
2954 if( flags & kXR_isManager )
2955 repr += "manager ";
2956
2957 else if( flags & kXR_isServer )
2958 repr += "server ";
2959
2960 repr += "[";
2961
2962 if( flags & kXR_attrMeta )
2963 repr += "meta ";
2964
2965 else if( flags & kXR_attrCache )
2966 repr += "cache ";
2967
2968 else if( flags & kXR_attrProxy )
2969 repr += "proxy ";
2970
2971 else if( flags & kXR_attrSuper )
2972 repr += "super ";
2973
2974 else
2975 repr += " ";
2976
2977 repr.erase( repr.length()-1, 1 );
2978
2979 repr += "]";
2980 return repr;
2981 }
2982}
2983
2984namespace
2985{
2986 // Extract file name from a request
2987 //----------------------------------------------------------------------------
2988 char *GetDataAsString( char *msg )
2989 {
2991 char *fn = new char[req->dlen+1];
2992 memcpy( fn, msg + 24, req->dlen );
2993 fn[req->dlen] = 0;
2994 return fn;
2995 }
2996}
2997
2998namespace XrdCl
2999{
3000 //----------------------------------------------------------------------------
3001 // Get the description of a message
3002 //----------------------------------------------------------------------------
3003 void XRootDTransport::GenerateDescription( char *msg, std::ostringstream &o )
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 {
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 {
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 {
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 {
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 {
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 {
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 }
3559
3560 //----------------------------------------------------------------------------
3561 // Get a string representation of file handle
3562 //----------------------------------------------------------------------------
3563 std::string XRootDTransport::FileHandleToStr( const unsigned char handle[4] )
3564 {
3565 std::ostringstream o;
3566 o << "0x";
3567 for( uint8_t i = 0; i < 4; ++i )
3568 {
3569 o << std::setbase(16) << std::setfill('0') << std::setw(2);
3570 o << (int)handle[i];
3571 }
3572 return o.str();
3573 }
3574}
static const int kXR_ckpRollback
Definition XProtocol.hh:215
@ kXR_NotFound
kXR_int16 arg1len
Definition XProtocol.hh:430
#define kXR_isManager
struct ClientTruncateRequest truncate
Definition XProtocol.hh:875
@ kXR_ecredir
Definition XProtocol.hh:371
#define kXR_tlsLogin
@ kXR_fattrDel
Definition XProtocol.hh:270
@ kXR_fattrSet
Definition XProtocol.hh:273
@ kXR_fattrList
Definition XProtocol.hh:272
@ kXR_fattrGet
Definition XProtocol.hh:271
#define kXR_ShortProtRespLen
#define kXR_suppgrw
kXR_char fhandle[4]
Definition XProtocol.hh:531
kXR_unt16 requestid
Definition XProtocol.hh:394
ServerResponseStatus status
kXR_char fhandle[4]
Definition XProtocol.hh:782
#define kXR_gotoTLS
#define kXR_attrMeta
struct ClientPgReadRequest pgread
Definition XProtocol.hh:861
union ServerResponse::@040373375333017131300127053271011057331004327334 body
kXR_char fhandle[4]
Definition XProtocol.hh:807
#define kXR_haveTLS
kXR_char streamid[2]
Definition XProtocol.hh:156
kXR_char fhandle[4]
Definition XProtocol.hh:771
struct ClientMkdirRequest mkdir
Definition XProtocol.hh:858
kXR_int32 dlen
Definition XProtocol.hh:431
struct ClientAuthRequest auth
Definition XProtocol.hh:847
kXR_char streamid[2]
Definition XProtocol.hh:914
kXR_unt16 options
Definition XProtocol.hh:481
static const int kXR_ckpXeq
Definition XProtocol.hh:216
struct ClientPgWriteRequest pgwrite
Definition XProtocol.hh:862
#define kXR_attrSuper
struct ClientReadVRequest readv
Definition XProtocol.hh:868
kXR_char pathid
Definition XProtocol.hh:653
kXR_char credtype[4]
Definition XProtocol.hh:170
kXR_char username[8]
Definition XProtocol.hh:396
@ 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
struct ClientOpenRequest open
Definition XProtocol.hh:860
@ kXR_waitresp
Definition XProtocol.hh:906
@ kXR_redirect
Definition XProtocol.hh:904
@ kXR_status
Definition XProtocol.hh:907
@ kXR_ok
Definition XProtocol.hh:899
@ kXR_authmore
Definition XProtocol.hh:902
@ kXR_attn
Definition XProtocol.hh:901
@ kXR_wait
Definition XProtocol.hh:905
@ kXR_error
Definition XProtocol.hh:903
struct ServerResponseBody_Status bdy
struct ClientRequestHdr header
Definition XProtocol.hh:846
kXR_char fhandle[4]
Definition XProtocol.hh:509
kXR_char fhandle[4]
Definition XProtocol.hh:645
kXR_char fhandle[4]
Definition XProtocol.hh:659
struct ClientWriteVRequest writev
Definition XProtocol.hh:877
kXR_char fhandle[4]
Definition XProtocol.hh:229
struct ClientLoginRequest login
Definition XProtocol.hh:857
kXR_unt16 requestid
Definition XProtocol.hh:157
kXR_char fhandle[4]
Definition XProtocol.hh:633
kXR_char sessid[16]
Definition XProtocol.hh:181
@ 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_bind
Definition XProtocol.hh:136
@ 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_login
Definition XProtocol.hh:119
@ kXR_auth
Definition XProtocol.hh:112
@ kXR_endsess
Definition XProtocol.hh:135
@ kXR_set
Definition XProtocol.hh:130
@ kXR_rmdir
Definition XProtocol.hh:127
@ kXR_1stRequest
Definition XProtocol.hh:111
@ 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
struct ClientChmodRequest chmod
Definition XProtocol.hh:850
#define kXR_isServer
#define kXR_attrCache
struct ClientQueryRequest query
Definition XProtocol.hh:866
struct ClientReadRequest read
Definition XProtocol.hh:867
struct ClientMvRequest mv
Definition XProtocol.hh:859
kXR_int32 rlen
Definition XProtocol.hh:660
kXR_unt16 requestid
Definition XProtocol.hh:180
kXR_char sessid[16]
Definition XProtocol.hh:259
struct ClientChkPointRequest chkpoint
Definition XProtocol.hh:849
struct ServerResponseHeader hdr
@ kXR_asyncap
Definition XProtocol.hh:378
#define kXR_attrProxy
kXR_char options[1]
Definition XProtocol.hh:416
#define kXR_PROTOCOLVERSION
Definition XProtocol.hh:70
static const int kXR_ckpCommit
Definition XProtocol.hh:213
kXR_int64 offset
Definition XProtocol.hh:661
@ kXR_vfs
Definition XProtocol.hh:763
struct ClientPrepareRequest prepare
Definition XProtocol.hh:864
@ 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
#define kXR_tlsSess
#define kXR_DataServer
struct ClientWriteRequest write
Definition XProtocol.hh:876
#define kXR_PROTTLSVERSION
Definition XProtocol.hh:72
kXR_char capver[1]
Definition XProtocol.hh:399
struct ClientProtocolRequest protocol
Definition XProtocol.hh:865
@ 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
struct ClientLocateRequest locate
Definition XProtocol.hh:856
@ kXR_ver005
Definition XProtocol.hh:389
#define kXR_tlsData
@ kXR_readrdok
Definition XProtocol.hh:360
@ kXR_fullurl
Definition XProtocol.hh:358
@ kXR_onlyprv4
Definition XProtocol.hh:362
@ kXR_lclfile
Definition XProtocol.hh:364
@ kXR_multipr
Definition XProtocol.hh:359
@ kXR_redirflags
Definition XProtocol.hh:365
@ kXR_hasipv64
Definition XProtocol.hh:361
@ kXR_onlyprv6
Definition XProtocol.hh:363
ServerResponseHeader hdr
static const int kXR_ckpBegin
Definition XProtocol.hh:212
long long kXR_int64
Definition XPtypes.hh:98
unsigned char kXR_char
Definition XPtypes.hh:65
XrdVERSIONINFOREF(XrdCl)
XrdSecBuffer XrdSecParameters
XrdSecProtocol *(* XrdSecGetProt_t)(const char *hostname, XrdNetAddrInfo &endPoint, XrdSecParameters &sectoken, XrdOucErrInfo *einfo)
Typedef to simplify the encoding of methods returning XrdSecProtocol.
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)
Definition XrdSysE2T.cc:104
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.
void Zero()
Zero.
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)
Definition XrdClEnv.cc:110
bool GetInt(const std::string &key, int &value)
Definition XrdClEnv.cc:89
Handle diagnostics.
Definition XrdClLog.hh:101
@ ErrorMsg
report errors
Definition XrdClLog.hh:109
void Error(uint64_t topic, const char *format,...)
Report an error.
Definition XrdClLog.cc:231
LogLevel GetLevel() const
Get the log level.
Definition XrdClLog.hh:258
void Dump(uint64_t topic, const char *format,...)
Print a dump message.
Definition XrdClLog.cc:299
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
Definition XrdClLog.cc:282
The message representation used throughout the system.
void SetIsMarshalled(bool isMarshalled)
Set the marshalling status.
bool IsMarshalled() const
Check if the message is marshalled.
static SIDMgrPool & Instance()
std::shared_ptr< SIDManager > GetSIDMgr(const URL &url)
A network socket.
virtual XRootDStatus Read(char *buffer, size_t size, int &bytesRead)
static void ClearErrorQueue()
Clear the error queue for the calling thread.
Definition XrdClTls.cc:422
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.
URL representation.
Definition XrdClURL.hh:31
std::string GetChannelId() const
Definition XrdClURL.cc:512
std::map< std::string, std::string > ParamsMap
Definition XrdClURL.hh:33
bool IsSecure() const
Does the protocol indicate encryption.
Definition XrdClURL.cc:482
bool IsTPC() const
Is the URL used in TPC context.
Definition XrdClURL.cc:490
std::string GetLoginToken() const
Get the login token if present in the opaque info.
Definition XrdClURL.cc:367
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.
Definition XrdClUtils.hh:56
const std::string & GetErrorMessage() const
Get error message.
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.
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)
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.
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)
void SetTLS(bool val)
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)
Definition XrdOucCRC.cc:190
static int UserName(uid_t uID, char *uName, int uNsz)
virtual int Secure(SecurityRequest *&newreq, ClientRequest &thereq, const char *thedata)
static int TimeZone()
const uint16_t suRetry
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 uint64_t XRootDMsg
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 suDone
const uint16_t suContinue
bool InitTLS()
Definition XrdClTls.cc:96
const int DefaultTlsNoData
const int DefaultNoTlsOK
const uint16_t errAuthFailed
const uint16_t errInvalidMessage
XrdSysError Log
Definition XrdConfig.cc:113
kXR_char fhandle[4]
Definition XProtocol.hh:832
struct ServerResponseBifs_Protocol bifReqs
struct ServerResponseReqs_Protocol secReqs
kXR_char fhandle[4]
Definition XProtocol.hh:288
BindPrefSelector(std::vector< std::string > &&bindprefs)
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.
static void UnloadHandler(const std::string &trProt)
void Register(const std::string &protocol)
std::set< std::string > protocols
Procedure execution status.
uint16_t code
Error type, or additional hints on what to do.
bool IsOK() const
We're fine.
Selects less loaded stream for read operation over multiple streams.
void AdjustQueues(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
std::unique_ptr< BindPrefSelector > bindSelector
std::atomic< uint32_t > finstcnt
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.