rippled
Classes | Public Types | Public Member Functions | Protected Member Functions | Private Types | Private Member Functions | Private Attributes | Friends | List of all members
ripple::JobQueue Class Reference

A pool of threads to perform work. More...

Inheritance diagram for ripple::JobQueue:
Inheritance graph
[legend]
Collaboration diagram for ripple::JobQueue:
Collaboration graph
[legend]

Classes

class  Coro
 Coroutines must run to completion. More...
 

Public Types

using JobFunction = std::function< void(Job &)>
 

Public Member Functions

 JobQueue (beast::insight::Collector::ptr const &collector, Stoppable &parent, beast::Journal journal, Logs &logs, perf::PerfLog &perfLog)
 
 ~JobQueue ()
 
template<typename JobHandler >
bool addJob (JobType type, std::string const &name, JobHandler &&jobHandler)
 Adds a job to the JobQueue. More...
 
template<class F >
std::shared_ptr< CoropostCoro (JobType t, std::string const &name, F &&f)
 Creates a coroutine and adds a job to the queue which will run it. More...
 
int getJobCount (JobType t) const
 Jobs waiting at this priority. More...
 
int getJobCountTotal (JobType t) const
 Jobs waiting plus running at this priority. More...
 
int getJobCountGE (JobType t) const
 All waiting jobs at or greater than this priority. More...
 
void setThreadCount (int c, bool const standaloneMode)
 Set the number of thread serving the job queue to precisely this number. More...
 
std::unique_ptr< LoadEventmakeLoadEvent (JobType t, std::string const &name)
 Return a scoped LoadEvent. More...
 
void addLoadEvents (JobType t, int count, std::chrono::milliseconds elapsed)
 Add multiple load events. More...
 
bool isOverloaded ()
 
Json::Value getJson (int c=0)
 
void rendezvous ()
 Block until no tasks running. More...
 
RootStoppablegetRoot ()
 
void setParent (Stoppable &parent)
 Set the parent of this Stoppable. More...
 
bool isStopping () const
 Returns true if the stoppable should stop. More...
 
bool isStopped () const
 Returns true if the requested stop has completed. More...
 
bool areChildrenStopped () const
 Returns true if all children have stopped. More...
 
JobCounterjobCounter ()
 
bool alertable_sleep_until (std::chrono::system_clock::time_point const &t)
 Sleep or wake up on stop. More...
 

Protected Member Functions

void stopped ()
 Called by derived classes to indicate that the stoppable has stopped. More...
 

Private Types

using JobDataMap = std::map< JobType, JobTypeData >
 
using Children = beast::LockFreeStack< Child >
 

Private Member Functions

void collect ()
 
JobTypeDatagetJobTypeData (JobType type)
 
void onStop () override
 Override called when the stop notification is issued. More...
 
void checkStopped (std::lock_guard< std::mutex > const &lock)
 
bool addRefCountedJob (JobType type, std::string const &name, JobFunction const &func)
 
void queueJob (Job const &job, std::lock_guard< std::mutex > const &lock)
 
void getNextJob (Job &job)
 
void finishJob (JobType type)
 
void processTask (int instance) override
 Perform a task. More...
 
int getJobLimit (JobType type)
 
void onChildrenStopped () override
 Override called when all children have stopped. More...
 
virtual void onPrepare ()
 Override called during preparation. More...
 
virtual void onStart ()
 Override called during start. More...
 
void prepareRecursive ()
 
void startRecursive ()
 
void stopAsyncRecursive (beast::Journal j)
 
void stopRecursive (beast::Journal j)
 

Private Attributes

beast::Journal m_journal
 
std::mutex m_mutex
 
std::uint64_t m_lastJob
 
std::set< Jobm_jobSet
 
JobDataMap m_jobData
 
JobTypeData m_invalidJobData
 
int m_processCount
 
int nSuspend_ = 0
 
Workers m_workers
 
Job::CancelCallback m_cancelCallback
 
perf::PerfLogperfLog_
 
beast::insight::Collector::ptr m_collector
 
beast::insight::Gauge job_count
 
beast::insight::Hook hook
 
std::condition_variable cv_
 
std::string m_name
 
RootStoppablem_root
 
Child m_child
 
std::atomic< bool > m_stopped {false}
 
std::atomic< bool > m_childrenStopped {false}
 
Children m_children
 
std::condition_variable m_cv
 
std::mutex m_mut
 
bool m_is_stopping = false
 
bool hasParent_ {false}
 

Friends

class Coro
 

Detailed Description

A pool of threads to perform work.

A job posted will always run to completion.

