XRootD
XrdCl::PollerBuiltIn Class Reference

A poller implementation using the build-in XRootD poller. More...

#include <XrdClPollerBuiltIn.hh>

+ Inheritance diagram for XrdCl::PollerBuiltIn:
+ Collaboration diagram for XrdCl::PollerBuiltIn:

Public Member Functions

 PollerBuiltIn ()
 Constructor. More...
 
 ~PollerBuiltIn ()
 
virtual bool AddSocket (Socket *socket, SocketHandler *handler)
 
virtual bool EnableReadNotification (Socket *socket, bool notify, time_t timeout=60)
 
virtual bool EnableWriteNotification (Socket *socket, bool notify, time_t timeout=60)
 
virtual bool Finalize ()
 Finalize the poller. More...
 
virtual bool Initialize ()
 Initialize the poller. More...
 
virtual bool IsRegistered (Socket *socket)
 Check whether the socket is registered with the poller. More...
 
virtual bool IsRunning () const
 Is the event loop running? More...
 
virtual bool RemoveSocket (Socket *socket)
 Remove the socket. More...
 
virtual void ShutdownEvents (Socket *socket)
 
virtual bool Start ()
 Start polling. More...
 
virtual bool Stop ()
 Stop polling. More...
 
- Public Member Functions inherited from XrdCl::Poller
virtual ~Poller ()
 Destructor. More...
 

Detailed Description

A poller implementation using the build-in XRootD poller.

Definition at line 40 of file XrdClPollerBuiltIn.hh.

Constructor & Destructor Documentation

◆ PollerBuiltIn()

XrdCl::PollerBuiltIn::PollerBuiltIn ( )
inline

Constructor.

Definition at line 46 of file XrdClPollerBuiltIn.hh.

46 : pNbPoller( GetNbPollerInit() ){}

◆ ~PollerBuiltIn()

XrdCl::PollerBuiltIn::~PollerBuiltIn ( )
inline

Definition at line 48 of file XrdClPollerBuiltIn.hh.

48 {}

Member Function Documentation

◆ AddSocket()

bool XrdCl::PollerBuiltIn::AddSocket ( Socket socket,
SocketHandler handler 
)
virtual

Add socket to the polling loop

Parameters
socketthe socket
handlerobject handling the events

Implements XrdCl::Poller.

Definition at line 348 of file XrdClPollerBuiltIn.cc.

350  {
351  Log *log = DefaultEnv::GetLog();
352  XrdSysMutexHelper scopedLock( pMutex );
353 
354  if( !socket )
355  {
356  log->Error( PollerMsg, "Invalid socket, impossible to poll" );
357  return false;
358  }
359 
360  if( socket->GetStatus() != Socket::Connected &&
361  socket->GetStatus() != Socket::Connecting )
362  {
363  log->Error( PollerMsg, "Socket is not in a state valid for polling" );
364  return false;
365  }
366 
367  log->Debug( PollerMsg, "Adding socket %p to the poller", (void*)socket );
368 
369  //--------------------------------------------------------------------------
370  // Check if the socket is already registered
371  //--------------------------------------------------------------------------
372  SocketMap::const_iterator it = pSocketMap.find( socket );
373  if( it != pSocketMap.end() )
374  {
375  log->Warning( PollerMsg, "%s Already registered with this poller",
376  socket->GetName().c_str() );
377  return false;
378  }
379 
380  //--------------------------------------------------------------------------
381  // Create the socket helper
382  //--------------------------------------------------------------------------
383  XrdSys::IOEvents::Poller* poller = RegisterAndGetPoller( socket );
384 
385  if( !poller )
386  {
387  log->Error( PollerMsg, "No poller available, can not add socket" );
388  return false;
389  }
390 
391  PollerHelper *helper = new PollerHelper();
392  helper->callBack = new ::SocketCallBack( socket, handler );
393 
394  if( poller )
395  {
396  helper->channel = new XrdSys::IOEvents::Channel( poller,
397  socket->GetFD(),
398  helper->callBack );
399  }
400 
401  handler->Initialize( this );
402  pSocketMap[socket] = helper;
403  return true;
404  }
static Log * GetLog()
Get default log.
Handle diagnostics.
Definition: XrdClLog.hh:101
void Error(uint64_t topic, const char *format,...)
Report an error.
Definition: XrdClLog.cc:231
void Warning(uint64_t topic, const char *format,...)
Report a warning.
Definition: XrdClLog.cc:248
void Debug(uint64_t topic, const char *format,...)
Print a debug message.
Definition: XrdClLog.cc:282
virtual void Initialize(Poller *)
Initializer.
Definition: XrdClPoller.hh:56
@ Connected
The socket is connected.
Definition: XrdClSocket.hh:51
@ Connecting
The connection process is in progress.
Definition: XrdClSocket.hh:52
SocketStatus GetStatus() const
Get the socket status.
Definition: XrdClSocket.hh:125
const uint64_t PollerMsg

