xrpld
Toggle main menu visibility
Loading...
Searching...
No Matches
include
xrpl
core
detail
Workers.h
1
#pragma once
2
3
#include <xrpl/beast/core/LockFreeStack.h>
4
#include <xrpl/core/detail/semaphore.h>
5
6
#include <
atomic
>
7
#include <
condition_variable
>
8
#include <
mutex
>
9
#include <
string
>
10
#include <
thread
>
11
12
namespace
xrpl
{
13
14
namespace
perf
{
15
class
PerfLog
;
16
}
// namespace perf
17
60
class
Workers
61
{
62
public
:
66
struct
Callback
67
{
68
virtual
~Callback
() =
default
;
69
Callback
() =
default
;
70
Callback
(
Callback
const
&) =
delete
;
71
Callback
&
72
operator=
(
Callback
const
&) =
delete
;
73
85
virtual
void
86
processTask
(
int
instance) = 0;
87
};
88
97
explicit
Workers
(
98
Callback
& callback,
99
perf::PerfLog
* perfLog,
100
std::string
threadNames =
"Worker"
,
101
int
numberOfThreads =
static_cast<
int
>
(
std::thread::hardware_concurrency
()));
102
103
~Workers
();
104
114
[[nodiscard]]
int
115
getNumberOfThreads
() const noexcept;
116
121
void
122
setNumberOfThreads
(
int
numberOfThreads);
123
133
void
134
stop
();
135
145
void
146
addTask
();
147
153
[[nodiscard]]
int
154
numberOfCurrentlyRunningTasks
() const noexcept;
155
156
//--------------------------------------------------------------------------
157
158
private:
159
struct
PausedTag
160
{
161
explicit
PausedTag
() =
default
;
162
};
163
164
/* A Worker executes tasks on its provided thread.
165
166
These are the states:
167
168
Active: Running the task processing loop.
169
Idle: Active, but blocked on waiting for a task.
170
Paused: Blocked waiting to exit or become active.
171
*/
172
class
Worker
:
public
beast::LockFreeStack
<Worker>
::Node
,
173
public
beast::LockFreeStack
<Worker, PausedTag>
::Node
174
{
175
public
:
176
Worker
(
Workers
& workers,
std::string
threadName,
int
const
instance);
177
178
~Worker
();
179
180
void
181
notify
();
182
183
private
:
184
void
185
run
();
186
187
private
:
188
Workers
&
workers_
;
189
std::string
const
threadName_
;
190
int
const
instance_
;
191
192
std::thread
thread_
;
193
std::mutex
mutex_
;
194
std::condition_variable
wakeup_
;
195
int
wakeCount_
{0};
// how many times to un-pause
196
bool
shouldExit_
{
false
};
197
};
198
199
private
:
200
static
void
201
deleteWorkers
(
beast::LockFreeStack<Worker>
& stack);
202
203
private
:
204
Callback
&
callback_
;
205
perf::PerfLog
*
perfLog_
;
206
std::string
threadNames_
;
// The name to give each thread
207
std::condition_variable
cv_
;
// signaled when all threads paused
208
std::mutex
mut_
;
209
bool
allPaused_
{
true
};
210
Semaphore
semaphore_
;
// each pending task is 1 resource
211
int
numberOfThreads_
{0};
// how many we want active now
212
std::atomic<int>
activeCount_
;
// to know when all are paused
213
std::atomic<int>
pauseCount_
;
// how many threads need to pause now
214
std::atomic<int>
runningTaskCount_
;
// how many calls to processTask() active
215
beast::LockFreeStack<Worker>
everyone_
;
// holds all created workers
216
beast::LockFreeStack<Worker, PausedTag>
paused_
;
// holds just paused workers
217
};
218
219
}
// namespace xrpl
atomic
std::string
beast::LockFreeStack::Node::Node
Node()
Definition
LockFreeStack.h:126
beast::LockFreeStack
Multiple Producer, Multiple Consumer (MPMC) intrusive stack.
Definition
LockFreeStack.h:121
xrpl::Workers::Worker::workers_
Workers & workers_
Definition
Workers.h:188
xrpl::Workers::Worker::~Worker
~Worker()
Definition
Workers.cpp:147
xrpl::Workers::Worker::run
void run()
Definition
Workers.cpp:168
xrpl::Workers::Worker::threadName_
std::string const threadName_
Definition
Workers.h:189
xrpl::Workers::Worker::thread_
std::thread thread_
Definition
Workers.h:192
xrpl::Workers::Worker::mutex_
std::mutex mutex_
Definition
Workers.h:193
xrpl::Workers::Worker::instance_
int const instance_
Definition
Workers.h:190
xrpl::Workers::Worker::wakeCount_
int wakeCount_
Definition
Workers.h:195
xrpl::Workers::Worker::notify
void notify()
Definition
Workers.cpp:160
xrpl::Workers::Worker::Worker
Worker(Workers &workers, std::string threadName, int const instance)
Definition
Workers.cpp:140
xrpl::Workers::Worker::wakeup_
std::condition_variable wakeup_
Definition
Workers.h:194
xrpl::Workers::Worker::shouldExit_
bool shouldExit_
Definition
Workers.h:196
xrpl::Workers::setNumberOfThreads
void setNumberOfThreads(int numberOfThreads)
Set the desired number of threads.
Definition
Workers.cpp:47
xrpl::Workers::semaphore_
Semaphore semaphore_
Definition
Workers.h:210
xrpl::Workers::stop
void stop()
Pause all threads and wait until they are paused.
Definition
Workers.cpp:98
xrpl::Workers::numberOfThreads_
int numberOfThreads_
Definition
Workers.h:211
xrpl::Workers::activeCount_
std::atomic< int > activeCount_
Definition
Workers.h:212
xrpl::Workers::mut_
std::mutex mut_
Definition
Workers.h:208
xrpl::Workers::paused_
beast::LockFreeStack< Worker, PausedTag > paused_
Definition
Workers.h:216
xrpl::Workers::numberOfCurrentlyRunningTasks
int numberOfCurrentlyRunningTasks() const noexcept
Get the number of currently executing calls of Callback::processTask.
Definition
Workers.cpp:118
xrpl::Workers::pauseCount_
std::atomic< int > pauseCount_
Definition
Workers.h:213
xrpl::Workers::callback_
Callback & callback_
Definition
Workers.h:204
xrpl::Workers::addTask
void addTask()
Add a task to be performed.
Definition
Workers.cpp:112
xrpl::Workers::~Workers
~Workers()
Definition
Workers.cpp:29
xrpl::Workers::Workers
Workers(Callback &callback, perf::PerfLog *perfLog, std::string threadNames="Worker", int numberOfThreads=static_cast< int >(std::thread::hardware_concurrency()))
Create the object.
Definition
Workers.cpp:13
xrpl::Workers::perfLog_
perf::PerfLog * perfLog_
Definition
Workers.h:205
xrpl::Workers::deleteWorkers
static void deleteWorkers(beast::LockFreeStack< Worker > &stack)
Definition
Workers.cpp:124
xrpl::Workers::threadNames_
std::string threadNames_
Definition
Workers.h:206
xrpl::Workers::allPaused_
bool allPaused_
Definition
Workers.h:209
xrpl::Workers::cv_
std::condition_variable cv_
Definition
Workers.h:207
xrpl::Workers::runningTaskCount_
std::atomic< int > runningTaskCount_
Definition
Workers.h:214
xrpl::Workers::everyone_
beast::LockFreeStack< Worker > everyone_
Definition
Workers.h:215
xrpl::Workers::getNumberOfThreads
int getNumberOfThreads() const noexcept
Retrieve the desired number of threads.
Definition
Workers.cpp:37
xrpl::perf::PerfLog
Singleton class that maintains performance counters and optionally writes Json-formatted data to a di...
Definition
PerfLog.h:33
condition_variable
std::thread::hardware_concurrency
T hardware_concurrency(T... args)
mutex
xrpl::perf
Dummy class for unit tests.
Definition
Workers.h:14
xrpl
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition
algorithm.h:5
xrpl::Semaphore
BasicSemaphore< std::mutex, std::condition_variable > Semaphore
Definition
semaphore.h:94
string
xrpl::Workers::Callback
Called to perform tasks as needed.
Definition
Workers.h:67
xrpl::Workers::Callback::Callback
Callback(Callback const &)=delete
xrpl::Workers::Callback::Callback
Callback()=default
xrpl::Workers::Callback::processTask
virtual void processTask(int instance)=0
Perform a task.
xrpl::Workers::Callback::operator=
Callback & operator=(Callback const &)=delete
xrpl::Workers::Callback::~Callback
virtual ~Callback()=default
xrpl::Workers::PausedTag::PausedTag
PausedTag()=default
thread
Generated by
1.17.0