XRootD
Loading...
Searching...
No Matches
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.
 ~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.
virtual bool Initialize ()
 Initialize the poller.
virtual bool IsRegistered (Socket *socket)
 Check whether the socket is registered with the poller.
virtual bool IsRunning () const
 Is the event loop running?
virtual bool RemoveSocket (Socket *socket)
 Remove the socket.
virtual void ShutdownEvents (Socket *socket)
virtual bool Start ()
 Start polling.
virtual bool Stop ()
 Stop polling.
Public Member Functions inherited from XrdCl::Poller
virtual ~Poller ()
 Destructor.

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.
@ Connected
The socket is connected.
@ Connecting
The connection process is in progress.
const uint64_t PollerMsg
XrdSysError Log
Definition XrdConfig.cc:113

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 }

References XrdSys::IOEvents::Channel::Disable(), XrdCl::Log::Dump(), XrdSys::IOEvents::Channel::Enable(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdCl::PollerMsg, XrdSys::IOEvents::Channel::readEvents, 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 XrdSys::IOEvents::Channel::Disable(), XrdCl::Log::Dump(), XrdSys::IOEvents::Channel::Enable(), XrdCl::Log::Error(), XrdCl::DefaultEnv::GetLog(), XrdCl::Socket::GetName(), XrdCl::PollerMsg, XrdCl::Log::Warning(), and XrdSys::IOEvents::Channel::writeEvents.

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 }

References XrdSys::IOEvents::Channel::Delete().

Here is the call graph for this function:

◆ 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 XrdSys::IOEvents::Channel::allEvents, XrdCl::Log::Debug(), XrdSys::IOEvents::Channel::Delete(), XrdSys::IOEvents::Channel::Disable(), 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 }
const char * XrdSysE2T(int errcode)
Definition XrdSysE2T.cc:104

References XrdSys::IOEvents::Poller::Create(), XrdCl::Log::Debug(), XrdSys::IOEvents::Channel::Enable(), XrdCl::Log::Error(), XrdCl::Socket::GetFD(), XrdCl::DefaultEnv::GetLog(), XrdCl::PollerMsg, XrdSys::IOEvents::Channel::readEvents, XrdSys::IOEvents::Channel::writeEvents, 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 XrdSys::IOEvents::Channel::allEvents, XrdCl::Log::Debug(), XrdSys::IOEvents::Channel::Delete(), XrdSys::IOEvents::Channel::Disable(), 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: