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 //--------------------------------------------------------------------------
127 serverFlags( 0 )
128 {
129 }
130
132 uint8_t pathId;
133 uint32_t serverFlags;
134 };
135
136 //----------------------------------------------------------------------------
138 //----------------------------------------------------------------------------
140 {
141 StreamSelector( uint16_t size )
142 {
143 //----------------------------------------------------------------------
144 // Subtract one because we shouldn't take into account the control
145 // stream.
146 //----------------------------------------------------------------------
147 strmqueues.resize( size - 1, 0 );
148 }
149
150 //------------------------------------------------------------------------
151 // @param size : number of streams
152 //------------------------------------------------------------------------
153 void AdjustQueues( uint16_t size )
154 {
155 strmqueues.resize( size - 1, 0);
156 }
157
158 //------------------------------------------------------------------------
159 // @param connected : bitarray stating if given sub-stream is connected
160 //
161 // @return : substream number
162 //------------------------------------------------------------------------
163 uint16_t Select( const std::vector<bool> &connected )
164 {
165 uint16_t ret = 0;
166 size_t minval = std::numeric_limits<size_t>::max();
167
168 for( size_t i = 0; i < connected.size() && i < strmqueues.size(); ++i )
169 {
170 if( !connected[i] ) continue;
171
172 if( strmqueues[i] < minval )
173 {
174 ret = i;
175 minval = strmqueues[i];
176 }
177 }
178
179 ++strmqueues[ret];
180 return ret + 1;
181 }
182
183 //--------------------------------------------------------------------------
184 // Update queue for given substream
185 //--------------------------------------------------------------------------
186 void MsgReceived( uint16_t substrm )
187 {
188 if( substrm > 0 )
189 --strmqueues[substrm - 1];
190 }
191
192 private:
193
194 std::vector<size_t> strmqueues;
195 };
196
198 {
199 BindPrefSelector( std::vector<std::string> && bindprefs ) :
200 bindprefs( std::move( bindprefs ) ), next( 0 )
201 {
202 }
203
204 inline const std::string& Get()
205 {
206 std::string &ret = bindprefs[next];
207 ++next;
208 if( next >= bindprefs.size() )
209 next = 0;
210 return ret;
211 }
212
213 private:
214 std::vector<std::string> bindprefs;
215 size_t next;
216 };
217
218 //----------------------------------------------------------------------------
220 //----------------------------------------------------------------------------
222 {
223 //--------------------------------------------------------------------------
224 // Constructor
225 //--------------------------------------------------------------------------
226 XRootDChannelInfo( const URL &url ):
227 serverFlags(0),
229 firstLogIn(true),
230 authBuffer(0),
231 authProtocol(0),
232 authParams(0),
233 authEnv(0),
234 finstcnt(0),
235 openFiles(0),
236 waitBarrier(0),
237 protection(0),
238 protRespSize(0),
239 encrypted(false),
240 istpc(false)
241 {
243 memset( sessionId, 0, 16 );
244 memset( oldSessionId, 0, 16 );
245 }
246
247 //--------------------------------------------------------------------------
248 // Destructor
249 //--------------------------------------------------------------------------
251 {
252 delete [] authBuffer;
253 }
254
255 typedef std::vector<XRootDStreamInfo> StreamInfoVector;
256
257 //--------------------------------------------------------------------------
258 // Data
259 //--------------------------------------------------------------------------
260 uint32_t serverFlags;
262 uint8_t sessionId[16];
263 uint8_t oldSessionId[16];
265 std::shared_ptr<SIDManager> sidManager;
271 std::string streamName;
272 std::string authProtocolName;
273 std::set<uint16_t> sentOpens;
274 std::set<uint16_t> sentCloses;
275 std::atomic<uint32_t> finstcnt; // file instance count
276 uint32_t openFiles;
279 std::vector<char> protRespBuff;
280 unsigned int protRespSize;
281 std::unique_ptr<StreamSelector> strmSelector;
283 bool istpc;
284 std::unique_ptr<BindPrefSelector> bindSelector;
285 std::string logintoken;
287 };
288
289 //----------------------------------------------------------------------------
290 // Constructor
291 //----------------------------------------------------------------------------
293 pSecUnloadHandler( new PluginUnloadHandler() )
294 {
295 }
296
297 //----------------------------------------------------------------------------
298 // Destructor
299 //----------------------------------------------------------------------------
301 {
302 delete pSecUnloadHandler; pSecUnloadHandler = 0;
303 }
304
305 //----------------------------------------------------------------------------
306 // Read message header from socket
307 //----------------------------------------------------------------------------
309 {
310 //--------------------------------------------------------------------------
311 // A new message - allocate the space needed for the header
312 //--------------------------------------------------------------------------
313 if( message.GetCursor() == 0 && message.GetSize() < 8 )
314 message.Allocate( 8 );
315
316 //--------------------------------------------------------------------------
317 // Read the message header
318 //--------------------------------------------------------------------------
319 if( message.GetCursor() < 8 )
320 {
321 size_t leftToBeRead = 8 - message.GetCursor();
322 while( leftToBeRead )
323 {
324 int bytesRead = 0;
325 XRootDStatus status = socket->Read( message.GetBufferAtCursor(),
326 leftToBeRead, bytesRead );
327 if( !status.IsOK() || status.code == suRetry )
328 return status;
329
330 leftToBeRead -= bytesRead;
331 message.AdvanceCursor( bytesRead );
332 }
333 UnMarshallHeader( message );
334
335 uint32_t bodySize = *(uint32_t*)(message.GetBuffer(4));
336 Log *log = DefaultEnv::GetLog();
337 log->Dump( XRootDTransportMsg, "[msg: %p] Expecting %d bytes of message "
338 "body", (void*)&message, bodySize );
339
340 return XRootDStatus( stOK, suDone );
341 }
343 }
344
345 //----------------------------------------------------------------------------
346 // Read message body from socket
347 //----------------------------------------------------------------------------
349 {
350 //--------------------------------------------------------------------------
351 // Retrieve the body
352 //--------------------------------------------------------------------------
353 size_t leftToBeRead = 0;
354 uint32_t bodySize = 0;
356 bodySize = rsphdr->dlen;
357
358 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
360 "Response body too large." );
361
362 if( message.GetSize() < bodySize + 8 )
363 message.ReAllocate( bodySize + 8 );
364
365 leftToBeRead = bodySize-(message.GetCursor()-8);
366 while( leftToBeRead )
367 {
368 int bytesRead = 0;
369 XRootDStatus status = socket->Read( message.GetBufferAtCursor(), leftToBeRead, bytesRead );
370
371 if( !status.IsOK() || status.code == suRetry )
372 return status;
373
374 leftToBeRead -= bytesRead;
375 message.AdvanceCursor( bytesRead );
376 }
377
378 return XRootDStatus( stOK, suDone );
379 }
380
381 //----------------------------------------------------------------------------
382 // Read more of the message body from socket
383 //----------------------------------------------------------------------------
385 {
387 if( rsphdr->status != kXR_status )
389
390 //--------------------------------------------------------------------------
391 // In case of non kXR_status responses we read all the response, including
392 // data. For kXR_status responses we first read only the remainder of the
393 // header. The header must then be unmarshalled, and then a second call to
394 // GetMore (repeated for suRetry as needed) will read the data.
395 //--------------------------------------------------------------------------
396
397 uint32_t bodySize = rsphdr->dlen;
398 if( bodySize > std::numeric_limits<uint32_t>::max() - 8 )
400 "kXR_status: response body too large." );
401 if( bodySize+8 < sizeof( ServerResponseStatus ) )
403 "kXR_status: invalid message size." );
404
406 uint32_t moreSize = static_cast<uint32_t>( rspst->bdy.dlen );
407 if( moreSize > std::numeric_limits<uint32_t>::max() - 8 - bodySize )
409 "kXR_status: response body too large." );
410 bodySize += moreSize;
411
412 if( message.GetSize() < bodySize + 8 )
413 message.ReAllocate( bodySize + 8 );
414
415 size_t leftToBeRead = bodySize-(message.GetCursor()-8);
416 while( leftToBeRead )
417 {
418 int bytesRead = 0;
419 XRootDStatus status = socket->Read( message.GetBufferAtCursor(), leftToBeRead, bytesRead );
420
421 if( !status.IsOK() || status.code == suRetry )
422 return status;
423
424 leftToBeRead -= bytesRead;
425 message.AdvanceCursor( bytesRead );
426 }
427
428 // Unmarchal to message body
429 Log *log = DefaultEnv::GetLog();
431 if( !st.IsOK() && st.code == errDataError )
432 {
433 log->Error( XRootDTransportMsg, "[msg: %p] %s", (void*)&message,
434 st.GetErrorMessage().c_str() );
435 return st;
436 }
437
438 if( !st.IsOK() )
439 {
440 log->Error( XRootDTransportMsg, "[msg: %p] Failed to unmarshall status body.",
441 (void*)&message );
442 return st;
443 }
444
445 return XRootDStatus( stOK, suDone );
446 }
447
448 //----------------------------------------------------------------------------
449 // Initialize channel
450 //----------------------------------------------------------------------------
452 AnyObject &channelData )
453 {
454 XRootDChannelInfo *info = new XRootDChannelInfo( url );
455 XrdSysMutexHelper scopedLock( info->mutex );
456 channelData.Set( info );
457
458 Env *env = DefaultEnv::GetEnv();
459 int streams = DefaultSubStreamsPerChannel;
460 env->GetInt( "SubStreamsPerChannel", streams );
461 if( streams < 1 ) streams = 1;
462 info->stream.resize( streams );
463 info->strmSelector.reset( new StreamSelector( streams ) );
464 info->encrypted = url.IsSecure();
465 info->istpc = url.IsTPC();
466 info->logintoken = url.GetLoginToken();
467 }
468
469 //----------------------------------------------------------------------------
470 // Finalize channel
471 //----------------------------------------------------------------------------
475
476 //----------------------------------------------------------------------------
477 // HandShake
478 //----------------------------------------------------------------------------
480 AnyObject &channelData )
481 {
482 XRootDChannelInfo *info = 0;
483 channelData.Get( info );
484
485 if (!info)
487
488 XrdSysMutexHelper scopedLock( info->mutex );
489
490 if( info->stream.size() <= handShakeData->subStreamId )
491 {
492 Log *log = DefaultEnv::GetLog();
494 "[%s] Internal error: not enough substreams",
495 handShakeData->streamName.c_str() );
497 }
498
499 if( handShakeData->subStreamId == 0 )
500 {
501 info->streamName = handShakeData->streamName;
502 return HandShakeMain( handShakeData, channelData );
503 }
504 return HandShakeParallel( handShakeData, channelData );
505 }
506
507 //----------------------------------------------------------------------------
508 // Hand shake the main stream
509 //----------------------------------------------------------------------------
510 XRootDStatus XRootDTransport::HandShakeMain( HandShakeData *handShakeData,
511 AnyObject &channelData )
512 {
513 XRootDChannelInfo *info = 0;
514 channelData.Get( info );
515
516 if (!info) {
518 "[%s] Internal error: no channel info",
519 handShakeData->streamName.c_str());
521 }
522
523 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
524
525 //--------------------------------------------------------------------------
526 // First step - we need to create and initial handshake and send it out
527 //--------------------------------------------------------------------------
528 if( sInfo.status == XRootDStreamInfo::Disconnected ||
529 sInfo.status == XRootDStreamInfo::Broken )
530 {
531 handShakeData->out = GenerateInitialHSProtocol( handShakeData, info,
533 sInfo.status = XRootDStreamInfo::HandShakeSent;
534 return XRootDStatus( stOK, suContinue );
535 }
536
537 //--------------------------------------------------------------------------
538 // Second step - we got the reply message to the initial handshake
539 //--------------------------------------------------------------------------
540 if( sInfo.status == XRootDStreamInfo::HandShakeSent )
541 {
542 XRootDStatus st = ProcessServerHS( handShakeData, info );
543 if( st.IsOK() )
545 else
546 sInfo.status = XRootDStreamInfo::Broken;
547 return st;
548 }
549
550 //--------------------------------------------------------------------------
551 // Third step - we got the response to the protocol request, we need
552 // to process it and send out a login request
553 //--------------------------------------------------------------------------
554 if( sInfo.status == XRootDStreamInfo::HandShakeReceived )
555 {
556 XRootDStatus st = ProcessProtocolResp( handShakeData, info );
557
558 if( !st.IsOK() )
559 {
560 sInfo.status = XRootDStreamInfo::Broken;
561 return st;
562 }
563
564 if( st.code == suRetry )
565 {
566 handShakeData->out = GenerateProtocol( handShakeData, info,
569 return XRootDStatus( stOK, suRetry );
570 }
571
572 handShakeData->out = GenerateLogIn( handShakeData, info );
573 sInfo.status = XRootDStreamInfo::LoginSent;
574 return XRootDStatus( stOK, suContinue );
575 }
576
577 //--------------------------------------------------------------------------
578 // Fourth step - handle the log in response and proceed with the
579 // authentication if required by the server
580 //--------------------------------------------------------------------------
581 if( sInfo.status == XRootDStreamInfo::LoginSent )
582 {
583 XRootDStatus st = ProcessLogInResp( handShakeData, info );
584
585 if( !st.IsOK() )
586 {
587 sInfo.status = XRootDStreamInfo::Broken;
588 return st;
589 }
590
591 if( st.IsOK() && st.code == suDone )
592 {
593 //----------------------------------------------------------------------
594 // If it's not our first log in we need to end the previous session
595 // to make sure that the server noticed our disconnection and closed
596 // all the writable handles that we owned
597 //----------------------------------------------------------------------
598 if( !info->firstLogIn )
599 {
600 handShakeData->out = GenerateEndSession( handShakeData, info );
602 return XRootDStatus( stOK, suContinue );
603 }
604
605 sInfo.status = XRootDStreamInfo::Connected;
606 info->firstLogIn = false;
607 return st;
608 }
609
610 st = DoAuthentication( handShakeData, info );
611 if( !st.IsOK() )
612 sInfo.status = XRootDStreamInfo::Broken;
613 else
614 sInfo.status = XRootDStreamInfo::AuthSent;
615 return st;
616 }
617
618 //--------------------------------------------------------------------------
619 // Fifth step and later - proceed with the authentication
620 //--------------------------------------------------------------------------
621 if( sInfo.status == XRootDStreamInfo::AuthSent )
622 {
623 XRootDStatus st = DoAuthentication( handShakeData, info );
624
625 if( !st.IsOK() )
626 {
627 sInfo.status = XRootDStreamInfo::Broken;
628 return st;
629 }
630
631 if( st.IsOK() && st.code == suDone )
632 {
633 //----------------------------------------------------------------------
634 // If it's not our first log in we need to end the previous session
635 //----------------------------------------------------------------------
636 if( !info->firstLogIn )
637 {
638 handShakeData->out = GenerateEndSession( handShakeData, info );
640 return XRootDStatus( stOK, suContinue );
641 }
642
643 sInfo.status = XRootDStreamInfo::Connected;
644 info->firstLogIn = false;
645 return st;
646 }
647
648 return st;
649 }
650
651 //--------------------------------------------------------------------------
652 // The last step - kXR_endsess returned
653 //--------------------------------------------------------------------------
654 if( sInfo.status == XRootDStreamInfo::EndSessionSent )
655 {
656 XRootDStatus st = ProcessEndSessionResp( handShakeData, info );
657
658 if( st.IsOK() && st.code == suDone )
659 {
660 sInfo.status = XRootDStreamInfo::Connected;
661 }
662 else if( !st.IsOK() )
663 {
664 sInfo.status = XRootDStreamInfo::Broken;
665 }
666
667 return st;
668 }
669
670 return XRootDStatus( stOK, suDone );
671 }
672
673 //----------------------------------------------------------------------------
674 // Hand shake parallel stream
675 //----------------------------------------------------------------------------
676 XRootDStatus XRootDTransport::HandShakeParallel( HandShakeData *handShakeData,
677 AnyObject &channelData )
678 {
679 XRootDChannelInfo *info = 0;
680 channelData.Get( info );
681
682 if (!info) {
684 "[%s] Internal error: no channel info",
685 handShakeData->streamName.c_str());
687 }
688
689 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
690
691 //--------------------------------------------------------------------------
692 // First step - we need to create and initial handshake and send it out
693 //--------------------------------------------------------------------------
694 if( sInfo.status == XRootDStreamInfo::Disconnected ||
695 sInfo.status == XRootDStreamInfo::Broken )
696 {
697 handShakeData->out = GenerateInitialHSProtocol( handShakeData, info,
699 sInfo.status = XRootDStreamInfo::HandShakeSent;
700 return XRootDStatus( stOK, suContinue );
701 }
702
703 //--------------------------------------------------------------------------
704 // Second step - we got the reply message to the initial handshake,
705 // if successful we need to send bind
706 //--------------------------------------------------------------------------
707 if( sInfo.status == XRootDStreamInfo::HandShakeSent )
708 {
709 XRootDStatus st = ProcessServerHS( handShakeData, info );
710 if( st.IsOK() )
712 else
713 sInfo.status = XRootDStreamInfo::Broken;
714 return st;
715 }
716
717 //--------------------------------------------------------------------------
718 // Second step bis - we got the response to the protocol request, we need
719 // to process it and send out a bind request
720 //--------------------------------------------------------------------------
721 if( sInfo.status == XRootDStreamInfo::HandShakeReceived )
722 {
723 XRootDStatus st = ProcessProtocolResp( handShakeData, info );
724
725 if( !st.IsOK() )
726 {
727 sInfo.status = XRootDStreamInfo::Broken;
728 return st;
729 }
730
731 handShakeData->out = GenerateBind( handShakeData, info );
732 sInfo.status = XRootDStreamInfo::BindSent;
733 return XRootDStatus( stOK, suContinue );
734 }
735
736 //--------------------------------------------------------------------------
737 // Third step - we got the response to the kXR_bind
738 //--------------------------------------------------------------------------
739 if( sInfo.status == XRootDStreamInfo::BindSent )
740 {
741 XRootDStatus st = ProcessBindResp( handShakeData, info );
742
743 if( !st.IsOK() )
744 {
745 sInfo.status = XRootDStreamInfo::Broken;
746 return st;
747 }
748 sInfo.status = XRootDStreamInfo::Connected;
749 return XRootDStatus();
750 }
751 return XRootDStatus();
752 }
753
754 //------------------------------------------------------------------------
755 // @return true if handshake has been done and stream is connected,
756 // false otherwise
757 //------------------------------------------------------------------------
759 AnyObject &channelData )
760 {
761 XRootDChannelInfo *info = 0;
762 channelData.Get( info );
763
764 if (!info) {
766 "[%s] Internal error: no channel info",
767 handShakeData->streamName.c_str());
768 return false;
769 }
770
771 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
772 return ( sInfo.status == XRootDStreamInfo::Connected );
773 }
774
775 //----------------------------------------------------------------------------
776 // Check if the stream should be disconnected
777 //----------------------------------------------------------------------------
778 bool XRootDTransport::IsStreamTTLElapsed( time_t inactiveTime,
779 AnyObject &channelData )
780 {
781 XRootDChannelInfo *info = 0;
782 channelData.Get( info );
783
784 Env *env = DefaultEnv::GetEnv();
785 Log *log = DefaultEnv::GetLog();
786
787 if (!info) {
789 "Internal error: no channel info, behaving as if TTL has elapsed");
790 return true;
791 }
792
793 //--------------------------------------------------------------------------
794 // Check the TTL settings for the current server
795 //--------------------------------------------------------------------------
796 int ttl;
797 if( info->serverFlags & kXR_isServer )
798 {
800 env->GetInt( "DataServerTTL", ttl );
801 }
802 else
803 {
805 env->GetInt( "LoadBalancerTTL", ttl );
806 }
807
808 //--------------------------------------------------------------------------
809 // See whether we can give a go-ahead for the disconnection
810 //--------------------------------------------------------------------------
811 XrdSysMutexHelper scopedLock( info->mutex );
812 uint16_t allocatedSIDs = info->sidManager->GetNumberOfAllocatedSIDs();
813 log->Dump( XRootDTransportMsg, "[%s] Stream inactive since %lld seconds, "
814 "TTL: %d, allocated SIDs: %d, open files: %d, bound file objects: %d",
815 info->streamName.c_str(), (long long) inactiveTime, ttl, allocatedSIDs,
816 info->openFiles, info->finstcnt.load( std::memory_order_relaxed ) );
817
818 if( info->openFiles != 0 && info->finstcnt.load( std::memory_order_relaxed ) != 0 )
819 return false;
820
821 if( !allocatedSIDs && inactiveTime > ttl )
822 return true;
823
824 return false;
825 }
826
827 //----------------------------------------------------------------------------
828 // Check the stream is broken - ie. TCP connection got broken and
829 // went undetected by the TCP stack
830 //----------------------------------------------------------------------------
832 AnyObject &channelData )
833 {
834 XRootDChannelInfo *info = 0;
835 channelData.Get( info );
836 Env *env = DefaultEnv::GetEnv();
837 Log *log = DefaultEnv::GetLog();
838
839 if (!info) {
841 "Internal error: no channel info, behaving as if stream is broken");
842 return true;
843 }
844
845 int streamTimeout = DefaultStreamTimeout;
846 env->GetInt( "StreamTimeout", streamTimeout );
847
848 XrdSysMutexHelper scopedLock( info->mutex );
849
850 const time_t now = time(0);
851 const bool anySID =
852 info->sidManager->IsAnySIDOldAs( now - streamTimeout );
853
854 log->Dump( XRootDTransportMsg, "[%s] Stream inactive since %lld seconds, "
855 "stream timeout: %d, any SID: %d, wait barrier: %s",
856 info->streamName.c_str(), (long long) inactiveTime, streamTimeout,
857 anySID, Utils::TimeToString(info->waitBarrier).c_str() );
858
859 if( inactiveTime < streamTimeout )
860 return Status();
861
862 if( now < info->waitBarrier )
863 return Status();
864
865 if( !anySID )
866 return Status();
867
869 }
870
871 //----------------------------------------------------------------------------
872 // Multiplex
873 //----------------------------------------------------------------------------
875 {
876 return PathID( 0, 0 );
877 }
878
879 //----------------------------------------------------------------------------
880 // Multiplex
881 //----------------------------------------------------------------------------
883 AnyObject &channelData,
884 PathID *hint )
885 {
886 XRootDChannelInfo *info = 0;
887 channelData.Get( info );
888
889 if (!info) {
891 "Internal error: no channel info, cannot multiplex");
892 return PathID(0,0);
893 }
894
895 XrdSysMutexHelper scopedLock( info->mutex );
896
897 //--------------------------------------------------------------------------
898 // If we're not connected to a data server or we don't know that yet
899 // we stream through 0
900 //--------------------------------------------------------------------------
901 if( !(info->serverFlags & kXR_isServer) || info->stream.size() == 0 )
902 return PathID( 0, 0 );
903
904 //--------------------------------------------------------------------------
905 // Select the streams
906 //--------------------------------------------------------------------------
907 Log *log = DefaultEnv::GetLog();
908 uint16_t upStream = 0;
909 uint16_t downStream = 0;
910
911 if( hint )
912 {
913 upStream = hint->up;
914 downStream = hint->down;
915 }
916 else
917 {
918 upStream = 0;
919 std::vector<bool> connected;
920 connected.reserve( info->stream.size() - 1 );
921 size_t nbConnected = 0;
922 for( size_t i = 1; i < info->stream.size(); ++i )
923 if( info->stream[i].status == XRootDStreamInfo::Connected )
924 {
925 connected.push_back( true );
926 ++nbConnected;
927 }
928 else
929 connected.push_back( false );
930
931 if( nbConnected == 0 )
932 downStream = 0;
933 else
934 downStream = info->strmSelector->Select( connected );
935 }
936
937 if( upStream >= info->stream.size() )
938 {
940 "[%s] Up link stream %d does not exist, using 0",
941 info->streamName.c_str(), upStream );
942 upStream = 0;
943 }
944
945 if( downStream >= info->stream.size() )
946 {
948 "[%s] Down link stream %d does not exist, using 0",
949 info->streamName.c_str(), downStream );
950 downStream = 0;
951 }
952
953 //--------------------------------------------------------------------------
954 // Modify the message
955 //--------------------------------------------------------------------------
956 UnMarshallRequest( msg );
958 switch( hdr->requestid )
959 {
960 //------------------------------------------------------------------------
961 // Read - we update the path id to tell the server where we want to
962 // get the response, but we still send the request through stream 0
963 // We need to allocate space for read_args if we don't have it
964 // included yet
965 //------------------------------------------------------------------------
966 case kXR_read:
967 {
968 if( msg->GetSize() < sizeof(ClientReadRequest) + 8 )
969 {
970 msg->ReAllocate( sizeof(ClientReadRequest) + 8 );
971 void *newBuf = msg->GetBuffer(sizeof(ClientReadRequest));
972 memset( newBuf, 0, 8 );
974 req->dlen += 8;
975 }
976 read_args *args = (read_args*)msg->GetBuffer(sizeof(ClientReadRequest));
977 args->pathid = info->stream[downStream].pathId;
978 break;
979 }
980
981
982 //------------------------------------------------------------------------
983 // PgRead - we update the path id to tell the server where we want to
984 // get the response, but we still send the request through stream 0
985 // We need to allocate space for ClientPgReadReqArgs if we don't have it
986 // included yet
987 //------------------------------------------------------------------------
988 case kXR_pgread:
989 {
990 if( msg->GetSize() < sizeof( ClientPgReadRequest ) + sizeof( ClientPgReadReqArgs ) )
991 {
992 msg->ReAllocate( sizeof( ClientPgReadRequest ) + sizeof( ClientPgReadReqArgs ) );
993 void *newBuf = msg->GetBuffer( sizeof( ClientPgReadRequest ) );
994 memset( newBuf, 0, sizeof( ClientPgReadReqArgs ) );
996 req->dlen += sizeof( ClientPgReadReqArgs );
997 }
998 ClientPgReadReqArgs *args = reinterpret_cast<ClientPgReadReqArgs*>(
999 msg->GetBuffer( sizeof( ClientPgReadRequest ) ) );
1000 args->pathid = info->stream[downStream].pathId;
1001 break;
1002 }
1003
1004 //------------------------------------------------------------------------
1005 // ReadV - the situation is identical to read but we don't need any
1006 // additional structures to specify the return path
1007 //------------------------------------------------------------------------
1008 case kXR_readv:
1009 {
1011 req->pathid = info->stream[downStream].pathId;
1012 break;
1013 }
1014
1015 //------------------------------------------------------------------------
1016 // Write - multiplexing writes doesn't work properly in the server
1017 //------------------------------------------------------------------------
1018 case kXR_write:
1019 {
1020// ClientWriteRequest *req = (ClientWriteRequest*)msg->GetBuffer();
1021// req->pathid = info->stream[downStream].pathId;
1022 break;
1023 }
1024
1025 //------------------------------------------------------------------------
1026 // WriteV - multiplexing writes doesn't work properly in the server
1027 //------------------------------------------------------------------------
1028 case kXR_writev:
1029 {
1030// ClientWriteVRequest *req = (ClientWriteVRequest*)msg->GetBuffer();
1031// req->pathid = info->stream[downStream].pathId;
1032 break;
1033 }
1034
1035 //------------------------------------------------------------------------
1036 // PgWrite - multiplexing writes doesn't work properly in the server
1037 //------------------------------------------------------------------------
1038 case kXR_pgwrite:
1039 {
1040// ClientWriteVRequest *req = (ClientWriteVRequest*)msg->GetBuffer();
1041// req->pathid = info->stream[downStream].pathId;
1042 break;
1043 }
1044 };
1045 MarshallRequest( msg );
1046 return PathID( upStream, downStream );
1047 }
1048
1049 //----------------------------------------------------------------------------
1050 // Return a number of substreams per stream that should be created
1051 // This depends on the environment and whether we are connected to
1052 // a data server or not
1053 //----------------------------------------------------------------------------
1055 {
1056 XRootDChannelInfo *info = 0;
1057 channelData.Get( info );
1058
1059 if (!info) {
1060 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1061 return 1;
1062 }
1063
1064 XrdSysMutexHelper scopedLock( info->mutex );
1065
1066 //--------------------------------------------------------------------------
1067 // If the connection has been opened in order to orchestrate a TPC or
1068 // the remote server is a Manager or Metamanager we will need only one
1069 // (control) stream.
1070 //--------------------------------------------------------------------------
1071 if( info->istpc || !(info->serverFlags & kXR_isServer ) ) return 1;
1072
1073 //--------------------------------------------------------------------------
1074 // Number of streams requested by user
1075 //--------------------------------------------------------------------------
1076 uint16_t ret = info->stream.size();
1077
1079 int nodata = DefaultTlsNoData;
1080 env->GetInt( "TlsNoData", nodata );
1081
1082 // Does the server require the stream 0 to be encrypted?
1083 bool srvTlsStrm0 = ( info->serverFlags & kXR_gotoTLS ) ||
1084 ( info->serverFlags & kXR_tlsLogin ) ||
1085 ( info->serverFlags & kXR_tlsSess );
1086 // Does the server NOT require the data streams to be encrypted?
1087 bool srvNoTlsData = !( info->serverFlags & kXR_tlsData );
1088 // Does the user require the stream 0 to be encrypted?
1089 bool usrTlsStrm0 = info->encrypted;
1090 // Does the user NOT require the data streams to be encrypted?
1091 bool usrNoTlsData = !info->encrypted || ( info->encrypted && nodata );
1092
1093 if( ( usrTlsStrm0 && usrNoTlsData && srvNoTlsData ) ||
1094 ( srvTlsStrm0 && srvNoTlsData && usrNoTlsData ) )
1095 {
1096 //------------------------------------------------------------------------
1097 // The server or user asked us to encrypt stream 0, but to send the data
1098 // (read/write) using a plain TCP connection
1099 //------------------------------------------------------------------------
1100 if( ret == 1 ) ++ret;
1101 }
1102
1103 if( ret > info->stream.size() )
1104 {
1105 info->stream.resize( ret );
1106 info->strmSelector->AdjustQueues( ret );
1107 }
1108
1109 return ret;
1110 }
1111
1112 //----------------------------------------------------------------------------
1113 // Marshall
1114 //----------------------------------------------------------------------------
1116 {
1117 ClientRequest *req = (ClientRequest*)msg;
1118 switch( req->header.requestid )
1119 {
1120 //------------------------------------------------------------------------
1121 // kXR_protocol
1122 //------------------------------------------------------------------------
1123 case kXR_protocol:
1124 req->protocol.clientpv = htonl( req->protocol.clientpv );
1125 break;
1126
1127 //------------------------------------------------------------------------
1128 // kXR_login
1129 //------------------------------------------------------------------------
1130 case kXR_login:
1131 req->login.pid = htonl( req->login.pid );
1132 break;
1133
1134 //------------------------------------------------------------------------
1135 // kXR_locate
1136 //------------------------------------------------------------------------
1137 case kXR_locate:
1138 req->locate.options = htons( req->locate.options );
1139 break;
1140
1141 //------------------------------------------------------------------------
1142 // kXR_query
1143 //------------------------------------------------------------------------
1144 case kXR_query:
1145 req->query.infotype = htons( req->query.infotype );
1146 break;
1147
1148 //------------------------------------------------------------------------
1149 // kXR_truncate
1150 //------------------------------------------------------------------------
1151 case kXR_truncate:
1152 req->truncate.offset = htonll( req->truncate.offset );
1153 break;
1154
1155 //------------------------------------------------------------------------
1156 // kXR_mkdir
1157 //------------------------------------------------------------------------
1158 case kXR_mkdir:
1159 req->mkdir.mode = htons( req->mkdir.mode );
1160 break;
1161
1162 //------------------------------------------------------------------------
1163 // kXR_chmod
1164 //------------------------------------------------------------------------
1165 case kXR_chmod:
1166 req->chmod.mode = htons( req->chmod.mode );
1167 break;
1168
1169 //------------------------------------------------------------------------
1170 // kXR_open
1171 //------------------------------------------------------------------------
1172 case kXR_open:
1173 req->open.mode = htons( req->open.mode );
1174 req->open.options = htons( req->open.options );
1175 req->open.optiont = htons( req->open.optiont );
1176 break;
1177
1178 //------------------------------------------------------------------------
1179 // kXR_read
1180 //------------------------------------------------------------------------
1181 case kXR_read:
1182 req->read.offset = htonll( req->read.offset );
1183 req->read.rlen = htonl( req->read.rlen );
1184 break;
1185
1186 //------------------------------------------------------------------------
1187 // kXR_write
1188 //------------------------------------------------------------------------
1189 case kXR_write:
1190 req->write.offset = htonll( req->write.offset );
1191 break;
1192
1193 //------------------------------------------------------------------------
1194 // kXR_mv
1195 //------------------------------------------------------------------------
1196 case kXR_mv:
1197 req->mv.arg1len = htons( req->mv.arg1len );
1198 break;
1199
1200 //------------------------------------------------------------------------
1201 // kXR_readv
1202 //------------------------------------------------------------------------
1203 case kXR_readv:
1204 {
1205 uint16_t numChunks = (req->readv.dlen)/16;
1206 readahead_list *dataChunk = (readahead_list*)( msg + 24 );
1207 for( size_t i = 0; i < numChunks; ++i )
1208 {
1209 dataChunk[i].rlen = htonl( dataChunk[i].rlen );
1210 dataChunk[i].offset = htonll( dataChunk[i].offset );
1211 }
1212 break;
1213 }
1214
1215 case kXR_clone:
1216 {
1217 uint32_t numChunks = (req->clone.dlen)/sizeof(XrdProto::clone_list);
1218 XrdProto::clone_list *dataChunk =
1219 (XrdProto::clone_list*)( msg + sizeof( ClientRequestHdr ) );
1220 for( size_t i = 0; i < numChunks; ++i )
1221 {
1222 dataChunk[i].srcOffs = htonll( dataChunk[i].srcOffs );
1223 dataChunk[i].srcLen = htonll( dataChunk[i].srcLen );
1224 dataChunk[i].dstOffs = htonll( dataChunk[i].dstOffs );
1225 }
1226 break;
1227 }
1228
1229 //------------------------------------------------------------------------
1230 // kXR_writev
1231 //------------------------------------------------------------------------
1232 case kXR_writev:
1233 {
1234 uint16_t numChunks = (req->writev.dlen)/16;
1235 XrdProto::write_list *wrtList =
1236 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
1237 for( size_t i = 0; i < numChunks; ++i )
1238 {
1239 wrtList[i].wlen = htonl( wrtList[i].wlen );
1240 wrtList[i].offset = htonll( wrtList[i].offset );
1241 }
1242
1243 break;
1244 }
1245
1246 case kXR_pgread:
1247 {
1248 req->pgread.offset = htonll( req->pgread.offset );
1249 req->pgread.rlen = htonl( req->pgread.rlen );
1250 break;
1251 }
1252
1253 case kXR_pgwrite:
1254 {
1255 req->pgwrite.offset = htonll( req->pgwrite.offset );
1256 break;
1257 }
1258
1259 //------------------------------------------------------------------------
1260 // kXR_prepare
1261 //------------------------------------------------------------------------
1262 case kXR_prepare:
1263 {
1264 req->prepare.optionX = htons( req->prepare.optionX );
1265 req->prepare.port = htons( req->prepare.port );
1266 break;
1267 }
1268
1269 case kXR_chkpoint:
1270 {
1271 if( req->chkpoint.opcode == kXR_ckpXeq )
1272 MarshallRequest( msg + 24 );
1273 break;
1274 }
1275 };
1276
1277 req->header.requestid = htons( req->header.requestid );
1278 req->header.dlen = htonl( req->header.dlen );
1279 return XRootDStatus();
1280 }
1281
1282 //----------------------------------------------------------------------------
1283 // Unmarshall the request - sometimes the requests need to be rewritten,
1284 // so we need to unmarshall them
1285 //----------------------------------------------------------------------------
1287 {
1288 if( !msg->IsMarshalled() ) return XRootDStatus( stOK, suAlreadyDone );
1289 // We rely on the marshaling process to be symmetric!
1290 // First we unmarshall the request ID and the length because
1291 // MarshallRequest() relies on these, and then we need to unmarshall these
1292 // two again, because they get marshalled in MarshallRequest().
1293 // All this is pretty damn ugly and should be rewritten.
1294 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1295 req->header.requestid = htons( req->header.requestid );
1296 req->header.dlen = htonl( req->header.dlen );
1297 XRootDStatus st = MarshallRequest( msg );
1298 req->header.requestid = htons( req->header.requestid );
1299 req->header.dlen = htonl( req->header.dlen );
1300 msg->SetIsMarshalled( false );
1301 return st;
1302 }
1303
1304 //----------------------------------------------------------------------------
1305 // Unmarshall the body of the incoming message
1306 //----------------------------------------------------------------------------
1308 {
1310
1311 //--------------------------------------------------------------------------
1312 // kXR_ok
1313 //--------------------------------------------------------------------------
1314 if( m->hdr.status == kXR_ok )
1315 {
1316 switch( reqType )
1317 {
1318 //----------------------------------------------------------------------
1319 // kXR_protocol
1320 //----------------------------------------------------------------------
1321 case kXR_protocol:
1322 if( m->hdr.dlen < 8 )
1323 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_protocol: body too short." );
1324 m->body.protocol.pval = ntohl( m->body.protocol.pval );
1325 m->body.protocol.flags = ntohl( m->body.protocol.flags );
1326 break;
1327 }
1328 }
1329 //--------------------------------------------------------------------------
1330 // kXR_error
1331 //--------------------------------------------------------------------------
1332 else if( m->hdr.status == kXR_error )
1333 {
1334 if( m->hdr.dlen < 4 )
1335 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_error: body too short." );
1336 m->body.error.errnum = ntohl( m->body.error.errnum );
1337 }
1338
1339 //--------------------------------------------------------------------------
1340 // kXR_wait
1341 //--------------------------------------------------------------------------
1342 else if( m->hdr.status == kXR_wait )
1343 {
1344 if( m->hdr.dlen < 4 )
1345 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_wait: body too short." );
1346 m->body.wait.seconds = htonl( m->body.wait.seconds );
1347 }
1348
1349 //--------------------------------------------------------------------------
1350 // kXR_redirect
1351 //--------------------------------------------------------------------------
1352 else if( m->hdr.status == kXR_redirect )
1353 {
1354 if( m->hdr.dlen < 4 )
1355 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_redirect: body too short." );
1356 m->body.redirect.port = htonl( m->body.redirect.port );
1357 }
1358
1359 //--------------------------------------------------------------------------
1360 // kXR_waitresp
1361 //--------------------------------------------------------------------------
1362 else if( m->hdr.status == kXR_waitresp )
1363 {
1364 if( m->hdr.dlen < 4 )
1365 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_waitresp: body too short." );
1366 m->body.waitresp.seconds = htonl( m->body.waitresp.seconds );
1367 }
1368
1369 //--------------------------------------------------------------------------
1370 // kXR_attn
1371 //--------------------------------------------------------------------------
1372 else if( m->hdr.status == kXR_attn )
1373 {
1374 if( m->hdr.dlen < 4 )
1375 return XRootDStatus( stError, errInvalidMessage, 0, "kXR_attn: body too short." );
1376 m->body.attn.actnum = htonl( m->body.attn.actnum );
1377 }
1378
1379 return XRootDStatus();
1380 }
1381
1382 //------------------------------------------------------------------------
1384 //------------------------------------------------------------------------
1386 {
1387 //--------------------------------------------------------------------------
1388 // Calculate the crc32c before the unmarshaling the body!
1389 //--------------------------------------------------------------------------
1391 char *buffer = msg.GetBuffer( 8 + sizeof( rspst->bdy.crc32c ) );
1392 size_t length = rspst->hdr.dlen - sizeof( rspst->bdy.crc32c );
1393 uint32_t crcval = XrdOucCRC::Calc32C( buffer, length );
1394
1395 size_t stlen = sizeof( ServerResponseStatus );
1396 switch( reqType )
1397 {
1398 case kXR_pgread:
1399 {
1400 stlen += sizeof( ServerResponseBody_pgRead );
1401 break;
1402 }
1403
1404 case kXR_pgwrite:
1405 {
1406 stlen += sizeof( ServerResponseBody_pgWrite );
1407 break;
1408 }
1409 }
1410
1411 if( msg.GetSize() < stlen ) return XRootDStatus( stError, errInvalidMessage, 0,
1412 "kXR_status: invalid message size." );
1413
1414 rspst->bdy.crc32c = ntohl( rspst->bdy.crc32c );
1415 rspst->bdy.dlen = ntohl( rspst->bdy.dlen );
1416
1417 switch( reqType )
1418 {
1419 case kXR_pgread:
1420 {
1422 pgrdbdy->offset = ntohll( pgrdbdy->offset );
1423 break;
1424 }
1425
1426 case kXR_pgwrite:
1427 {
1429 pgwrtbdy->offset = ntohll( pgwrtbdy->offset );
1430 break;
1431 }
1432 }
1433
1434 //--------------------------------------------------------------------------
1435 // Do the integrity checks
1436 //--------------------------------------------------------------------------
1437 if( crcval != rspst->bdy.crc32c )
1438 {
1439 return XRootDStatus( stError, errDataError, 0, "kXR_status response header "
1440 "corrupted (crc32c integrity check failed)." );
1441 }
1442
1443 if( rspst->hdr.streamid[0] != rspst->bdy.streamID[0] ||
1444 rspst->hdr.streamid[1] != rspst->bdy.streamID[1] )
1445 {
1446 return XRootDStatus( stError, errDataError, 0, "response header corrupted "
1447 "(stream ID mismatch)." );
1448 }
1449
1450
1451
1452 if( rspst->bdy.requestid + kXR_1stRequest != reqType )
1453 {
1454 return XRootDStatus( stError, errDataError, 0, "kXR_status response header corrupted "
1455 "(request ID mismatch)." );
1456 }
1457
1458 return XRootDStatus();
1459 }
1460
1462 {
1464 uint16_t reqType = rsp->status.bdy.requestid + kXR_1stRequest;
1465
1466 switch( reqType )
1467 {
1468 case kXR_pgwrite:
1469 {
1470 //--------------------------------------------------------------------------
1471 // If there's no additional data there's nothing to unmarshal
1472 //--------------------------------------------------------------------------
1473 if( rsp->status.bdy.dlen == 0 ) return XRootDStatus();
1474 //--------------------------------------------------------------------------
1475 // If there's not enough data to form correction-segment report an error
1476 //--------------------------------------------------------------------------
1477 if( size_t( rsp->status.bdy.dlen ) < sizeof( ServerResponseBody_pgWrCSE ) )
1479 "kXR_status: invalid message size." );
1480
1481 //--------------------------------------------------------------------------
1482 // Calculate the crc32c for the additional data
1483 //--------------------------------------------------------------------------
1485 cse->cseCRC = ntohl( cse->cseCRC );
1486 size_t length = rsp->status.bdy.dlen - sizeof( uint32_t );
1487 void* buffer = msg.GetBuffer( sizeof( ServerResponseV2 ) + sizeof( uint32_t ) );
1488 uint32_t crcval = XrdOucCRC::Calc32C( buffer, length );
1489
1490 //--------------------------------------------------------------------------
1491 // Do the integrity checks
1492 //--------------------------------------------------------------------------
1493 if( crcval != cse->cseCRC )
1494 {
1495 return XRootDStatus( stError, errDataError, 0, "kXR_status response header "
1496 "corrupted (crc32c integrity check failed)." );
1497 }
1498
1499 cse->dlFirst = ntohs( cse->dlFirst );
1500 cse->dlLast = ntohs( cse->dlLast );
1501
1502 size_t pgcnt = ( rsp->status.bdy.dlen - sizeof( ServerResponseBody_pgWrCSE ) ) /
1503 sizeof( kXR_int64 );
1504 kXR_int64 *pgoffs = (kXR_int64*)msg.GetBuffer( sizeof( ServerResponseV2 ) +
1505 sizeof( ServerResponseBody_pgWrCSE ) );
1506
1507 for( size_t i = 0; i < pgcnt; ++i )
1508 pgoffs[i] = ntohll( pgoffs[i] );
1509
1510 return XRootDStatus();
1511 break;
1512 }
1513
1514 default:
1515 break;
1516 }
1517
1519 }
1520
1521 //----------------------------------------------------------------------------
1522 // Unmarshall the header of the incoming message
1523 //----------------------------------------------------------------------------
1525 {
1527 header->status = ntohs( header->status );
1528 header->dlen = ntohl( header->dlen );
1529 }
1530
1531 //----------------------------------------------------------------------------
1532 // Log server error response
1533 //----------------------------------------------------------------------------
1535 {
1536 Log *log = DefaultEnv::GetLog();
1537 ServerResponse *rsp = (ServerResponse *)msg.GetBuffer();
1538 char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
1539 memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
1540 log->Error( XRootDTransportMsg, "Server responded with an error [%d]: %s",
1541 rsp->body.error.errnum, errmsg );
1542 delete [] errmsg;
1543 }
1544
1545 //------------------------------------------------------------------------
1546 // Number of currently connected data streams
1547 //------------------------------------------------------------------------
1549 {
1550 XRootDChannelInfo *info = 0;
1551 channelData.Get( info );
1552
1553 if (!info) {
1554 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1555 return 0;
1556 }
1557
1558 XrdSysMutexHelper scopedLock( info->mutex );
1559
1560 uint16_t nbConnected = 0;
1561 for( size_t i = 1; i < info->stream.size(); ++i )
1562 if( info->stream[i].status == XRootDStreamInfo::Connected )
1563 ++nbConnected;
1564
1565 return nbConnected;
1566 }
1567
1568 //----------------------------------------------------------------------------
1569 // The stream has been disconnected, do the cleanups
1570 //----------------------------------------------------------------------------
1572 uint16_t subStreamId )
1573 {
1574 XRootDChannelInfo *info = 0;
1575 channelData.Get( info );
1576
1577 if (!info) {
1578 DefaultEnv::GetLog()->Error(XRootDTransportMsg, "Internal error: no channel info");
1579 return;
1580 }
1581
1582 XrdSysMutexHelper scopedLock( info->mutex );
1583
1584 if( !info->stream.empty() )
1585 {
1586 XRootDStreamInfo &sInfo = info->stream[subStreamId];
1588 }
1589
1590 if( subStreamId == 0 )
1591 {
1592 CleanUpProtection( info );
1593 info->sidManager->ReleaseAllTimedOut();
1594 info->sentOpens.clear();
1595 info->sentCloses.clear();
1596 info->openFiles = 0;
1597 info->waitBarrier = 0;
1598 }
1599 }
1600
1601 //------------------------------------------------------------------------
1602 // Query the channel
1603 //------------------------------------------------------------------------
1605 AnyObject &result,
1606 AnyObject &channelData )
1607 {
1608 XRootDChannelInfo *info = 0;
1609 channelData.Get( info );
1610
1611 if (!info)
1613
1614 XrdSysMutexHelper scopedLock( info->mutex );
1615
1616 switch( query )
1617 {
1618 //------------------------------------------------------------------------
1619 // Protocol name
1620 //------------------------------------------------------------------------
1622 result.Set( (const char*)"XRootD", false );
1623 return Status();
1624
1625 //------------------------------------------------------------------------
1626 // Authentication
1627 //------------------------------------------------------------------------
1629 result.Set( new std::string( info->authProtocolName ), false );
1630 return Status();
1631
1632 //------------------------------------------------------------------------
1633 // Server flags
1634 //------------------------------------------------------------------------
1636 result.Set( new int( info->serverFlags ), false );
1637 return Status();
1638
1639 //------------------------------------------------------------------------
1640 // Protocol version
1641 //------------------------------------------------------------------------
1643 result.Set( new int( info->protocolVersion ), false );
1644 return Status();
1645
1647 result.Set( new bool( info->encrypted ), false );
1648 return Status();
1649 };
1651 }
1652
1653 //----------------------------------------------------------------------------
1654 // Check whether the transport can hijack the message
1655 //----------------------------------------------------------------------------
1657 uint16_t subStream,
1658 AnyObject &channelData )
1659 {
1660 XRootDChannelInfo *info = 0;
1661 channelData.Get( info );
1662 if( !info ) return NoAction;
1663 XrdSysMutexHelper scopedLock( info->mutex );
1664 Log *log = DefaultEnv::GetLog();
1665
1666 //--------------------------------------------------------------------------
1667 // Update the substream queues
1668 //--------------------------------------------------------------------------
1669 info->strmSelector->MsgReceived( subStream );
1670
1671 //--------------------------------------------------------------------------
1672 // Check whether this message is a response to a request that has
1673 // timed out, and if so, drop it
1674 //--------------------------------------------------------------------------
1676 if( rsp->hdr.status == kXR_attn )
1677 {
1678 return NoAction;
1679 }
1680
1681 if( info->sidManager->IsTimedOut( rsp->hdr.streamid ) )
1682 {
1683 log->Error( XRootDTransportMsg, "Message %p, stream [%d, %d] is a "
1684 "response that we're no longer interested in (timed out)",
1685 (void*)&msg, rsp->hdr.streamid[0], rsp->hdr.streamid[1] );
1686 //------------------------------------------------------------------------
1687 // If it is kXR_waitresp there will be another one,
1688 // so we don't release the sid yet
1689 //------------------------------------------------------------------------
1690 if( rsp->hdr.status != kXR_waitresp )
1691 info->sidManager->ReleaseTimedOut( rsp->hdr.streamid );
1692 //------------------------------------------------------------------------
1693 // If it is a successful response to an open request
1694 // that timed out, we need to send a close
1695 //------------------------------------------------------------------------
1696 uint16_t sid; memcpy( &sid, rsp->hdr.streamid, 2 );
1697 std::set<uint16_t>::iterator sidIt = info->sentOpens.find( sid );
1698 if( sidIt != info->sentOpens.end() )
1699 {
1700 info->sentOpens.erase( sidIt );
1701 if( rsp->hdr.status == kXR_ok ) return RequestClose;
1702 }
1703 return DigestMsg;
1704 }
1705
1706 //--------------------------------------------------------------------------
1707 // If we have a wait or waitresp
1708 //--------------------------------------------------------------------------
1709 uint32_t seconds = 0;
1710 if( rsp->hdr.status == kXR_wait )
1711 seconds = ntohl( rsp->body.wait.seconds ) + 5; // we need extra time
1712 // to re-send the request
1713 else if( rsp->hdr.status == kXR_waitresp )
1714 {
1715 seconds = ntohl( rsp->body.waitresp.seconds );
1716
1717 log->Dump( XRootDMsg, "[%s] Got kXR_waitresp response of %u seconds, "
1718 "setting up wait barrier.",
1719 info->streamName.c_str(),
1720 seconds );
1721 }
1722
1723 time_t barrier = time(0) + seconds;
1724 if( info->waitBarrier < barrier )
1725 info->waitBarrier = barrier;
1726
1727 //--------------------------------------------------------------------------
1728 // If we got a response to an open request, we may need to bump the counter
1729 // of open files
1730 //--------------------------------------------------------------------------
1731 uint16_t sid; memcpy( &sid, rsp->hdr.streamid, 2 );
1732 std::set<uint16_t>::iterator sidIt = info->sentOpens.find( sid );
1733 if( sidIt != info->sentOpens.end() )
1734 {
1735 if( rsp->hdr.status == kXR_waitresp )
1736 return NoAction;
1737 info->sentOpens.erase( sidIt );
1738 if( rsp->hdr.status == kXR_ok )
1739 {
1740 ++info->openFiles;
1741 info->finstcnt.fetch_add( 1, std::memory_order_relaxed ); // another file File object instance has been bound with this connection
1742 }
1743 return NoAction;
1744 }
1745
1746 //--------------------------------------------------------------------------
1747 // If we got a response to a close, we may need to decrement the counter of
1748 // open files
1749 //--------------------------------------------------------------------------
1750 sidIt = info->sentCloses.find( sid );
1751 if( sidIt != info->sentCloses.end() )
1752 {
1753 if( rsp->hdr.status == kXR_waitresp )
1754 return NoAction;
1755 info->sentCloses.erase( sidIt );
1756 --info->openFiles;
1757 return NoAction;
1758 }
1759 return NoAction;
1760 }
1761
1762 //----------------------------------------------------------------------------
1763 // Notify the transport about a message having been sent
1764 //----------------------------------------------------------------------------
1766 uint16_t subStream,
1767 uint32_t bytesSent,
1768 AnyObject &channelData )
1769 {
1770 // Called when a message has been sent. For messages that return on a
1771 // different pathid (and hence may use a different poller) it is possible
1772 // that the server has already replied and the reply will trigger
1773 // MessageReceived() before this method has been called. However for open
1774 // and close this is never the case and this method is used for tracking
1775 // only those.
1776 XRootDChannelInfo *info = 0;
1777 channelData.Get( info );
1778 if( !info ) return;
1779 XrdSysMutexHelper scopedLock( info->mutex );
1780 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
1781 uint16_t reqid = ntohs( req->header.requestid );
1782
1783
1784 //--------------------------------------------------------------------------
1785 // We need to track opens to know if we can close streams due to idleness
1786 //--------------------------------------------------------------------------
1787 uint16_t sid;
1788 memcpy( &sid, req->header.streamid, 2 );
1789
1790 if( reqid == kXR_open )
1791 info->sentOpens.insert( sid );
1792 else if( reqid == kXR_close )
1793 info->sentCloses.insert( sid );
1794 }
1795
1796
1797 //----------------------------------------------------------------------------
1798 // Get signature for given message
1799 //----------------------------------------------------------------------------
1801 {
1802 XRootDChannelInfo *info = 0;
1803 channelData.Get( info );
1804 return GetSignature( toSign, sign, info );
1805 }
1806
1807 //------------------------------------------------------------------------
1809 //------------------------------------------------------------------------
1811 Message *&sign,
1812 XRootDChannelInfo *info )
1813 {
1814 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
1815 if( pSecUnloadHandler->unloaded ) return Status( stError, errInvalidOp );
1816
1817 ClientRequest *thereq = reinterpret_cast<ClientRequest*>( toSign->GetBuffer() );
1818 if( !info ) return Status( stError, errInternal );
1819 if( info->protection )
1820 {
1821 SecurityRequest *newreq = 0;
1822 // check if we have to secure the request in the first place
1823 if( !( NEED2SECURE ( info->protection )( *thereq ) ) ) return Status();
1824 // secure (sign/encrypt) the request
1825 int rc = info->protection->Secure( newreq, *thereq, 0 );
1826 // there was an error
1827 if( rc < 0 )
1828 return Status( stError, errInternal, -rc );
1829
1830 sign = new Message();
1831 sign->Grab( reinterpret_cast<char*>( newreq ), rc );
1832 }
1833
1834 return Status();
1835 }
1836
1837 //------------------------------------------------------------------------
1839 //------------------------------------------------------------------------
1841 {
1842 XRootDChannelInfo *info = 0;
1843 channelData.Get( info );
1844 if( info->finstcnt.load( std::memory_order_relaxed ) > 0 )
1845 info->finstcnt.fetch_sub( 1, std::memory_order_relaxed );
1846 }
1847
1848 //----------------------------------------------------------------------------
1849 // Wait before exit
1850 //----------------------------------------------------------------------------
1852 {
1853 XrdSysRWLockHelper scope( pSecUnloadHandler->lock, false ); // obtain write lock
1854 pSecUnloadHandler->unloaded = true;
1855 }
1856
1857 //----------------------------------------------------------------------------
1858 // @return : true if encryption should be turned on, false otherwise
1859 //----------------------------------------------------------------------------
1861 AnyObject &channelData )
1862 {
1863 XRootDChannelInfo *info = 0;
1864 channelData.Get( info );
1865
1867 int notlsok = DefaultNoTlsOK;
1868 env->GetInt( "NoTlsOK", notlsok );
1869
1870
1871 if( notlsok )
1872 return info->encrypted;
1873
1874 XRootDStreamInfo &sInfo = info->stream[handShakeData->subStreamId];
1875
1876 // Did the server instructed us to switch to TLS right away?
1877 if( sInfo.serverFlags & kXR_gotoTLS )
1878 {
1879 if( handShakeData->subStreamId == 0 ) info->encrypted = true;
1880 return true ;
1881 }
1882
1883 //--------------------------------------------------------------------------
1884 // The control stream (sub-stream 0) might need to switch to TLS before
1885 // login or after login
1886 //--------------------------------------------------------------------------
1887 if( handShakeData->subStreamId == 0 )
1888 {
1889 //------------------------------------------------------------------------
1890 // We are about to login and the server asked to start encrypting
1891 // before login
1892 //------------------------------------------------------------------------
1893 if( ( sInfo.status == XRootDStreamInfo::LoginSent ) &&
1894 ( info->serverFlags & kXR_tlsLogin ) )
1895 {
1896 info->encrypted = true;
1897 return true;
1898 }
1899
1900 //--------------------------------------------------------------------
1901 // The hand-shake is done and the server requested to encrypt the session
1902 //--------------------------------------------------------------------
1903 if( (sInfo.status == XRootDStreamInfo::Connected ||
1904 //--------------------------------------------------------------------
1905 // we really need to turn on TLS before we sent kXR_endsess and we
1906 // are about to do so (1st enable encryption, then send kXR_endsess)
1907 //--------------------------------------------------------------------
1909 ( info->serverFlags & kXR_tlsSess ) )
1910 {
1911 info->encrypted = true;
1912 return true;
1913 }
1914 }
1915 //--------------------------------------------------------------------------
1916 // A data stream (sub-stream > 0) if need be will be switched to TLS before
1917 // bind.
1918 //--------------------------------------------------------------------------
1919 else
1920 {
1921 //------------------------------------------------------------------------
1922 // We are about to bind a data stream and the server asked to start
1923 // encrypting before bind
1924 //------------------------------------------------------------------------
1925 if( ( sInfo.status == XRootDStreamInfo::BindSent ) &&
1926 ( info->serverFlags & kXR_tlsData ) )
1927 {
1928 return true;
1929 }
1930 }
1931
1932 return false;
1933 }
1934
1935 //------------------------------------------------------------------------
1936 // Get bind preference for the next data stream
1937 //------------------------------------------------------------------------
1939 AnyObject &channelData )
1940 {
1941 XRootDChannelInfo *info = 0;
1942 channelData.Get( info );
1943
1944 if(!info || !info->bindSelector)
1945 return url;
1946
1947 return URL( info->bindSelector->Get() );
1948 }
1949
1950 //----------------------------------------------------------------------------
1951 // Generate the message to be sent as an initial handshake
1952 // (handshake+kXR_protocol)
1953 //----------------------------------------------------------------------------
1954 Message *XRootDTransport::GenerateInitialHSProtocol( HandShakeData *hsData,
1955 XRootDChannelInfo *info,
1956 kXR_char expect )
1957 {
1958 Log *log = DefaultEnv::GetLog();
1960 "[%s] Sending out the initial hand shake + kXR_protocol",
1961 hsData->streamName.c_str() );
1962
1963 Message *msg = new Message();
1964
1965 msg->Allocate( 20+sizeof(ClientProtocolRequest) );
1966 msg->Zero();
1967
1969 init->fourth = htonl(4);
1970 init->fifth = htonl(2012);
1971
1973 InitProtocolReq( proto, info, expect );
1974
1975 return msg;
1976 }
1977
1978 //------------------------------------------------------------------------
1979 // Generate the protocol message
1980 //------------------------------------------------------------------------
1981 Message *XRootDTransport::GenerateProtocol( HandShakeData *hsData,
1982 XRootDChannelInfo *info,
1983 kXR_char expect )
1984 {
1985 Log *log = DefaultEnv::GetLog();
1986 log->Debug( XRootDTransportMsg,
1987 "[%s] Sending out the kXR_protocol",
1988 hsData->streamName.c_str() );
1989
1990 Message *msg = new Message();
1991 msg->Allocate( sizeof(ClientProtocolRequest) );
1992 msg->Zero();
1993
1994 ClientProtocolRequest *proto = (ClientProtocolRequest *)msg->GetBuffer();
1995 InitProtocolReq( proto, info, expect );
1996
1997 return msg;
1998 }
1999
2000 //------------------------------------------------------------------------
2001 // Initialize protocol request
2002 //------------------------------------------------------------------------
2003 void XRootDTransport::InitProtocolReq( ClientProtocolRequest *request,
2004 XRootDChannelInfo *info,
2005 kXR_char expect )
2006 {
2007 request->requestid = htons(kXR_protocol);
2008 request->clientpv = htonl(kXR_PROTOCOLVERSION);
2011
2012 int notlsok = DefaultNoTlsOK;
2013 int tlsnodata = DefaultTlsNoData;
2014
2015 XrdCl::Env *env = XrdCl::DefaultEnv::GetEnv();
2016
2017 env->GetInt( "NoTlsOK", notlsok );
2018
2020 env->GetInt( "TlsNoData", tlsnodata );
2021
2022 if (info->encrypted || InitTLS())
2024
2025 if (info->encrypted && !(notlsok || tlsnodata))
2027
2028 request->expect = expect;
2029
2030 //--------------------------------------------------------------------------
2031 // If we are in the curse of establishing a connection in the context of
2032 // TPC update the expect! (this will be never followed be a bind)
2033 //--------------------------------------------------------------------------
2034 if( info->istpc )
2036 }
2037
2038 //----------------------------------------------------------------------------
2039 // Process the server initial handshake response
2040 //----------------------------------------------------------------------------
2041 XRootDStatus XRootDTransport::ProcessServerHS( HandShakeData *hsData,
2042 XRootDChannelInfo *info )
2043 {
2044 Log *log = DefaultEnv::GetLog();
2045
2046 Message *msg = hsData->in;
2047 ServerResponseHeader *respHdr = (ServerResponseHeader *)msg->GetBuffer();
2048 ServerInitHandShake *hs = (ServerInitHandShake *)msg->GetBuffer(4);
2049
2050 if( respHdr->status != kXR_ok )
2051 {
2052 log->Error( XRootDTransportMsg, "[%s] Invalid hand shake response",
2053 hsData->streamName.c_str() );
2054
2055 return XRootDStatus( stFatal, errHandShakeFailed, 0, "Invalid hand shake response." );
2056 }
2057
2058 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2059 const uint32_t pv = ntohl(hs->protover);
2060 sInfo.serverFlags = ntohl(hs->msgval) == kXR_DataServer ?
2063
2064 if( hsData->subStreamId == 0 )
2065 {
2066 info->protocolVersion = pv;
2067 info->serverFlags = sInfo.serverFlags;
2068 }
2069
2070 log->Debug( XRootDTransportMsg,
2071 "[%s] Got the server hand shake response (%s, protocol "
2072 "version %x)",
2073 hsData->streamName.c_str(),
2074 ServerFlagsToStr( sInfo.serverFlags ).c_str(),
2075 info->protocolVersion );
2076
2077 return XRootDStatus( stOK, suContinue );
2078 }
2079
2080 //----------------------------------------------------------------------------
2081 // Process the protocol response
2082 //----------------------------------------------------------------------------
2083 XRootDStatus XRootDTransport::ProcessProtocolResp( HandShakeData *hsData,
2084 XRootDChannelInfo *info )
2085 {
2086 Log *log = DefaultEnv::GetLog();
2087
2088 XRootDStatus st = UnMarshallBody( hsData->in, kXR_protocol );
2089 if( !st.IsOK() )
2090 return st;
2091
2092 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2093
2094
2095 if( rsp->hdr.status != kXR_ok )
2096 {
2097 log->Error( XRootDTransportMsg, "[%s] kXR_protocol request failed",
2098 hsData->streamName.c_str() );
2099
2100 return XRootDStatus( stFatal, errHandShakeFailed, 0, "kXR_protocol request failed" );
2101 }
2102
2103 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2104 if( rsp->body.protocol.pval >= 0x297 )
2105 sInfo.serverFlags = rsp->body.protocol.flags;
2106
2107 if( hsData->subStreamId > 0 )
2108 return XRootDStatus( stOK, suContinue );
2109
2110 info->serverFlags = sInfo.serverFlags;
2111
2112 XrdCl::Env *env = XrdCl::DefaultEnv::GetEnv();
2113 int notlsok = DefaultNoTlsOK;
2114 env->GetInt( "NoTlsOK", notlsok );
2115
2116 if( rsp->body.protocol.pval < kXR_PROTTLSVERSION && info->encrypted )
2117 {
2118 //------------------------------------------------------------------------
2119 // User requested an encrypted connection but the server is to old to
2120 // support it!
2121 //------------------------------------------------------------------------
2122 if( !notlsok ) return XRootDStatus( stFatal, errTlsError, ENOTSUP, "TLS not supported" );
2123
2124 //------------------------------------------------------------------------
2125 // We are falling back to unencrypted data transmission, as configured
2126 // in XRD_NOTLSOK environment variable
2127 //------------------------------------------------------------------------
2128 log->Info( XRootDTransportMsg,
2129 "[%s] Falling back to unencrypted transmission, server does "
2130 "not support TLS encryption.",
2131 hsData->streamName.c_str() );
2132 info->encrypted = false;
2133 }
2134
2135 if( rsp->body.protocol.pval >= 0x297 )
2136 info->serverFlags = rsp->body.protocol.flags;
2137
2138 if( rsp->hdr.dlen > 8 )
2139 {
2140 info->protRespBuff.assign( sizeof( ServerResponseBody_Protocol ), 0 );
2141 info->protRespSize = 0;
2142 ServerResponseBody_Protocol *protRespBody =
2143 reinterpret_cast<ServerResponseBody_Protocol*>( info->protRespBuff.data() );
2144 protRespBody->flags = rsp->body.protocol.flags;
2145 protRespBody->pval = rsp->body.protocol.pval;
2146
2147 char* bodybuff = reinterpret_cast<char*>( &rsp->body.protocol.secreq );
2148 size_t bodysize = rsp->hdr.dlen - 8;
2149 XRootDStatus st = ProcessProtocolBody( bodybuff, bodysize, info );
2150 if( !st.IsOK() )
2151 return st;
2152 }
2153
2154 log->Debug( XRootDTransportMsg,
2155 "[%s] kXR_protocol successful (%s, protocol version %x)",
2156 hsData->streamName.c_str(),
2157 ServerFlagsToStr( info->serverFlags ).c_str(),
2158 info->protocolVersion );
2159
2160 if( !( info->serverFlags & kXR_haveTLS ) && info->encrypted )
2161 {
2162 //------------------------------------------------------------------------
2163 // User requested an encrypted connection but the server was not configured
2164 // to support encryption!
2165 //------------------------------------------------------------------------
2166 return XRootDStatus( stFatal, errTlsError, ECONNREFUSED,
2167 "Server was not configured to support encryption." );
2168 }
2169
2170 //--------------------------------------------------------------------------
2171 // Now see if we have to enforce encryption in case the server does not
2172 // support PgRead/PgWrite
2173 //--------------------------------------------------------------------------
2174 int tlsOnNoPgrw = DefaultWantTlsOnNoPgrw;
2175 env->GetInt( "WantTlsOnNoPgrw", tlsOnNoPgrw );
2176 if( !( info->serverFlags & kXR_suppgrw ) && tlsOnNoPgrw )
2177 {
2178 //------------------------------------------------------------------------
2179 // If user requested encryption just make sure it is not switched off for
2180 // data
2181 //------------------------------------------------------------------------
2182 if( info->encrypted )
2183 {
2184 log->Debug( XRootDTransportMsg,
2185 "[%s] Server does not support PgRead/PgWrite and"
2186 " WantTlsOnNoPgrw is on; enforcing encryption for data.",
2187 hsData->streamName.c_str() );
2188 env->PutInt( "TlsNoData", DefaultTlsNoData );
2189 }
2190 //------------------------------------------------------------------------
2191 // Otherwise, if server is not enforcing data encryption, we will need to
2192 // redo the protocol request with kXR_wantTLS set.
2193 //------------------------------------------------------------------------
2194 else if( !( info->serverFlags & kXR_tlsData ) &&
2195 ( info->serverFlags & kXR_haveTLS ) )
2196 {
2197 info->encrypted = true;
2198 return XRootDStatus( stOK, suRetry );
2199 }
2200 }
2201
2202 return XRootDStatus( stOK, suContinue );
2203 }
2204
2205 XRootDStatus XRootDTransport::ProcessProtocolBody( char *bodybuff,
2206 size_t bodysize,
2207 XRootDChannelInfo *info )
2208 {
2209 //--------------------------------------------------------------------------
2210 // Parse bind preferences
2211 //--------------------------------------------------------------------------
2212 XrdProto::bifReqs *bifreq = reinterpret_cast<XrdProto::bifReqs*>( bodybuff );
2213 if( bodysize >= sizeof( XrdProto::bifReqs ) && bifreq->theTag == 'B' )
2214 {
2215 bodybuff += sizeof( XrdProto::bifReqs );
2216 bodysize -= sizeof( XrdProto::bifReqs );
2217
2218 if( bodysize < bifreq->bifILen )
2219 return XRootDStatus( stError, errDataError, 0, "Received incomplete "
2220 "protocol response." );
2221 std::string bindprefs_str( bodybuff, bifreq->bifILen );
2222 std::vector<std::string> bindprefs;
2223 Utils::splitString( bindprefs, bindprefs_str, "," );
2224 info->bindSelector.reset( new BindPrefSelector( std::move( bindprefs ) ) );
2225 bodybuff += bifreq->bifILen;
2226 bodysize -= bifreq->bifILen;
2227 }
2228 //--------------------------------------------------------------------------
2229 // Parse security requirements
2230 //--------------------------------------------------------------------------
2231 XrdProto::secReqs *secreq = reinterpret_cast<XrdProto::secReqs*>( bodybuff );
2232 static const size_t secHdrLen = sizeof( XrdProto::secReqs ) -
2233 sizeof( ServerResponseSVec_Protocol );
2234 if( bodysize >= secHdrLen && secreq->theTag == 'S' )
2235 {
2236 //------------------------------------------------------------------------
2237 // Copy only the header and the secvsz entries of the security vector,
2238 // the server may send fewer bytes than declared or trailing garbage
2239 //------------------------------------------------------------------------
2240 size_t secsize = secHdrLen + secreq->secvsz *
2241 sizeof( ServerResponseSVec_Protocol );
2242 if( bodysize < secsize )
2243 return XRootDStatus( stError, errDataError, 0, "Received incomplete "
2244 "protocol response." );
2245
2246 size_t respsize = kXR_ShortProtRespLen + secsize;
2247 if( info->protRespBuff.size() < respsize )
2248 info->protRespBuff.resize( respsize, 0 );
2249 memcpy( info->protRespBuff.data() + kXR_ShortProtRespLen, secreq, secsize );
2250 info->protRespSize = respsize;
2251 }
2252
2253 return XRootDStatus();
2254 }
2255
2256 //----------------------------------------------------------------------------
2257 // Generate the bind message
2258 //----------------------------------------------------------------------------
2259 Message *XRootDTransport::GenerateBind( HandShakeData *hsData,
2260 XRootDChannelInfo *info )
2261 {
2262 Log *log = DefaultEnv::GetLog();
2263
2264 log->Debug( XRootDTransportMsg,
2265 "[%s] Sending out the bind request",
2266 hsData->streamName.c_str() );
2267
2268
2269 Message *msg = new Message( sizeof( ClientBindRequest ) );
2270 ClientBindRequest *bindReq = (ClientBindRequest *)msg->GetBuffer();
2271
2272 bindReq->requestid = kXR_bind;
2273 memcpy( bindReq->sessid, info->sessionId, 16 );
2274 bindReq->dlen = 0;
2275 MarshallRequest( msg );
2276 return msg;
2277 }
2278
2279 //----------------------------------------------------------------------------
2280 // Generate the bind message
2281 //----------------------------------------------------------------------------
2282 XRootDStatus XRootDTransport::ProcessBindResp( HandShakeData *hsData,
2283 XRootDChannelInfo *info )
2284 {
2285 Log *log = DefaultEnv::GetLog();
2286
2287 XRootDStatus st = UnMarshallBody( hsData->in, kXR_bind );
2288 if( !st.IsOK() )
2289 return st;
2290
2291 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2292
2293 if( rsp->hdr.status != kXR_ok )
2294 {
2295 log->Error( XRootDTransportMsg, "[%s] kXR_bind request failed",
2296 hsData->streamName.c_str() );
2297 return XRootDStatus( stFatal, errHandShakeFailed, 0, "kXR_bind request failed" );
2298 }
2299
2300 info->stream[hsData->subStreamId].pathId = rsp->body.bind.substreamid;
2301 log->Debug( XRootDTransportMsg, "[%s] kXR_bind successful",
2302 hsData->streamName.c_str() );
2303
2304 return XRootDStatus();
2305 }
2306
2307 //----------------------------------------------------------------------------
2308 // Generate the login message
2309 //----------------------------------------------------------------------------
2310 Message *XRootDTransport::GenerateLogIn( HandShakeData *hsData,
2311 XRootDChannelInfo *info )
2312 {
2313 Log *log = DefaultEnv::GetLog();
2314 Env *env = DefaultEnv::GetEnv();
2315
2316 //--------------------------------------------------------------------------
2317 // Compute the login cgi
2318 //--------------------------------------------------------------------------
2319 int timeZone = XrdSysTimer::TimeZone();
2320 char *hostName = XrdNetUtils::MyHostName();
2321 std::string countryCode = Utils::FQDNToCC( hostName );
2322 char *cgiBuffer = new char[1024 + info->logintoken.size()];
2323 std::string appName;
2324 std::string monInfo;
2325 env->GetString( "AppName", appName );
2326 env->GetString( "MonInfo", monInfo );
2327 if( info->logintoken.empty() )
2328 {
2329 snprintf( cgiBuffer, 1024,
2330 "xrd.cc=%s&xrd.tz=%d&xrd.appname=%s&xrd.info=%s&"
2331 "xrd.hostname=%s&xrd.rn=%s", countryCode.c_str(), timeZone,
2332 appName.c_str(), monInfo.c_str(), hostName, XrdVERSION );
2333 }
2334 else
2335 {
2336 snprintf( cgiBuffer, 1024,
2337 "xrd.cc=%s&xrd.tz=%d&xrd.appname=%s&xrd.info=%s&"
2338 "xrd.hostname=%s&xrd.rn=%s&%s", countryCode.c_str(), timeZone,
2339 appName.c_str(), monInfo.c_str(), hostName, XrdVERSION, info->logintoken.c_str() );
2340 }
2341 uint16_t cgiLen = strlen( cgiBuffer );
2342 free( hostName );
2343
2344 //--------------------------------------------------------------------------
2345 // Generate the message
2346 //--------------------------------------------------------------------------
2347 Message *msg = new Message( sizeof(ClientLoginRequest) + cgiLen );
2348 ClientLoginRequest *loginReq = (ClientLoginRequest *)msg->GetBuffer();
2349
2350 loginReq->requestid = kXR_login;
2351 loginReq->pid = ::getpid();
2352 loginReq->capver[0] = (kXR_char) kXR_asyncap | (kXR_char) kXR_ver005;
2353 loginReq->dlen = cgiLen;
2355#ifdef WITH_XRDEC
2356 loginReq->ability2 = kXR_ecredir;
2357#endif
2358
2359 int multiProtocol = 0;
2360 env->GetInt( "MultiProtocol", multiProtocol );
2361 if(multiProtocol)
2362 loginReq->ability |= kXR_multipr;
2363
2364 //--------------------------------------------------------------------------
2365 // Check the IP stacks
2366 //--------------------------------------------------------------------------
2368 bool dualStack = false;
2369 bool privateIPv6 = false;
2370 bool privateIPv4 = false;
2371
2372 if( (stacks & XrdNetUtils::hasIP64) == XrdNetUtils::hasIP64 )
2373 {
2374 dualStack = true;
2375 loginReq->ability |= kXR_hasipv64;
2376 }
2377
2378 if( (stacks & XrdNetUtils::hasIPv6) && !(stacks & XrdNetUtils::hasPub6) )
2379 {
2380 privateIPv6 = true;
2381 loginReq->ability |= kXR_onlyprv6;
2382 }
2383
2384 if( (stacks & XrdNetUtils::hasIPv4) && !(stacks & XrdNetUtils::hasPub4) )
2385 {
2386 privateIPv4 = true;
2387 loginReq->ability |= kXR_onlyprv4;
2388 }
2389
2390 // The following code snippet tries to overcome the problem that this host
2391 // may still be dual-stacked but we don't know it because one of the
2392 // interfaces was not registered in DNS.
2393 //
2394 if( !dualStack && hsData->serverAddr )
2395 {if ( ( ( stacks & XrdNetUtils::hasIPv4 )
2396 && hsData->serverAddr->isIPType(XrdNetAddrInfo::IPv6))
2397 || ( ( stacks & XrdNetUtils::hasIPv6 )
2398 && hsData->serverAddr->isIPType(XrdNetAddrInfo::IPv4)))
2399 {dualStack = true;
2400 loginReq->ability |= kXR_hasipv64;
2401 }
2402 }
2403
2404 //--------------------------------------------------------------------------
2405 // Check the username
2406 //--------------------------------------------------------------------------
2407 std::string buffer( 8, 0 );
2408 if( hsData->url->GetUserName().length() )
2409 buffer = hsData->url->GetUserName();
2410 else
2411 {
2412 char *name = new char[1024];
2413 if( !XrdOucUtils::UserName( geteuid(), name, 1024 ) )
2414 buffer = name;
2415 else
2416 buffer = "_anon_";
2417 delete [] name;
2418 }
2419 buffer.resize( 8, 0 );
2420 std::copy( buffer.begin(), buffer.end(), (char*)loginReq->username );
2421
2422 msg->Append( cgiBuffer, cgiLen, 24 );
2423
2424 log->Debug( XRootDTransportMsg, "[%s] Sending out kXR_login request, "
2425 "username: %s, cgi: %s, dual-stack: %s, private IPv4: %s, "
2426 "private IPv6: %s", hsData->streamName.c_str(),
2427 loginReq->username, cgiBuffer, dualStack ? "true" : "false",
2428 privateIPv4 ? "true" : "false",
2429 privateIPv6 ? "true" : "false" );
2430
2431 delete [] cgiBuffer;
2432 MarshallRequest( msg );
2433 return msg;
2434 }
2435
2436 //----------------------------------------------------------------------------
2437 // Process the protocol response
2438 //----------------------------------------------------------------------------
2439 XRootDStatus XRootDTransport::ProcessLogInResp( HandShakeData *hsData,
2440 XRootDChannelInfo *info )
2441 {
2442 Log *log = DefaultEnv::GetLog();
2443
2444 XRootDStatus st = UnMarshallBody( hsData->in, kXR_login );
2445 if( !st.IsOK() )
2446 return st;
2447
2448 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2449
2450 if( rsp->hdr.status != kXR_ok )
2451 {
2452 log->Error( XRootDTransportMsg, "[%s] Got invalid login response",
2453 hsData->streamName.c_str() );
2454 return XRootDStatus( stFatal, errLoginFailed, 0, "Got invalid login response." );
2455 }
2456
2457 if( !info->firstLogIn )
2458 memcpy( info->oldSessionId, info->sessionId, 16 );
2459
2460 if( rsp->hdr.dlen == 0 && info->protocolVersion <= 0x289 )
2461 {
2462 //--------------------------------------------------------------------------
2463 // This if statement is there only to support dCache inaccurate
2464 // implementation of XRoot protocol, that in some cases returns
2465 // an empty login response for protocol version <= 2.8.9.
2466 //--------------------------------------------------------------------------
2467 memset( info->sessionId, 0, 16 );
2468 log->Warning( XRootDTransportMsg,
2469 "[%s] Logged in, accepting empty login response.",
2470 hsData->streamName.c_str() );
2471 return XRootDStatus();
2472 }
2473
2474 if( rsp->hdr.dlen < 16 )
2475 return XRootDStatus( stError, errDataError, 0, "Login response too short." );
2476
2477 memcpy( info->sessionId, rsp->body.login.sessid, 16 );
2478
2479 std::string sessId = Utils::Char2Hex( rsp->body.login.sessid, 16 );
2480
2481 log->Debug( XRootDTransportMsg, "[%s] Logged in, session: %s",
2482 hsData->streamName.c_str(), sessId.c_str() );
2483
2484 //--------------------------------------------------------------------------
2485 // We have an authentication info to process
2486 //--------------------------------------------------------------------------
2487 if( rsp->hdr.dlen > 16 )
2488 {
2489 size_t len = rsp->hdr.dlen-16;
2490 info->authBuffer = new char[len+1];
2491 info->authBuffer[len] = 0;
2492 memcpy( info->authBuffer, rsp->body.login.sec, len );
2493 log->Debug( XRootDTransportMsg, "[%s] Authentication is required: %s",
2494 hsData->streamName.c_str(), info->authBuffer );
2495
2496 return XRootDStatus( stOK, suContinue );
2497 }
2498
2499 return XRootDStatus();
2500 }
2501
2502 //----------------------------------------------------------------------------
2503 // Do the authentication
2504 //----------------------------------------------------------------------------
2505 XRootDStatus XRootDTransport::DoAuthentication( HandShakeData *hsData,
2506 XRootDChannelInfo *info )
2507 {
2508 //--------------------------------------------------------------------------
2509 // Prepare
2510 //--------------------------------------------------------------------------
2511 Log *log = DefaultEnv::GetLog();
2512 XRootDStreamInfo &sInfo = info->stream[hsData->subStreamId];
2513 XrdSecCredentials *credentials = 0;
2514 std::string protocolName;
2515
2516 //--------------------------------------------------------------------------
2517 // We're doing this for the first time
2518 //--------------------------------------------------------------------------
2519 if( sInfo.status == XRootDStreamInfo::LoginSent )
2520 {
2521 log->Debug( XRootDTransportMsg, "[%s] Sending authentication data",
2522 hsData->streamName.c_str() );
2523
2524 //------------------------------------------------------------------------
2525 // Set up the authentication environment
2526 //------------------------------------------------------------------------
2527 info->authEnv = new XrdOucEnv();
2528 info->authEnv->Put( "sockname", hsData->clientName.c_str() );
2529 info->authEnv->Put( "username", hsData->url->GetUserName().c_str() );
2530 info->authEnv->Put( "password", hsData->url->GetPassword().c_str() );
2531
2532 const URL::ParamsMap &urlParams = hsData->url->GetParams();
2533 URL::ParamsMap::const_iterator it;
2534 for( it = urlParams.begin(); it != urlParams.end(); ++it )
2535 {
2536 if( it->first.compare( 0, 4, "xrd." ) == 0 ||
2537 it->first.compare( 0, 6, "xrdcl." ) == 0 )
2538 info->authEnv->Put( it->first.c_str(), it->second.c_str() );
2539 }
2540
2541 //------------------------------------------------------------------------
2542 // Initialize some other structs
2543 //------------------------------------------------------------------------
2544 size_t authBuffLen = strlen( info->authBuffer );
2545 char *pars = (char *)malloc( authBuffLen + 1 );
2546 memcpy( pars, info->authBuffer, authBuffLen );
2547 info->authParams = new XrdSecParameters( pars, authBuffLen );
2548 sInfo.status = XRootDStreamInfo::AuthSent;
2549 delete [] info->authBuffer;
2550 info->authBuffer = 0;
2551
2552 //------------------------------------------------------------------------
2553 // Find a protocol that gives us valid credentials
2554 //------------------------------------------------------------------------
2555 XRootDStatus st = GetCredentials( credentials, hsData, info );
2556 if( !st.IsOK() )
2557 {
2558 CleanUpAuthentication( info );
2559 return st;
2560 }
2561 protocolName = info->authProtocol->Entity.prot;
2562 }
2563
2564 //--------------------------------------------------------------------------
2565 // We've been here already
2566 //--------------------------------------------------------------------------
2567 else
2568 {
2569 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2570 protocolName = info->authProtocol->Entity.prot;
2571
2572 //------------------------------------------------------------------------
2573 // We're required to send out more authentication data
2574 //------------------------------------------------------------------------
2575 if( rsp->hdr.status == kXR_authmore )
2576 {
2577 log->Debug( XRootDTransportMsg,
2578 "[%s] Sending more authentication data for %s",
2579 hsData->streamName.c_str(), protocolName.c_str() );
2580
2581 uint32_t len = rsp->hdr.dlen;
2582 char *secTokenData = (char*)malloc( len );
2583 memcpy( secTokenData, rsp->body.authmore.data, len );
2584 XrdSecParameters *secToken = new XrdSecParameters( secTokenData, len );
2585 XrdOucErrInfo ei( "", info->authEnv);
2586 credentials = info->authProtocol->getCredentials( secToken, &ei );
2587 delete secToken;
2588
2589 //----------------------------------------------------------------------
2590 // The protocol handler refuses to give us the data
2591 //----------------------------------------------------------------------
2592 if( !credentials )
2593 {
2594 log->Error( XRootDTransportMsg,
2595 "[%s] Auth protocol handler for %s refuses to give "
2596 "us more credentials %s",
2597 hsData->streamName.c_str(), protocolName.c_str(),
2598 ei.getErrText() );
2599 CleanUpAuthentication( info );
2600 return XRootDStatus( stFatal, errAuthFailed, 0, ei.getErrText() );
2601 }
2602 }
2603
2604 //------------------------------------------------------------------------
2605 // We have succeeded
2606 //------------------------------------------------------------------------
2607 else if( rsp->hdr.status == kXR_ok )
2608 {
2609 info->authProtocolName = info->authProtocol->Entity.prot;
2610
2611 //----------------------------------------------------------------------
2612 // Do we need protection?
2613 //----------------------------------------------------------------------
2614 if( !info->protRespBuff.empty() )
2615 {
2616 ServerResponseBody_Protocol *protRespBody =
2617 reinterpret_cast<ServerResponseBody_Protocol*>( info->protRespBuff.data() );
2618 int rc = XrdSecGetProtection( info->protection, *info->authProtocol, *protRespBody, info->protRespSize );
2619 if( rc > 0 )
2620 {
2621 log->Debug( XRootDTransportMsg,
2622 "[%s] XrdSecProtect loaded.", hsData->streamName.c_str() );
2623 }
2624 else if( rc == 0 )
2625 {
2626 log->Debug( XRootDTransportMsg,
2627 "[%s] XrdSecProtect: no protection needed.",
2628 hsData->streamName.c_str() );
2629 }
2630 else
2631 {
2632 log->Debug( XRootDTransportMsg,
2633 "[%s] Failed to load XrdSecProtect: %s",
2634 hsData->streamName.c_str(), XrdSysE2T( -rc ) );
2635 CleanUpAuthentication( info );
2636
2637 return XRootDStatus( stError, errAuthFailed, -rc, XrdSysE2T( -rc ) );
2638 }
2639 }
2640
2641 if( !info->protection )
2642 CleanUpAuthentication( info );
2643 else
2644 pSecUnloadHandler->Register( info->authProtocolName );
2645
2646 log->Debug( XRootDTransportMsg,
2647 "[%s] Authenticated with %s.", hsData->streamName.c_str(),
2648 protocolName.c_str() );
2649
2650 //--------------------------------------------------------------------
2651 // Clear the SSL error queue of the calling thread, as there might be
2652 // some leftover from the authentication!
2653 //--------------------------------------------------------------------
2655
2656 return XRootDStatus();
2657 }
2658 //------------------------------------------------------------------------
2659 // Failure
2660 //------------------------------------------------------------------------
2661 else if( rsp->hdr.status == kXR_error )
2662 {
2663 char *errmsg = new char[rsp->hdr.dlen-3]; errmsg[rsp->hdr.dlen-4] = 0;
2664 memcpy( errmsg, rsp->body.error.errmsg, rsp->hdr.dlen-4 );
2665 log->Error( XRootDTransportMsg,
2666 "[%s] Authentication with %s failed: %s",
2667 hsData->streamName.c_str(), protocolName.c_str(),
2668 errmsg );
2669 delete [] errmsg;
2670
2671 info->authProtocol->Delete();
2672 info->authProtocol = 0;
2673
2674 //----------------------------------------------------------------------
2675 // Find another protocol that gives us valid credentials
2676 //----------------------------------------------------------------------
2677 XRootDStatus st = GetCredentials( credentials, hsData, info );
2678 if( !st.IsOK() )
2679 {
2680 CleanUpAuthentication( info );
2681 return st;
2682 }
2683 protocolName = info->authProtocol->Entity.prot;
2684 }
2685 //------------------------------------------------------------------------
2686 // God knows what
2687 //------------------------------------------------------------------------
2688 else
2689 {
2690 info->authProtocolName = info->authProtocol->Entity.prot;
2691 CleanUpAuthentication( info );
2692
2693 log->Error( XRootDTransportMsg,
2694 "[%s] Authentication with %s failed: unexpected answer",
2695 hsData->streamName.c_str(), protocolName.c_str() );
2696 return XRootDStatus( stFatal, errAuthFailed, 0, "Authentication failed: unexpected answer." );
2697 }
2698 }
2699
2700 //--------------------------------------------------------------------------
2701 // Generate the client request
2702 //--------------------------------------------------------------------------
2703 Message *msg = new Message( sizeof(ClientAuthRequest)+credentials->size );
2704 msg->Zero();
2705 ClientRequest *req = (ClientRequest*)msg->GetBuffer();
2706 char *reqBuffer = msg->GetBuffer(sizeof(ClientAuthRequest));
2707
2708 req->header.requestid = kXR_auth;
2709 req->auth.dlen = credentials->size;
2710 memcpy( req->auth.credtype, protocolName.c_str(),
2711 protocolName.length() > 4 ? 4 : protocolName.length() );
2712
2713 memcpy( reqBuffer, credentials->buffer, credentials->size );
2714 hsData->out = msg;
2715 MarshallRequest( msg );
2716 delete credentials;
2717
2718 //------------------------------------------------------------------------
2719 // Clear the SSL error queue of the calling thread, as there might be
2720 // some leftover from the authentication!
2721 //------------------------------------------------------------------------
2723
2724 return XRootDStatus( stOK, suContinue );
2725 }
2726
2727 //------------------------------------------------------------------------
2728 // Get the initial credentials using one of the protocols
2729 //------------------------------------------------------------------------
2730 XRootDStatus XRootDTransport::GetCredentials( XrdSecCredentials *&credentials,
2731 HandShakeData *hsData,
2732 XRootDChannelInfo *info )
2733 {
2734 //--------------------------------------------------------------------------
2735 // Set up the auth handler
2736 //--------------------------------------------------------------------------
2737 Log *log = DefaultEnv::GetLog();
2738 XrdOucErrInfo ei( "", info->authEnv);
2739 XrdSecGetProt_t authHandler = GetAuthHandler();
2740 if( !authHandler )
2741 return XRootDStatus( stFatal, errAuthFailed, 0, "Could not load authentication handler." );
2742
2743 //--------------------------------------------------------------------------
2744 // Retrieve secuid and secgid, if available. These will override the fsuid
2745 // and fsgid of the current thread reading the credentials to prevent
2746 // security holes in case this process is running with elevated permissions.
2747 //--------------------------------------------------------------------------
2748 char *secuidc = (ei.getEnv()) ? ei.getEnv()->Get("xrdcl.secuid") : 0;
2749 char *secgidc = (ei.getEnv()) ? ei.getEnv()->Get("xrdcl.secgid") : 0;
2750
2751 int secuid = -1;
2752 int secgid = -1;
2753
2754 if(secuidc) secuid = atoi(secuidc);
2755 if(secgidc) secgid = atoi(secgidc);
2756
2757#ifdef __linux__
2758 ScopedFsUidSetter uidSetter(secuid, secgid, hsData->streamName);
2759 if(!uidSetter.IsOk()) {
2760 log->Error( XRootDTransportMsg, "[%s] Error while setting (fsuid, fsgid) to (%d, %d)",
2761 hsData->streamName.c_str(), secuid, secgid );
2762 return XRootDStatus( stFatal, errAuthFailed, 0, "Error while setting (fsuid, fsgid)." );
2763 }
2764#else
2765 if(secuid >= 0 || secgid >= 0) {
2766 log->Error( XRootDTransportMsg, "[%s] xrdcl.secuid and xrdcl.secgid only supported on Linux.",
2767 hsData->streamName.c_str() );
2768 return XRootDStatus( stFatal, errAuthFailed, 0, "xrdcl.secuid and xrdcl.secgid"
2769 " only supported on Linux" );
2770 }
2771#endif
2772
2773 //--------------------------------------------------------------------------
2774 // Loop over the possible protocols to find one that gives us valid
2775 // credentials
2776 //--------------------------------------------------------------------------
2777 XrdNetAddr &srvAddrInfo = *const_cast<XrdNetAddr *>(hsData->serverAddr);
2778 srvAddrInfo.SetTLS( info->encrypted );
2779 while(1)
2780 {
2781 //------------------------------------------------------------------------
2782 // Get the protocol
2783 //------------------------------------------------------------------------
2784 info->authProtocol = (*authHandler)( hsData->url->GetHostName().c_str(),
2785 srvAddrInfo,
2786 *info->authParams,
2787 &ei );
2788 if( !info->authProtocol )
2789 {
2790 log->Error( XRootDTransportMsg, "[%s] No protocols left to try",
2791 hsData->streamName.c_str() );
2792 return XRootDStatus( stFatal, errAuthFailed, 0, "No protocols left to try" );
2793 }
2794
2795 std::string protocolName = info->authProtocol->Entity.prot;
2796 log->Debug( XRootDTransportMsg, "[%s] Trying to authenticate using %s",
2797 hsData->streamName.c_str(), protocolName.c_str() );
2798
2799 //------------------------------------------------------------------------
2800 // Get the credentials from the current protocol
2801 //------------------------------------------------------------------------
2802 credentials = info->authProtocol->getCredentials( 0, &ei );
2803 if( !credentials )
2804 {
2805 log->Debug( XRootDTransportMsg,
2806 "[%s] Cannot get credentials for protocol %s: %s",
2807 hsData->streamName.c_str(), protocolName.c_str(),
2808 ei.getErrText() );
2809 info->authProtocol->Delete();
2810 continue;
2811 }
2812 return XRootDStatus( stOK, suContinue );
2813 }
2814 }
2815
2816 //------------------------------------------------------------------------
2817 // Clean up the data structures created for the authentication process
2818 //------------------------------------------------------------------------
2819 Status XRootDTransport::CleanUpAuthentication( XRootDChannelInfo *info )
2820 {
2821 if( info->authProtocol )
2822 info->authProtocol->Delete();
2823 delete info->authParams;
2824 delete info->authEnv;
2825 info->authProtocol = 0;
2826 info->authParams = 0;
2827 info->authEnv = 0;
2829 return Status();
2830 }
2831
2832 //------------------------------------------------------------------------
2833 // Clean up the data structures created for the protection purposes
2834 //------------------------------------------------------------------------
2835 Status XRootDTransport::CleanUpProtection( XRootDChannelInfo *info )
2836 {
2837 XrdSysRWLockHelper scope( pSecUnloadHandler->lock );
2838 if( pSecUnloadHandler->unloaded ) return Status( stError, errInvalidOp );
2839
2840 if( info->protection )
2841 {
2842 info->protection->Delete();
2843 info->protection = 0;
2844
2845 CleanUpAuthentication( info );
2846 }
2847
2848 info->protRespBuff.clear();
2849 info->protRespSize = 0;
2850
2851 return Status();
2852 }
2853
2854 //----------------------------------------------------------------------------
2855 // Get the authentication function handle
2856 //----------------------------------------------------------------------------
2857 XrdSecGetProt_t XRootDTransport::GetAuthHandler()
2858 {
2859 Log *log = DefaultEnv::GetLog();
2860 char errorBuff[1024];
2861
2862 // the static constructor is invoked only once and it is guaranteed that this
2863 // is thread safe
2864 static std::atomic<XrdSecGetProt_t> authHandler( XrdSecLoadSecFactory( errorBuff, 1024 ) );
2865 auto ret = authHandler.load( std::memory_order_relaxed );
2866 if( ret ) return ret;
2867
2868 // if we are here it means we failed to load the security library for the
2869 // first time and we hope the environment changed
2870
2871 // obtain a lock
2872 static XrdSysMutex mtx;
2873 XrdSysMutexHelper lck( mtx );
2874 // check if in the meanwhile some else didn't load the library
2875 ret = authHandler.load( std::memory_order_relaxed );
2876 if( ret ) return ret;
2877
2878 // load the library
2879 ret = XrdSecLoadSecFactory( errorBuff, 1024 );
2880 authHandler.store( ret, std::memory_order_relaxed );
2881 // if we failed report an error
2882 if( !ret )
2883 {
2884 log->Error( XRootDTransportMsg,
2885 "Unable to get the security framework: %s", errorBuff );
2886 return 0;
2887 }
2888 return ret;
2889 }
2890
2891 //----------------------------------------------------------------------------
2892 // Generate the end session message
2893 //----------------------------------------------------------------------------
2894 Message *XRootDTransport::GenerateEndSession( HandShakeData *hsData,
2895 XRootDChannelInfo *info )
2896 {
2897 Log *log = DefaultEnv::GetLog();
2898
2899 //--------------------------------------------------------------------------
2900 // Generate the message
2901 //--------------------------------------------------------------------------
2902 Message *msg = new Message( sizeof(ClientEndsessRequest) );
2903 ClientEndsessRequest *endsessReq = (ClientEndsessRequest *)msg->GetBuffer();
2904
2905 endsessReq->requestid = kXR_endsess;
2906 memcpy( endsessReq->sessid, info->oldSessionId, 16 );
2907 std::string sessId = Utils::Char2Hex( endsessReq->sessid, 16 );
2908
2909 log->Debug( XRootDTransportMsg, "[%s] Sending out kXR_endsess for session:"
2910 " %s", hsData->streamName.c_str(), sessId.c_str() );
2911
2912 MarshallRequest( msg );
2913
2914 Message *sign = 0;
2915 GetSignature( msg, sign, info );
2916 if( sign )
2917 {
2918 //------------------------------------------------------------------------
2919 // Now place both the signature and the request in a single buffer
2920 //------------------------------------------------------------------------
2921 uint32_t size = sign->GetSize();
2922 sign->ReAllocate( size + msg->GetSize() );
2923 char* buffer = sign->GetBuffer( size );
2924 memcpy( buffer, msg->GetBuffer(), msg->GetSize() );
2925 msg->Grab( sign->GetBuffer(), sign->GetSize() );
2926 }
2927
2928 return msg;
2929 }
2930
2931 //----------------------------------------------------------------------------
2932 // Process the protocol response
2933 //----------------------------------------------------------------------------
2934 Status XRootDTransport::ProcessEndSessionResp( HandShakeData *hsData,
2935 XRootDChannelInfo *info )
2936 {
2937 Log *log = DefaultEnv::GetLog();
2938
2939 Status st = UnMarshallBody( hsData->in, kXR_endsess );
2940 if( !st.IsOK() )
2941 return st;
2942
2943 ServerResponse *rsp = (ServerResponse*)hsData->in->GetBuffer();
2944
2945 // If we're good, we're good!
2946 if( rsp->hdr.status == kXR_ok )
2947 return Status();
2948
2949 // we ignore not found errors as such an error means the connection
2950 // has been already terminated
2951 if( rsp->hdr.status == kXR_error && rsp->body.error.errnum == kXR_NotFound )
2952 return Status();
2953
2954 // other errors
2955 if( rsp->hdr.status == kXR_error )
2956 {
2957 std::string errorMsg( rsp->body.error.errmsg, rsp->hdr.dlen - 4 );
2958 log->Error( XRootDTransportMsg, "[%s] Got error response to "
2959 "kXR_endsess: %s", hsData->streamName.c_str(),
2960 errorMsg.c_str() );
2962 }
2963
2964 // Wait Response.
2965 if( rsp->hdr.status == kXR_wait )
2966 {
2967 std::string msg( rsp->body.wait.infomsg, rsp->hdr.dlen - 4 );
2968 log->Info( XRootDTransportMsg, "[%s] Got wait response to "
2969 "kXR_endsess: %s", hsData->streamName.c_str(),
2970 msg.c_str() );
2971 hsData->out = GenerateEndSession( hsData, info );
2972 return Status( stOK, suRetry );
2973 }
2974
2975 // Any other response is protocol violation
2976 return Status( stError, errDataError );
2977 }
2978
2979 //----------------------------------------------------------------------------
2980 // Get a string representation of the server flags
2981 //----------------------------------------------------------------------------
2982 std::string XRootDTransport::ServerFlagsToStr( uint32_t flags )
2983 {
2984 std::string repr = "type: ";
2985 if( flags & kXR_isManager )
2986 repr += "manager ";
2987
2988 else if( flags & kXR_isServer )
2989 repr += "server ";
2990
2991 repr += "[";
2992
2993 if( flags & kXR_attrMeta )
2994 repr += "meta ";
2995
2996 else if( flags & kXR_attrCache )
2997 repr += "cache ";
2998
2999 else if( flags & kXR_attrProxy )
3000 repr += "proxy ";
3001
3002 else if( flags & kXR_attrSuper )
3003 repr += "super ";
3004
3005 else
3006 repr += " ";
3007
3008 repr.erase( repr.length()-1, 1 );
3009
3010 repr += "]";
3011 return repr;
3012 }
3013}
3014
3015namespace
3016{
3017 // Extract file name from a request
3018 //----------------------------------------------------------------------------
3019 char *GetDataAsString( char *msg )
3020 {
3022 char *fn = new char[req->dlen+1];
3023 memcpy( fn, msg + 24, req->dlen );
3024 fn[req->dlen] = 0;
3025 return fn;
3026 }
3027}
3028
3029namespace XrdCl
3030{
3031 //----------------------------------------------------------------------------
3032 // Get the description of a message
3033 //----------------------------------------------------------------------------
3034 void XRootDTransport::GenerateDescription( char *msg, std::ostringstream &o )
3035 {
3036 Log *log = DefaultEnv::GetLog();
3037 if( log->GetLevel() < Log::ErrorMsg )
3038 return;
3039
3040 ClientRequestHdr *req = (ClientRequestHdr *)msg;
3041 switch( req->requestid )
3042 {
3043 //------------------------------------------------------------------------
3044 // kXR_open
3045 //------------------------------------------------------------------------
3046 case kXR_open:
3047 {
3048 ClientOpenRequest *sreq = (ClientOpenRequest *)msg;
3049 o << "kXR_open (";
3050 char *fn = GetDataAsString( msg );
3051 o << "file: " << fn << ", ";
3052 delete [] fn;
3053 o << "mode: 0" << std::setbase(8) << sreq->mode << ", ";
3054 o << std::setbase(10);
3055 o << "flags: ";
3056 if( sreq->options == 0 )
3057 o << "none ";
3058 else
3059 {
3060 if( sreq->options & kXR_compress )
3061 o << "kXR_compress ";
3062 if( sreq->options & kXR_delete )
3063 o << "kXR_delete ";
3064 if( sreq->options & kXR_force )
3065 o << "kXR_force ";
3066 if( sreq->options & kXR_mkpath )
3067 o << "kXR_mkpath ";
3068 if( sreq->options & kXR_new )
3069 o << "kXR_new ";
3070 if( sreq->options & kXR_nowait )
3071 o << "kXR_nowait ";
3072 if( sreq->options & kXR_open_apnd )
3073 o << "kXR_open_apnd ";
3074 if( sreq->options & kXR_open_read )
3075 o << "kXR_open_read ";
3076 if( sreq->options & kXR_open_updt )
3077 o << "kXR_open_updt ";
3078 if( sreq->options & kXR_open_wrto )
3079 o << "kXR_open_wrto ";
3080 if( sreq->options & kXR_posc )
3081 o << "kXR_posc ";
3082 if( sreq->options & kXR_prefname )
3083 o << "kXR_prefname ";
3084 if( sreq->options & kXR_refresh )
3085 o << "kXR_refresh ";
3086 if( sreq->options & kXR_4dirlist )
3087 o << "kXR_4dirlist ";
3088 if( sreq->options & kXR_replica )
3089 o << "kXR_replica ";
3090 if( sreq->options & kXR_seqio )
3091 o << "kXR_seqio ";
3092 if( sreq->options & kXR_async )
3093 o << "kXR_async ";
3094 if( sreq->options & kXR_retstat )
3095 o << "kXR_retstat ";
3096 }
3097 o << "flagt: ";
3098 if( sreq->optiont == 0 )
3099 o << "none ";
3100 else
3101 {
3102 if( sreq->optiont & kXR_dup )
3103 o << "kXR_dup ";
3104 if( sreq->options & kXR_samefs )
3105 o << "kXR_samefs ";
3106 }
3107 o << "fhtemplt: " << FileHandleToStr( sreq->fhtemplt );
3108 o << ")";
3109 break;
3110 }
3111
3112 //------------------------------------------------------------------------
3113 // kXR_clone
3114 //------------------------------------------------------------------------
3115 case kXR_clone:
3116 {
3118 XrdProto::clone_list *dataChunk = (XrdProto::clone_list*)(msg + 24 );
3119 o << "kXR_clone ( ";
3120 o << "handle: " << FileHandleToStr( sreq->fhandle );
3121 o << std::setbase(10);
3122 o << " list [ ";
3123 for( size_t i = 0; i < req->dlen/sizeof(XrdProto::clone_list); ++i )
3124 {
3125 o << "(src_handle: ";
3126 o << FileHandleToStr( dataChunk[i].srcFH );
3127 o << ", ";
3128 o << std::setbase(10);
3129 o << "src_offset: " << dataChunk[i].srcOffs;
3130 o << ", src_length: " << dataChunk[i].srcLen;
3131 o << ", dst_offset: " << dataChunk[i].dstOffs << "); ";
3132 }
3133
3134 o << " ] )";
3135 break;
3136 }
3137
3138 //------------------------------------------------------------------------
3139 // kXR_close
3140 //------------------------------------------------------------------------
3141 case kXR_close:
3142 {
3144 o << "kXR_close (";
3145 o << "handle: " << FileHandleToStr( sreq->fhandle );
3146 o << ")";
3147 break;
3148 }
3149
3150 //------------------------------------------------------------------------
3151 // kXR_stat
3152 //------------------------------------------------------------------------
3153 case kXR_stat:
3154 {
3155 ClientStatRequest *sreq = (ClientStatRequest *)msg;
3156 o << "kXR_stat (";
3157 if( sreq->dlen )
3158 {
3159 char *fn = GetDataAsString( msg );;
3160 o << "path: " << fn << ", ";
3161 delete [] fn;
3162 }
3163 else
3164 {
3165 o << "handle: " << FileHandleToStr( sreq->fhandle );
3166 o << ", ";
3167 }
3168 o << "flags: ";
3169 if( sreq->options == 0 )
3170 o << "none";
3171 else
3172 {
3173 if( sreq->options & kXR_vfs )
3174 o << "kXR_vfs";
3175 }
3176 o << ")";
3177 break;
3178 }
3179
3180 //------------------------------------------------------------------------
3181 // kXR_read
3182 //------------------------------------------------------------------------
3183 case kXR_read:
3184 {
3185 ClientReadRequest *sreq = (ClientReadRequest *)msg;
3186 o << "kXR_read (";
3187 o << "handle: " << FileHandleToStr( sreq->fhandle );
3188 o << std::setbase(10);
3189 o << ", ";
3190 o << "offset: " << sreq->offset << ", ";
3191 o << "size: " << sreq->rlen << ")";
3192 break;
3193 }
3194
3195 //------------------------------------------------------------------------
3196 // kXR_pgread
3197 //------------------------------------------------------------------------
3198 case kXR_pgread:
3199 {
3201 o << "kXR_pgread (";
3202 o << "handle: " << FileHandleToStr( sreq->fhandle );
3203 o << std::setbase(10);
3204 o << ", ";
3205 o << "offset: " << sreq->offset << ", ";
3206 o << "size: " << sreq->rlen << ")";
3207 break;
3208 }
3209
3210 //------------------------------------------------------------------------
3211 // kXR_write
3212 //------------------------------------------------------------------------
3213 case kXR_write:
3214 {
3216 o << "kXR_write (";
3217 o << "handle: " << FileHandleToStr( sreq->fhandle );
3218 o << std::setbase(10);
3219 o << ", ";
3220 o << "offset: " << sreq->offset << ", ";
3221 o << "size: " << sreq->dlen << ")";
3222 break;
3223 }
3224
3225 //------------------------------------------------------------------------
3226 // kXR_pgwrite
3227 //------------------------------------------------------------------------
3228 case kXR_pgwrite:
3229 {
3231 o << "kXR_pgwrite (";
3232 o << "handle: " << FileHandleToStr( sreq->fhandle );
3233 o << std::setbase(10);
3234 o << ", ";
3235 o << "offset: " << sreq->offset << ", ";
3236 o << "size: " << sreq->dlen << ")";
3237 break;
3238 }
3239
3240 //------------------------------------------------------------------------
3241 // kXR_fattr
3242 //------------------------------------------------------------------------
3243 case kXR_fattr:
3244 {
3246 int nattr = sreq->numattr;
3247 int options = sreq->options;
3248 o << "kXR_fattr";
3249 switch (sreq->subcode) {
3250 case kXR_fattrGet:
3251 o << "Get";
3252 break;
3253 case kXR_fattrSet:
3254 o << "Set";
3255 break;
3256 case kXR_fattrList:
3257 o << "List";
3258 break;
3259 case kXR_fattrDel:
3260 o << "Delete";
3261 break;
3262 default:
3263 o << " unknown subcode: " << sreq->subcode;
3264 break;
3265 }
3266 o << " (handle: " << FileHandleToStr( sreq->fhandle );
3267 o << std::setbase(10);
3268 if (nattr)
3269 o << ", numattr: " << nattr;
3270 if (options) {
3271 o << ", options: ";
3272 if (options & 0x01)
3273 o << "new";
3274 if (options & 0x10)
3275 o << "list values";
3276 }
3277 o << ", total size: " << req->dlen << ")";
3278 break;
3279 }
3280
3281 //------------------------------------------------------------------------
3282 // kXR_sync
3283 //------------------------------------------------------------------------
3284 case kXR_sync:
3285 {
3286 ClientSyncRequest *sreq = (ClientSyncRequest *)msg;
3287 o << "kXR_sync (";
3288 o << "handle: " << FileHandleToStr( sreq->fhandle );
3289 o << ")";
3290 break;
3291 }
3292
3293 //------------------------------------------------------------------------
3294 // kXR_truncate
3295 //------------------------------------------------------------------------
3296 case kXR_truncate:
3297 {
3299 o << "kXR_truncate (";
3300 if( !sreq->dlen )
3301 o << "handle: " << FileHandleToStr( sreq->fhandle );
3302 else
3303 {
3304 char *fn = GetDataAsString( msg );
3305 o << "file: " << fn;
3306 delete [] fn;
3307 }
3308 o << std::setbase(10);
3309 o << ", ";
3310 o << "offset: " << sreq->offset;
3311 o << ")";
3312 break;
3313 }
3314
3315 //------------------------------------------------------------------------
3316 // kXR_readv
3317 //------------------------------------------------------------------------
3318 case kXR_readv:
3319 {
3320 unsigned char *fhandle = 0;
3321 o << "kXR_readv (";
3322
3323 o << "handle: ";
3324 readahead_list *dataChunk = (readahead_list*)(msg + 24 );
3325 fhandle = dataChunk[0].fhandle;
3326 if( fhandle )
3327 o << FileHandleToStr( fhandle );
3328 else
3329 o << "unknown";
3330 o << ", ";
3331 o << std::setbase(10);
3332 o << "chunks: [";
3333 uint64_t size = 0;
3334 for( size_t i = 0; i < req->dlen/sizeof(readahead_list); ++i )
3335 {
3336 size += dataChunk[i].rlen;
3337 o << "(offset: " << dataChunk[i].offset;
3338 o << ", size: " << dataChunk[i].rlen << "); ";
3339 }
3340 o << "], ";
3341 o << "total size: " << size << ")";
3342 break;
3343 }
3344
3345 //------------------------------------------------------------------------
3346 // kXR_writev
3347 //------------------------------------------------------------------------
3348 case kXR_writev:
3349 {
3350 unsigned char *fhandle = 0;
3351 o << "kXR_writev (";
3352
3353 XrdProto::write_list *wrtList =
3354 reinterpret_cast<XrdProto::write_list*>( msg + 24 );
3355 uint64_t size = 0;
3356 uint32_t numChunks = 0;
3357 for( size_t i = 0; i < req->dlen/sizeof(XrdProto::write_list); ++i )
3358 {
3359 fhandle = wrtList[i].fhandle;
3360 size += wrtList[i].wlen;
3361 ++numChunks;
3362 }
3363 o << "handle: ";
3364 if( fhandle )
3365 o << FileHandleToStr( fhandle );
3366 else
3367 o << "unknown";
3368 o << ", ";
3369 o << std::setbase(10);
3370 o << "chunks: " << numChunks << ", ";
3371 o << "total size: " << size << ")";
3372 break;
3373 }
3374
3375 //------------------------------------------------------------------------
3376 // kXR_locate
3377 //------------------------------------------------------------------------
3378 case kXR_locate:
3379 {
3381 char *fn = GetDataAsString( msg );;
3382 o << "kXR_locate (";
3383 o << "path: " << fn << ", ";
3384 delete [] fn;
3385 o << "flags: ";
3386 if( sreq->options == 0 )
3387 o << "none";
3388 else
3389 {
3390 if( sreq->options & kXR_refresh )
3391 o << "kXR_refresh ";
3392 if( sreq->options & kXR_prefname )
3393 o << "kXR_prefname ";
3394 if( sreq->options & kXR_nowait )
3395 o << "kXR_nowait ";
3396 if( sreq->options & kXR_force )
3397 o << "kXR_force ";
3398 if( sreq->options & kXR_compress )
3399 o << "kXR_compress ";
3400 }
3401 o << ")";
3402 break;
3403 }
3404
3405 //------------------------------------------------------------------------
3406 // kXR_mv
3407 //------------------------------------------------------------------------
3408 case kXR_mv:
3409 {
3410 ClientMvRequest *sreq = (ClientMvRequest *)msg;
3411 o << "kXR_mv (";
3412 o << "source: ";
3413 o.write( msg + sizeof( ClientMvRequest ), sreq->arg1len );
3414 o << ", ";
3415 o << "destination: ";
3416 o.write( msg + sizeof( ClientMvRequest ) + sreq->arg1len + 1, sreq->dlen - sreq->arg1len - 1 );
3417 o << ")";
3418 break;
3419 }
3420
3421 //------------------------------------------------------------------------
3422 // kXR_query
3423 //------------------------------------------------------------------------
3424 case kXR_query:
3425 {
3427 o << "kXR_query (";
3428 o << "code: ";
3429 switch( sreq->infotype )
3430 {
3431 case kXR_Qconfig: o << "kXR_Qconfig"; break;
3432 case kXR_Qckscan: o << "kXR_Qckscan"; break;
3433 case kXR_Qcksum: o << "kXR_Qcksum"; break;
3434 case kXR_Qopaque: o << "kXR_Qopaque"; break;
3435 case kXR_Qopaquf: o << "kXR_Qopaquf"; break;
3436 case kXR_Qopaqug: o << "kXR_Qopaqug"; break;
3437 case kXR_QPrep: o << "kXR_QPrep"; break;
3438 case kXR_Qspace: o << "kXR_Qspace"; break;
3439 case kXR_QStats: o << "kXR_QStats"; break;
3440 case kXR_Qvisa: o << "kXR_Qvisa"; break;
3441 case kXR_Qxattr: o << "kXR_Qxattr"; break;
3442 default: o << sreq->infotype; break;
3443 }
3444 o << ", ";
3445
3446 if( sreq->infotype == kXR_Qopaqug || sreq->infotype == kXR_Qvisa )
3447 {
3448 o << "handle: " << FileHandleToStr( sreq->fhandle );
3449 o << ", ";
3450 }
3451
3452 o << "arg length: " << sreq->dlen << ")";
3453 break;
3454 }
3455
3456 //------------------------------------------------------------------------
3457 // kXR_rm
3458 //------------------------------------------------------------------------
3459 case kXR_rm:
3460 {
3461 o << "kXR_rm (";
3462 char *fn = GetDataAsString( msg );;
3463 o << "path: " << fn << ")";
3464 delete [] fn;
3465 break;
3466 }
3467
3468 //------------------------------------------------------------------------
3469 // kXR_mkdir
3470 //------------------------------------------------------------------------
3471 case kXR_mkdir:
3472 {
3474 o << "kXR_mkdir (";
3475 char *fn = GetDataAsString( msg );
3476 o << "path: " << fn << ", ";
3477 delete [] fn;
3478 o << "mode: 0" << std::setbase(8) << sreq->mode << ", ";
3479 o << std::setbase(10);
3480 o << "flags: ";
3481 if( sreq->options[0] == 0 )
3482 o << "none";
3483 else
3484 {
3485 if( sreq->options[0] & kXR_mkdirpath )
3486 o << "kXR_mkdirpath";
3487 }
3488 o << ")";
3489 break;
3490 }
3491
3492 //------------------------------------------------------------------------
3493 // kXR_rmdir
3494 //------------------------------------------------------------------------
3495 case kXR_rmdir:
3496 {
3497 o << "kXR_rmdir (";
3498 char *fn = GetDataAsString( msg );
3499 o << "path: " << fn << ")";
3500 delete [] fn;
3501 break;
3502 }
3503
3504 //------------------------------------------------------------------------
3505 // kXR_chmod
3506 //------------------------------------------------------------------------
3507 case kXR_chmod:
3508 {
3510 o << "kXR_chmod (";
3511 char *fn = GetDataAsString( msg );
3512 o << "path: " << fn << ", ";
3513 delete [] fn;
3514 o << "mode: 0" << std::setbase(8) << sreq->mode << ")";
3515 break;
3516 }
3517
3518 //------------------------------------------------------------------------
3519 // kXR_ping
3520 //------------------------------------------------------------------------
3521 case kXR_ping:
3522 {
3523 o << "kXR_ping ()";
3524 break;
3525 }
3526
3527 //------------------------------------------------------------------------
3528 // kXR_protocol
3529 //------------------------------------------------------------------------
3530 case kXR_protocol:
3531 {
3533 o << "kXR_protocol (";
3534 o << "clientpv: 0x" << std::setbase(16) << sreq->clientpv << ")";
3535 break;
3536 }
3537
3538 //------------------------------------------------------------------------
3539 // kXR_dirlist
3540 //------------------------------------------------------------------------
3541 case kXR_dirlist:
3542 {
3543 o << "kXR_dirlist (";
3544 char *fn = GetDataAsString( msg );;
3545 o << "path: " << fn << ")";
3546 delete [] fn;
3547 break;
3548 }
3549
3550 //------------------------------------------------------------------------
3551 // kXR_set
3552 //------------------------------------------------------------------------
3553 case kXR_set:
3554 {
3555 o << "kXR_set (";
3556 char *fn = GetDataAsString( msg );;
3557 o << "data: " << fn << ")";
3558 delete [] fn;
3559 break;
3560 }
3561
3562 //------------------------------------------------------------------------
3563 // kXR_prepare
3564 //------------------------------------------------------------------------
3565 case kXR_prepare:
3566 {
3568 o << "kXR_prepare (";
3569 o << "flags: ";
3570
3571 if( sreq->options == 0 )
3572 o << "none";
3573 else
3574 {
3575 if( sreq->options & kXR_stage )
3576 o << "kXR_stage ";
3577 if( sreq->options & kXR_wmode )
3578 o << "kXR_wmode ";
3579 if( sreq->options & kXR_coloc )
3580 o << "kXR_coloc ";
3581 if( sreq->options & kXR_fresh )
3582 o << "kXR_fresh ";
3583 }
3584
3585 o << ", priority: " << (int) sreq->prty << ", ";
3586
3587 char *fn = GetDataAsString( msg );
3588 char *cursor;
3589 for( cursor = fn; *cursor; ++cursor )
3590 if( *cursor == '\n' ) *cursor = ' ';
3591
3592 o << "paths: " << fn << ")";
3593 delete [] fn;
3594 break;
3595 }
3596
3597 case kXR_chkpoint:
3598 {
3600 o << "kXR_chkpoint (";
3601 o << "opcode: ";
3602 if( sreq->opcode == kXR_ckpBegin ) o << "kXR_ckpBegin)";
3603 else if( sreq->opcode == kXR_ckpCommit ) o << "kXR_ckpCommit)";
3604 else if( sreq->opcode == kXR_ckpQuery ) o << "kXR_ckpQuery)";
3605 else if( sreq->opcode == kXR_ckpRollback ) o << "kXR_ckpRollback)";
3606 else if( sreq->opcode == kXR_ckpXeq )
3607 {
3608 o << "kXR_ckpXeq) ";
3609 // In this case our request body will be one of kXR_pgwrite,
3610 // kXR_truncate, kXR_write, or kXR_writev request.
3611 GenerateDescription( msg + sizeof( ClientChkPointRequest ), o );
3612 }
3613
3614 break;
3615 }
3616
3617 //------------------------------------------------------------------------
3618 // Default
3619 //------------------------------------------------------------------------
3620 default:
3621 {
3622 o << "kXR_unknown (length: " << req->dlen << ")";
3623 break;
3624 }
3625 };
3626 }
3627
3628 //----------------------------------------------------------------------------
3629 // Get a string representation of file handle
3630 //----------------------------------------------------------------------------
3631 std::string XRootDTransport::FileHandleToStr( const unsigned char handle[4] )
3632 {
3633 std::ostringstream o;
3634 o << "0x";
3635 for( uint8_t i = 0; i < 4; ++i )
3636 {
3637 o << std::setbase(16) << std::setfill('0') << std::setw(2);
3638 o << (int)handle[i];
3639 }
3640 return o.str();
3641 }
3642}
@ kXR_NotFound
kXR_int16 arg1len
Definition XProtocol.hh:460
#define kXR_isManager
struct ClientTruncateRequest truncate
Definition XProtocol.hh:917
@ kXR_ecredir
Definition XProtocol.hh:401
#define kXR_tlsLogin
@ kXR_fattrDel
Definition XProtocol.hh:300
@ kXR_fattrSet
Definition XProtocol.hh:303
@ kXR_fattrList
Definition XProtocol.hh:302
@ kXR_fattrGet
Definition XProtocol.hh:301
#define kXR_ShortProtRespLen
#define kXR_suppgrw
kXR_char fhandle[4]
Definition XProtocol.hh:565
kXR_unt16 requestid
Definition XProtocol.hh:424
ServerResponseStatus status
kXR_char fhandle[4]
Definition XProtocol.hh:823
#define kXR_gotoTLS
#define kXR_attrMeta
struct ClientPgReadRequest pgread
Definition XProtocol.hh:903
union ServerResponse::@040373375333017131300127053271011057331004327334 body
kXR_char fhandle[4]
Definition XProtocol.hh:848
#define kXR_haveTLS
kXR_char streamid[2]
Definition XProtocol.hh:158
kXR_char fhandle[4]
Definition XProtocol.hh:812
struct ClientMkdirRequest mkdir
Definition XProtocol.hh:900
kXR_int32 dlen
Definition XProtocol.hh:461
struct ClientAuthRequest auth
Definition XProtocol.hh:888
kXR_char streamid[2]
Definition XProtocol.hh:956
kXR_char fhtemplt[4]
Definition XProtocol.hh:516
kXR_unt16 options
Definition XProtocol.hh:513
struct ClientPgWriteRequest pgwrite
Definition XProtocol.hh:904
#define kXR_attrSuper
struct ClientReadVRequest readv
Definition XProtocol.hh:910
kXR_char pathid
Definition XProtocol.hh:689
kXR_char credtype[4]
Definition XProtocol.hh:172
kXR_char username[8]
Definition XProtocol.hh:426
@ kXR_open_wrto
Definition XProtocol.hh:499
@ kXR_compress
Definition XProtocol.hh:482
@ kXR_async
Definition XProtocol.hh:488
@ kXR_delete
Definition XProtocol.hh:483
@ kXR_prefname
Definition XProtocol.hh:491
@ kXR_nowait
Definition XProtocol.hh:497
@ kXR_open_read
Definition XProtocol.hh:486
@ kXR_open_updt
Definition XProtocol.hh:487
@ kXR_mkpath
Definition XProtocol.hh:490
@ kXR_seqio
Definition XProtocol.hh:498
@ kXR_replica
Definition XProtocol.hh:495
@ kXR_posc
Definition XProtocol.hh:496
@ kXR_refresh
Definition XProtocol.hh:489
@ kXR_new
Definition XProtocol.hh:485
@ kXR_force
Definition XProtocol.hh:484
@ kXR_4dirlist
Definition XProtocol.hh:494
@ kXR_open_apnd
Definition XProtocol.hh:492
@ kXR_retstat
Definition XProtocol.hh:493
struct ClientOpenRequest open
Definition XProtocol.hh:902
@ kXR_waitresp
Definition XProtocol.hh:948
@ kXR_redirect
Definition XProtocol.hh:946
@ kXR_status
Definition XProtocol.hh:949
@ kXR_ok
Definition XProtocol.hh:941
@ kXR_authmore
Definition XProtocol.hh:944
@ kXR_attn
Definition XProtocol.hh:943
@ kXR_wait
Definition XProtocol.hh:947
@ kXR_error
Definition XProtocol.hh:945
struct ServerResponseBody_Status bdy
struct ClientRequestHdr header
Definition XProtocol.hh:887
kXR_char fhandle[4]
Definition XProtocol.hh:543
kXR_unt16 optiont
Definition XProtocol.hh:514
kXR_char fhandle[4]
Definition XProtocol.hh:681
kXR_char fhandle[4]
Definition XProtocol.hh:695
struct ClientWriteVRequest writev
Definition XProtocol.hh:919
kXR_char fhandle[4]
Definition XProtocol.hh:258
struct ClientLoginRequest login
Definition XProtocol.hh:899
kXR_unt16 requestid
Definition XProtocol.hh:159
kXR_char fhandle[4]
Definition XProtocol.hh:669
kXR_char sessid[16]
Definition XProtocol.hh:183
@ kXR_read
Definition XProtocol.hh:126
@ kXR_open
Definition XProtocol.hh:123
@ kXR_writev
Definition XProtocol.hh:144
@ kXR_clone
Definition XProtocol.hh:145
@ kXR_readv
Definition XProtocol.hh:138
@ kXR_mkdir
Definition XProtocol.hh:121
@ kXR_sync
Definition XProtocol.hh:129
@ kXR_chmod
Definition XProtocol.hh:115
@ kXR_bind
Definition XProtocol.hh:137
@ kXR_dirlist
Definition XProtocol.hh:117
@ kXR_fattr
Definition XProtocol.hh:133
@ kXR_rm
Definition XProtocol.hh:127
@ kXR_query
Definition XProtocol.hh:114
@ kXR_write
Definition XProtocol.hh:132
@ kXR_login
Definition XProtocol.hh:120
@ kXR_auth
Definition XProtocol.hh:113
@ kXR_endsess
Definition XProtocol.hh:136
@ kXR_set
Definition XProtocol.hh:131
@ kXR_rmdir
Definition XProtocol.hh:128
@ kXR_1stRequest
Definition XProtocol.hh:112
@ kXR_truncate
Definition XProtocol.hh:141
@ kXR_protocol
Definition XProtocol.hh:119
@ kXR_mv
Definition XProtocol.hh:122
@ kXR_ping
Definition XProtocol.hh:124
@ kXR_stat
Definition XProtocol.hh:130
@ kXR_pgread
Definition XProtocol.hh:143
@ kXR_chkpoint
Definition XProtocol.hh:125
@ kXR_locate
Definition XProtocol.hh:140
@ kXR_close
Definition XProtocol.hh:116
@ kXR_pgwrite
Definition XProtocol.hh:139
@ kXR_prepare
Definition XProtocol.hh:134
struct ClientChmodRequest chmod
Definition XProtocol.hh:891
#define kXR_isServer
#define kXR_attrCache
struct ClientQueryRequest query
Definition XProtocol.hh:908
struct ClientReadRequest read
Definition XProtocol.hh:909
struct ClientMvRequest mv
Definition XProtocol.hh:901
kXR_int32 rlen
Definition XProtocol.hh:696
kXR_unt16 requestid
Definition XProtocol.hh:182
kXR_char sessid[16]
Definition XProtocol.hh:289
struct ClientChkPointRequest chkpoint
Definition XProtocol.hh:890
struct ServerResponseHeader hdr
@ kXR_asyncap
Definition XProtocol.hh:408
#define kXR_attrProxy
kXR_char options[1]
Definition XProtocol.hh:446
#define kXR_PROTOCOLVERSION
Definition XProtocol.hh:70
kXR_int64 offset
Definition XProtocol.hh:697
@ kXR_vfs
Definition XProtocol.hh:799
struct ClientPrepareRequest prepare
Definition XProtocol.hh:906
@ kXR_mkdirpath
Definition XProtocol.hh:440
@ kXR_wmode
Definition XProtocol.hh:625
@ kXR_fresh
Definition XProtocol.hh:627
@ kXR_coloc
Definition XProtocol.hh:626
@ kXR_stage
Definition XProtocol.hh:624
#define kXR_tlsSess
#define kXR_DataServer
@ kXR_dup
Definition XProtocol.hh:503
@ kXR_samefs
Definition XProtocol.hh:504
struct ClientWriteRequest write
Definition XProtocol.hh:918
#define kXR_PROTTLSVERSION
Definition XProtocol.hh:72
kXR_char capver[1]
Definition XProtocol.hh:429
struct ClientProtocolRequest protocol
Definition XProtocol.hh:907
@ kXR_QPrep
Definition XProtocol.hh:650
@ kXR_Qopaqug
Definition XProtocol.hh:661
@ kXR_Qconfig
Definition XProtocol.hh:655
@ kXR_Qopaquf
Definition XProtocol.hh:660
@ kXR_Qckscan
Definition XProtocol.hh:654
@ kXR_Qxattr
Definition XProtocol.hh:652
@ kXR_Qspace
Definition XProtocol.hh:653
@ kXR_Qvisa
Definition XProtocol.hh:656
@ kXR_QStats
Definition XProtocol.hh:649
@ kXR_Qcksum
Definition XProtocol.hh:651
@ kXR_Qopaque
Definition XProtocol.hh:659
struct ClientLocateRequest locate
Definition XProtocol.hh:898
kXR_char fhandle[4]
Definition XProtocol.hh:231
@ kXR_ver005
Definition XProtocol.hh:419
#define kXR_tlsData
@ kXR_readrdok
Definition XProtocol.hh:390
@ kXR_fullurl
Definition XProtocol.hh:388
@ kXR_onlyprv4
Definition XProtocol.hh:392
@ kXR_lclfile
Definition XProtocol.hh:394
@ kXR_multipr
Definition XProtocol.hh:389
@ kXR_redirflags
Definition XProtocol.hh:395
@ kXR_hasipv64
Definition XProtocol.hh:391
@ kXR_onlyprv6
Definition XProtocol.hh:393
ServerResponseHeader hdr
struct ClientCloneRequest clone
Definition XProtocol.hh:892
long long kXR_int64
Definition XPtypes.hh:98
unsigned char kXR_char
Definition XPtypes.hh:65
XrdVERSIONINFOREF(XrdCl)
XrdSecProtocol *(*) XrdSecGetProt_t(const char *hostname, XrdNetAddrInfo &endPoint, XrdSecParameters &sectoken, XrdOucErrInfo *einfo)
Typedef to simplify the encoding of methods returning XrdSecProtocol.
XrdSecBuffer XrdSecParameters
XrdSecBuffer XrdSecCredentials
XrdSecGetProt_t XrdSecLoadSecFactory(char *eBuff, int eBlen, const char *seclib)
int XrdSecGetProtection(XrdSecProtect *&protP, XrdSecProtocol &aprot, ServerResponseBody_Protocol &resp, unsigned int resplen)
#define NEED2SECURE(protP)
This class implements the XRootD protocol security protection.
const char * XrdSysE2T(int errcode)
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:136
bool GetInt(const std::string &key, int &value)
Definition XrdClEnv.cc:115
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.
Message(uint32_t size=0)
Constructor.
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:488
std::map< std::string, std::string > ParamsMap
Definition XrdClURL.hh:33
URL()
Default constructor.
Definition XrdClURL.cc:36
bool IsSecure() const
Does the protocol indicate encryption.
Definition XrdClURL.cc:458
bool IsTPC() const
Is the URL used in TPC context.
Definition XrdClURL.cc:466
std::string GetLoginToken() const
Get the login token if present in the opaque info.
Definition XrdClURL.cc:343
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.
XRootDStatus(uint16_t st=0, uint16_t code=0, uint32_t errN=0, const std::string &message="")
Constructor.
static uint16_t NbConnectedStrm(AnyObject &channelData)
Number of currently connected data streams.
virtual bool IsStreamTTLElapsed(time_t time, AnyObject &channelData)
Check if the stream should be disconnected.
virtual void Disconnect(AnyObject &channelData, uint16_t subStreamId)
The stream has been disconnected, do the cleanups.
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
XrdOucEnv(const char *vardata=0, int vardlen=0, const XrdSecEntity *secent=0)
Definition XrdOucEnv.cc:42
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:873
struct ServerResponseBifs_Protocol bifReqs
struct ServerResponseReqs_Protocol secReqs
kXR_char fhandle[4]
Definition XProtocol.hh:318
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.
PathID(uint16_t u=0, uint16_t d=0)
static void UnloadHandler(const std::string &trProt)
void Register(const std::string &protocol)
std::set< std::string > protocols
Procedure execution status.
Status(uint16_t st=stOK, uint16_t cod=errNone, uint32_t errN=0)
Constructor.
uint16_t code
Error type, or additional hints on what to do.
bool IsOK() const
We're fine.
void AdjustQueues(uint16_t size)
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.