Coroutines that are suspended must be resumed, and run to completion.

When the JobQueue stops, it waits for all jobs and coroutines to finish.

Definition at line 55 of file JobQueue.h.

Member Typedef Documentation

◆ JobFunction

Definition at line 141 of file JobQueue.h.

◆ JobDataMap

Definition at line 234 of file JobQueue.h.

◆ Children

Definition at line 319 of file Stoppable.h.

Constructor & Destructor Documentation

◆ JobQueue()

ripple::JobQueue::JobQueue ( beast::insight::Collector::ptr const &  collector,
Stoppable parent,
beast::Journal  journal,
Logs logs,
perf::PerfLog perfLog 
)

Definition at line 26 of file JobQueue.cpp.

◆ ~JobQueue()

ripple::JobQueue::~JobQueue ( )

Definition at line 63 of file JobQueue.cpp.

Member Function Documentation

◆ addJob()

template<typename JobHandler >
bool ripple::JobQueue::addJob ( JobType  type,
std::string const &  name,
JobHandler &&  jobHandler 
)

Adds a job to the JobQueue.

Parameters
typeThe type of job.
nameName of the job.
jobHandlerLambda with signature void (Job&). Called when the job is executed.
Returns
true if jobHandler added to queue.

Definition at line 166 of file JobQueue.h.

◆ postCoro()

template<class F >
std::shared_ptr< JobQueue::Coro > ripple::JobQueue::postCoro ( JobType  t,
std::string const &  name,
F &&  f 
)

Creates a coroutine and adds a job to the queue which will run it.

Parameters
tThe type of job.
nameName of the job.
fHas a signature of void(std::shared_ptr<Coro>). Called when the job executes.
Returns
shared_ptr to posted Coro. nullptr if post was not successful.

Definition at line 427 of file JobQueue.h.

◆ getJobCount()

int ripple::JobQueue::getJobCount ( JobType  t) const

Jobs waiting at this priority.

Definition at line 121 of file JobQueue.cpp.

◆ getJobCountTotal()

int ripple::JobQueue::getJobCountTotal ( JobType  t) const

Jobs waiting plus running at this priority.

Definition at line 131 of file JobQueue.cpp.

◆ getJobCountGE()

int ripple::JobQueue::getJobCountGE ( JobType  t) const

All waiting jobs at or greater than this priority.

Definition at line 141 of file JobQueue.cpp.

◆ setThreadCount()

void ripple::JobQueue::setThreadCount ( int  c,
bool const  standaloneMode 
)

Set the number of thread serving the job queue to precisely this number.

Definition at line 158 of file JobQueue.cpp.

◆ makeLoadEvent()

std::unique_ptr< LoadEvent > ripple::JobQueue::makeLoadEvent ( JobType  t,
std::string const &  name 
)

Return a scoped LoadEvent.

Definition at line 181 of file JobQueue.cpp.

◆ addLoadEvents()

void ripple::JobQueue::addLoadEvents ( JobType  t,
int  count,
std::chrono::milliseconds  elapsed 
)

Add multiple load events.

Definition at line 193 of file JobQueue.cpp.

◆ isOverloaded()

bool ripple::JobQueue::isOverloaded ( )

Definition at line 204 of file JobQueue.cpp.

◆ getJson()

Json::Value ripple::JobQueue::getJson ( int  c = 0)

Definition at line 218 of file JobQueue.cpp.

◆ rendezvous()

void ripple::JobQueue::rendezvous ( )

Block until no tasks running.

Definition at line 276 of file JobQueue.cpp.

◆ collect()

void ripple::JobQueue::collect ( )
private

Definition at line 70 of file JobQueue.cpp.

◆ getJobTypeData()

JobTypeData & ripple::JobQueue::getJobTypeData ( JobType  type)
private

Definition at line 283 of file JobQueue.cpp.

◆ onStop()

void ripple::JobQueue::onStop ( )
overrideprivatevirtual

Override called when the stop notification is issued.

The call is made on an unspecified, implementation-specific thread. onStop and onChildrenStopped will never be called concurrently, across all Stoppable objects descended from the same root, inclusive of the root.

It is safe to call isStopping, isStopped, and areChildrenStopped from within this function; The values returned will always be valid and never change during the callback.

The default implementation simply calls stopped(). This is applicable when the Stoppable has a trivial stop operation (or no stop operation), and we are merely using the Stoppable API to position it as a dependency of some parent service.

Thread safety: May not block for long periods. Guaranteed only to be called once. Must be safe to call from any thread at any time.

Reimplemented from ripple::Stoppable.

Definition at line 297 of file JobQueue.cpp.

◆ checkStopped()

void ripple::JobQueue::checkStopped ( std::lock_guard< std::mutex > const &  lock)
private

Definition at line 304 of file JobQueue.cpp.

◆ addRefCountedJob()

bool ripple::JobQueue::addRefCountedJob ( JobType  type,
std::string const &  name,
JobFunction const &  func 
)
private

Definition at line 77 of file JobQueue.cpp.

◆ queueJob()

void ripple::JobQueue::queueJob ( Job const &  job,
std::lock_guard< std::mutex > const &  lock 
)
private

Definition at line 322 of file JobQueue.cpp.

◆ getNextJob()

void ripple::JobQueue::getNextJob ( Job job)
private

Definition at line 345 of file JobQueue.cpp.

◆ finishJob()

void ripple::JobQueue::finishJob ( JobType  type)
private

Definition at line 379 of file JobQueue.cpp.

◆ processTask()

void ripple::JobQueue::processTask ( int  instance)
overrideprivatevirtual

Perform a task.

The call is made on a thread owned by Workers. It is important that you only process one task from inside your callback. Each call to addTask will result in exactly one call to processTask.

Parameters
instanceThe worker thread instance.
See also
Workers::addTask

Implements ripple::Workers::Callback.

Definition at line 398 of file JobQueue.cpp.

◆ getJobLimit()

int ripple::JobQueue::getJobLimit ( JobType  type)
private

Definition at line 452 of file JobQueue.cpp.

◆ onChildrenStopped()

void ripple::JobQueue::onChildrenStopped ( )
overrideprivatevirtual

Override called when all children have stopped.

The call is made on an unspecified, implementation-specific thread. onStop and onChildrenStopped will never be called concurrently, across all Stoppable objects descended from the same root, inclusive of the root.

It is safe to call isStopping, isStopped, and areChildrenStopped from within this function; The values returned will always be valid and never change during the callback.

The default implementation does nothing.

Thread safety: May not block for long periods. Guaranteed only to be called once. Must be safe to call from any thread at any time.

Reimplemented from ripple::Stoppable.

Definition at line 461 of file JobQueue.cpp.

◆ getRoot()

RootStoppable& ripple::Stoppable::getRoot ( )
inherited

Definition at line 214 of file Stoppable.h.

◆ setParent()

void ripple::Stoppable::setParent ( Stoppable parent)
inherited

Set the parent of this Stoppable.

Note
The Stoppable must not already have a parent. The parent to be set cannot not be stopping. Both roots must match.

Definition at line 43 of file Stoppable.cpp.

◆ isStopping()

bool ripple::Stoppable::isStopping ( ) const
inherited

Returns true if the stoppable should stop.

Definition at line 54 of file Stoppable.cpp.

◆ isStopped()

bool ripple::Stoppable::isStopped ( ) const
inherited

Returns true if the requested stop has completed.

Definition at line 60 of file Stoppable.cpp.

◆ areChildrenStopped()

bool ripple::Stoppable::areChildrenStopped ( ) const
inherited

Returns true if all children have stopped.

Definition at line 66 of file Stoppable.cpp.

◆ jobCounter()

JobCounter & ripple::Stoppable::jobCounter ( )
inherited

Definition at line 437 of file Stoppable.h.

◆ alertable_sleep_until()

bool ripple::Stoppable::alertable_sleep_until ( std::chrono::system_clock::time_point const &  t)
inherited

Sleep or wake up on stop.

Returns
true if we are stopping

Definition at line 455 of file Stoppable.h.

◆ stopped()

void ripple::Stoppable::stopped ( )
protectedinherited

Called by derived classes to indicate that the stoppable has stopped.

Definition at line 72 of file Stoppable.cpp.

◆ onPrepare()

void ripple::Stoppable::onPrepare ( )
privatevirtualinherited

Override called during preparation.

Since all other Stoppable objects in the tree have already been constructed, this provides an opportunity to perform initialization which depends on calling into other Stoppable objects. This call is made on the same thread that called prepare(). The default implementation does nothing. Guaranteed to only be called once.

Reimplemented in ripple::ApplicationImp, ripple::OverlayImpl, ripple::test::Stoppable_test::Root, ripple::test::Stoppable_test::C, ripple::test::Stoppable_test::I, ripple::test::Stoppable_test::B, ripple::test::Stoppable_test::H, ripple::test::Stoppable_test::G, ripple::SHAMapStoreImp, ripple::test::Stoppable_test::A, ripple::PeerFinder::ManagerImp, ripple::perf::PerfLogImp, ripple::test::Stoppable_test::F, ripple::test::Stoppable_test::E, ripple::detail::LedgerCleanerImp, ripple::test::Stoppable_test::J, ripple::LoadManager, ripple::PerfLog_test::PerfLogParent, and ripple::test::Stoppable_test::D.

Definition at line 80 of file Stoppable.cpp.

◆ onStart()

void ripple::Stoppable::onStart ( )
privatevirtualinherited

◆ prepareRecursive()

void ripple::Stoppable::prepareRecursive ( )
privateinherited

Definition at line 103 of file Stoppable.cpp.

◆ startRecursive()

void ripple::Stoppable::startRecursive ( )
privateinherited

Definition at line 113 of file Stoppable.cpp.

◆ stopAsyncRecursive()

void ripple::Stoppable::stopAsyncRecursive ( beast::Journal  j)
privateinherited

Definition at line 123 of file Stoppable.cpp.

◆ stopRecursive()

void ripple::Stoppable::stopRecursive ( beast::Journal  j)
privateinherited

Definition at line 134 of file Stoppable.cpp.

Friends And Related Function Documentation

◆ Coro

friend class Coro
friend

Definition at line 232 of file JobQueue.h.

Member Data Documentation

◆ m_journal

beast::Journal ripple::JobQueue::m_journal
private

Definition at line 236 of file JobQueue.h.

◆ m_mutex

std::mutex ripple::JobQueue::m_mutex
mutableprivate

Definition at line 237 of file JobQueue.h.

◆ m_lastJob

std::uint64_t ripple::JobQueue::m_lastJob
private

Definition at line 238 of file JobQueue.h.

◆ m_jobSet

std::set<Job> ripple::JobQueue::m_jobSet
private

Definition at line 239 of file JobQueue.h.

◆ m_jobData

JobDataMap ripple::JobQueue::m_jobData
private

Definition at line 240 of file JobQueue.h.

◆ m_invalidJobData

JobTypeData ripple::JobQueue::m_invalidJobData
private

Definition at line 241 of file JobQueue.h.

◆ m_processCount

int ripple::JobQueue::m_processCount
private

Definition at line 244 of file JobQueue.h.

◆ nSuspend_

int ripple::JobQueue::nSuspend_ = 0
private

Definition at line 247 of file JobQueue.h.

◆ m_workers

Workers ripple::JobQueue::m_workers
private

Definition at line 249 of file JobQueue.h.

◆ m_cancelCallback

Job::CancelCallback ripple::JobQueue::m_cancelCallback
private

Definition at line 250 of file JobQueue.h.

◆ perfLog_

perf::PerfLog& ripple::JobQueue::perfLog_
private

Definition at line 253 of file JobQueue.h.

◆ m_collector

beast::insight::Collector::ptr ripple::JobQueue::m_collector
private

Definition at line 254 of file JobQueue.h.

◆ job_count

beast::insight::Gauge ripple::JobQueue::job_count
private

Definition at line 255 of file JobQueue.h.

◆ hook

beast::insight::Hook ripple::JobQueue::hook
private

Definition at line 256 of file JobQueue.h.

◆ cv_

std::condition_variable ripple::JobQueue::cv_
private

Definition at line 258 of file JobQueue.h.

◆ m_name

std::string ripple::Stoppable::m_name
privateinherited

Definition at line 339 of file Stoppable.h.

◆ m_root

RootStoppable& ripple::Stoppable::m_root
privateinherited

Definition at line 340 of file Stoppable.h.

◆ m_child

Child ripple::Stoppable::m_child
privateinherited

Definition at line 341 of file Stoppable.h.

◆ m_stopped

std::atomic<bool> ripple::Stoppable::m_stopped {false}
privateinherited

Definition at line 342 of file Stoppable.h.

◆ m_childrenStopped

std::atomic<bool> ripple::Stoppable::m_childrenStopped {false}
privateinherited

Definition at line 343 of file Stoppable.h.

◆ m_children

Children ripple::Stoppable::m_children
privateinherited

Definition at line 344 of file Stoppable.h.

◆ m_cv

std::condition_variable ripple::Stoppable::m_cv
privateinherited

Definition at line 345 of file Stoppable.h.

◆ m_mut

std::mutex ripple::Stoppable::m_mut
privateinherited

Definition at line 346 of file Stoppable.h.

◆ m_is_stopping

bool ripple::Stoppable::m_is_stopping = false
privateinherited

Definition at line 347 of file Stoppable.h.

◆ hasParent_

bool ripple::Stoppable::hasParent_ {false}
privateinherited

Definition at line 348 of file Stoppable.h.