References XrdCl::Socket::Connected, XrdCl::Socket::Connecting, XrdCl::Log::Debug(), XrdCl::Log::Error(), XrdCl::Socket::GetFD(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdCl::Socket::GetStatus(), XrdCl::SocketHandler::Initialize(), XrdCl::PollerMsg, and XrdCl::Log::Warning().

+ Here is the call graph for this function:

◆ EnableReadNotification()

bool XrdCl::PollerBuiltIn::EnableReadNotification ( Socket socket,
bool  notify,
time_t  timeout = 60 
)
virtual

Notify the handler about read events

Parameters
socketthe socket
notifyspecify if the handler should be notified
timeoutif no read event occurred after this time a timeout event will be generated

Implements XrdCl::Poller.

Definition at line 479 of file XrdClPollerBuiltIn.cc.

482  {
483  using namespace XrdSys::IOEvents;
484  Log *log = DefaultEnv::GetLog();
485 
486  if( !socket )
487  {
488  log->Error( PollerMsg, "Invalid socket, read events unavailable" );
489  return false;
490  }
491 
492  //--------------------------------------------------------------------------
493  // Check if the socket is registered
494  //--------------------------------------------------------------------------
495  XrdSysMutexHelper scopedLock( pMutex );
496  SocketMap::const_iterator it = pSocketMap.find( socket );
497  if( it == pSocketMap.end() )
498  {
499  log->Warning( PollerMsg, "%s Socket is not registered",
500  socket->GetName().c_str() );
501  return false;
502  }
503 
504  PollerHelper *helper = (PollerHelper*)it->second;
505  XrdSys::IOEvents::Poller *poller = GetPoller( socket );
506 
507  //--------------------------------------------------------------------------
508  // Enable read notifications
509  //--------------------------------------------------------------------------
510  if( notify )
511  {
512  if( helper->readEnabled )
513  return true;
514  helper->readTimeout = timeout;
515 
516  log->Dump( PollerMsg, "%s Enable read notifications, timeout: %lld",
517  socket->GetName().c_str(), (long long)timeout );
518 
519  if( poller )
520  {
521  const char *errMsg;
522  bool status = helper->channel->Enable( Channel::readEvents, timeout,
523  &errMsg );
524  if( !status )
525  {
526  log->Error( PollerMsg, "%s Unable to enable read notifications: %s",
527  socket->GetName().c_str(), errMsg );
528  return false;
529  }
530  }
531  helper->readEnabled = true;
532  }
533 
534  //--------------------------------------------------------------------------
535  // Disable read notifications
536  //--------------------------------------------------------------------------
537  else
538  {
539  if( !helper->readEnabled )
540  return true;
541 
542  log->Dump( PollerMsg, "%s Disable read notifications",
543  socket->GetName().c_str() );
544 
545  if( poller )
546  {
547  const char *errMsg;
548  bool status = helper->channel->Disable( Channel::readEvents, &errMsg );
549  if( !status )
550  {
551  log->Error( PollerMsg, "%s Unable to disable read notifications: %s",
552  socket->GetName().c_str(), errMsg );
553  return false;
554  }
555  }
556  helper->readEnabled = false;
557  }
558  return true;
559  }
void Dump(uint64_t topic, const char *format,...)
Print a dump message.
Definition: XrdClLog.cc:299
std::string GetName() const
Get the string representation of the socket.
Definition: XrdClSocket.cc:672

References XrdCl::Log::Dump(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdCl::PollerMsg, and XrdCl::Log::Warning().

+ Here is the call graph for this function:

◆ EnableWriteNotification()

bool XrdCl::PollerBuiltIn::EnableWriteNotification ( Socket socket,
bool  notify,
time_t  timeout = 60 
)
virtual

Notify the handler about write events

Parameters
socketthe socket
notifyspecify if the handler should be notified
timeoutif no write event occurred after this time a timeout event will be generated

Implements XrdCl::Poller.

Definition at line 564 of file XrdClPollerBuiltIn.cc.

567  {
568  using namespace XrdSys::IOEvents;
569  Log *log = DefaultEnv::GetLog();
570 
571  if( !socket )
572  {
573  log->Error( PollerMsg, "Invalid socket, write events unavailable" );
574  return false;
575  }
576 
577  //--------------------------------------------------------------------------
578  // Check if the socket is registered
579  //--------------------------------------------------------------------------
580  XrdSysMutexHelper scopedLock( pMutex );
581  SocketMap::const_iterator it = pSocketMap.find( socket );
582  if( it == pSocketMap.end() )
583  {
584  log->Warning( PollerMsg, "%s Socket is not registered",
585  socket->GetName().c_str() );
586  return false;
587  }
588 
589  PollerHelper *helper = (PollerHelper*)it->second;
590  XrdSys::IOEvents::Poller *poller = GetPoller( socket );
591 
592  //--------------------------------------------------------------------------
593  // Enable write notifications
594  //--------------------------------------------------------------------------
595  if( notify )
596  {
597  if( helper->writeEnabled )
598  return true;
599 
600  helper->writeTimeout = timeout;
601 
602  log->Dump( PollerMsg, "%s Enable write notifications, timeout: %lld",
603  socket->GetName().c_str(), (long long)timeout );
604 
605  if( poller )
606  {
607  const char *errMsg;
608  bool status = helper->channel->Enable( Channel::writeEvents, timeout,
609  &errMsg );
610  if( !status )
611  {
612  log->Error( PollerMsg, "%s Unable to enable write notifications: %s",
613  socket->GetName().c_str(), errMsg );
614  return false;
615  }
616  }
617  helper->writeEnabled = true;
618  }
619 
620  //--------------------------------------------------------------------------
621  // Disable read notifications
622  //--------------------------------------------------------------------------
623  else
624  {
625  if( !helper->writeEnabled )
626  return true;
627 
628  log->Dump( PollerMsg, "%s Disable write notifications",
629  socket->GetName().c_str() );
630  if( poller )
631  {
632  const char *errMsg;
633  bool status = helper->channel->Disable( Channel::writeEvents, &errMsg );
634  if( !status )
635  {
636  log->Error( PollerMsg, "%s Unable to disable write notifications: %s",
637  socket->GetName().c_str(), errMsg );
638  return false;
639  }
640  }
641  helper->writeEnabled = false;
642  }
643  return true;
644  }

References XrdCl::Log::Dump(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdCl::PollerMsg, and XrdCl::Log::Warning().

+ Here is the call graph for this function:

◆ Finalize()

bool XrdCl::PollerBuiltIn::Finalize ( )
virtual

Finalize the poller.

Implements XrdCl::Poller.

Definition at line 190 of file XrdClPollerBuiltIn.cc.

191  {
192  //--------------------------------------------------------------------------
193  // Clean up the channels
194  //--------------------------------------------------------------------------
195  SocketMap::iterator it;
196  for( it = pSocketMap.begin(); it != pSocketMap.end(); ++it )
197  {
198  PollerHelper *helper = (PollerHelper*)it->second;
199  if( helper->channel ) helper->channel->Delete();
200  delete helper->callBack;
201  delete helper;
202  }
203  pSocketMap.clear();
204 
205  return true;
206  }

◆ Initialize()

bool XrdCl::PollerBuiltIn::Initialize ( )
virtual

Initialize the poller.

Implements XrdCl::Poller.

Definition at line 182 of file XrdClPollerBuiltIn.cc.

183  {
184  return true;
185  }

◆ IsRegistered()

bool XrdCl::PollerBuiltIn::IsRegistered ( Socket socket)
virtual

Check whether the socket is registered with the poller.

Implements XrdCl::Poller.

Definition at line 649 of file XrdClPollerBuiltIn.cc.

650  {
651  XrdSysMutexHelper scopedLock( pMutex );
652  SocketMap::iterator it = pSocketMap.find( socket );
653  return it != pSocketMap.end();
654  }

◆ IsRunning()

virtual bool XrdCl::PollerBuiltIn::IsRunning ( ) const
inlinevirtual

Is the event loop running?

Implements XrdCl::Poller.

Definition at line 123 of file XrdClPollerBuiltIn.hh.

124  {
125  return !pPollerPool.empty();
126  }

◆ RemoveSocket()

bool XrdCl::PollerBuiltIn::RemoveSocket ( Socket socket)
virtual

Remove the socket.

Implements XrdCl::Poller.

Definition at line 433 of file XrdClPollerBuiltIn.cc.

434  {
435  using namespace XrdSys::IOEvents;
436  Log *log = DefaultEnv::GetLog();
437 
438  //--------------------------------------------------------------------------
439  // Find the right socket
440  //--------------------------------------------------------------------------
441  XrdSysMutexHelper scopedLock( pMutex );
442  SocketMap::iterator it = pSocketMap.find( socket );
443  if( it == pSocketMap.end() )
444  return true;
445 
446  log->Debug( PollerMsg, "%s Removing socket from the poller",
447  socket->GetName().c_str() );
448 
449  // unregister from the poller it's currently associated with
450  UnregisterFromPoller( socket );
451 
452  //--------------------------------------------------------------------------
453  // Remove the socket
454  //--------------------------------------------------------------------------
455  PollerHelper *helper = (PollerHelper*)it->second;
456  pSocketMap.erase( it );
457  scopedLock.UnLock();
458 
459  if( helper->channel )
460  {
461  const char *errMsg;
462  bool status = helper->channel->Disable( Channel::allEvents, &errMsg );
463  if( !status )
464  {
465  log->Error( PollerMsg, "%s Unable to disable write notifications: %s",
466  socket->GetName().c_str(), errMsg );
467  return false;
468  }
469  helper->channel->Delete();
470  }
471  delete helper->callBack;
472  delete helper;
473  return true;
474  }

References XrdCl::Log::Debug(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdCl::PollerMsg, and XrdSysMutexHelper::UnLock().

+ Here is the call graph for this function:

◆ ShutdownEvents()

void XrdCl::PollerBuiltIn::ShutdownEvents ( Socket socket)
virtual

Disables further callbacks to the socket's event handler. If callback is currently running wait for it, unless it is the current thread.

Implements XrdCl::Poller.

Definition at line 412 of file XrdClPollerBuiltIn.cc.

413  {
414  XrdSysMutexHelper scopedLock( pMutex );
415  SocketMap::iterator it = pSocketMap.find( socket );
416  if( it == pSocketMap.end() )
417  return;
418 
419  PollerHelper *helper = (PollerHelper*)it->second;
420  if( !helper ) return;
421  XrdSys::IOEvents::CallBack *cb = helper->callBack;
422  if( !cb ) return;
423  SocketCallBack *scb = dynamic_cast<SocketCallBack*>( cb );
424  if( !scb ) return;
425  auto dc = scb->GetControl();
426  scopedLock.UnLock();
427  dc->DisableCallBack();
428  }

References XrdSysMutexHelper::UnLock().

+ Here is the call graph for this function:

◆ Start()

bool XrdCl::PollerBuiltIn::Start ( )
virtual

Start polling.

Implements XrdCl::Poller.

Definition at line 211 of file XrdClPollerBuiltIn.cc.

212  {
213  //--------------------------------------------------------------------------
214  // Start the poller
215  //--------------------------------------------------------------------------
216  using namespace XrdSys;
217 
218  Log *log = DefaultEnv::GetLog();
219  log->Debug( PollerMsg, "Creating and starting the built-in poller..." );
220  XrdSysMutexHelper scopedLock( pMutex );
221  int errNum = 0;
222  const char *errMsg = 0;
223 
224  for( int i = 0; i < pNbPoller; ++i )
225  {
226  XrdSys::IOEvents::Poller* poller = IOEvents::Poller::Create( errNum, &errMsg );
227  if( !poller )
228  {
229  log->Error( PollerMsg, "Unable to create the internal poller object: "
230  "%s (%s)", XrdSysE2T( errno ), errMsg );
231  return false;
232  }
233  pPollerPool.push_back( poller );
234  }
235 
236  pNext = pPollerPool.begin();
237 
238  log->Debug( PollerMsg, "Using %d poller threads", pNbPoller );
239 
240  //--------------------------------------------------------------------------
241  // Check if we have any descriptors to reinsert from the last time we
242  // were started
243  //--------------------------------------------------------------------------
244  SocketMap::iterator it;
245  for( it = pSocketMap.begin(); it != pSocketMap.end(); ++it )
246  {
247  PollerHelper *helper = (PollerHelper*)it->second;
248  Socket *socket = it->first;
249 
250  // static cast for the downcast as we're sure it's a SocketCallBack
251  auto *scb = static_cast<SocketCallBack*>( helper->callBack );
252  if( scb )
253  {
254  auto dc = scb->GetControl();
255  if( dc ) dc->Restart();
256  }
257 
258  helper->channel = new IOEvents::Channel( RegisterAndGetPoller( socket ), socket->GetFD(),
259  helper->callBack );
260  if( helper->readEnabled )
261  {
262  bool status = helper->channel->Enable( IOEvents::Channel::readEvents,
263  helper->readTimeout, &errMsg );
264  if( !status )
265  {
266  log->Error( PollerMsg, "Unable to enable read notifications "
267  "while re-starting %s (%s)", XrdSysE2T( errno ), errMsg );
268 
269  return false;
270  }
271  }
272 
273  if( helper->writeEnabled )
274  {
275  bool status = helper->channel->Enable( IOEvents::Channel::writeEvents,
276  helper->writeTimeout, &errMsg );
277  if( !status )
278  {
279  log->Error( PollerMsg, "Unable to enable write notifications "
280  "while re-starting %s (%s)", XrdSysE2T( errno ), errMsg );
281 
282  return false;
283  }
284  }
285  }
286  return true;
287  }
bool Create
const char * XrdSysE2T(int errcode)
Definition: XrdSysE2T.cc:104
A network socket.
Definition: XrdClSocket.hh:43

References Create, XrdCl::Log::Debug(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::PollerMsg, and XrdSysE2T().

+ Here is the call graph for this function:

◆ Stop()

bool XrdCl::PollerBuiltIn::Stop ( )
virtual

Stop polling.

Implements XrdCl::Poller.

Definition at line 292 of file XrdClPollerBuiltIn.cc.

293  {
294  using namespace XrdSys::IOEvents;
295 
296  Log *log = DefaultEnv::GetLog();
297  log->Debug( PollerMsg, "Stopping the poller..." );
298 
299  XrdSysMutexHelper scopedLock( pMutex );
300 
301  if( pPollerPool.empty() )
302  {
303  log->Debug( PollerMsg, "Stopping a poller that has not been started" );
304  return true;
305  }
306 
307  while( !pPollerPool.empty() )
308  {
309  XrdSys::IOEvents::Poller *poller = pPollerPool.back();
310  if( *pNext == poller )
311  pNext = pPollerPool.begin();
312  pPollerPool.pop_back();
313 
314  if( !poller ) continue;
315 
316  scopedLock.UnLock();
317  poller->Stop();
318  delete poller;
319  scopedLock.Lock( &pMutex );
320  }
321  pNext = pPollerPool.end();
322  pPollerMap.clear();
323 
324  SocketMap::iterator it;
325  const char *errMsg = 0;
326 
327  for( it = pSocketMap.begin(); it != pSocketMap.end(); ++it )
328  {
329  PollerHelper *helper = (PollerHelper*)it->second;
330  if( !helper->channel ) continue;
331  bool status = helper->channel->Disable( Channel::allEvents, &errMsg );
332  if( !status )
333  {
334  Socket *socket = it->first;
335  log->Error( PollerMsg, "%s Unable to disable write notifications: %s",
336  socket->GetName().c_str(), errMsg );
337  }
338  helper->channel->Delete();
339  helper->channel = 0;
340  }
341 
342  return true;
343  }

References XrdCl::Log::Debug(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdSysMutexHelper::Lock(), XrdCl::PollerMsg, XrdSys::IOEvents::Poller::Stop(), and XrdSysMutexHelper::UnLock().

+ Here is the call graph for this function:

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