20#include "pedigree/kernel/LockGuard.h"
21#include "pedigree/kernel/Log.h"
22#include "pedigree/kernel/machine/Machine.h"
23#include "pedigree/kernel/machine/Timer.h"
24#include "pedigree/kernel/process/PerProcessorScheduler.h"
25#include "pedigree/kernel/process/Scheduler.h"
26#include "pedigree/kernel/process/TerminationDeferral.h"
27#include "pedigree/kernel/process/Thread.h"
28#include "pedigree/kernel/processor/Processor.h"
29#include "pedigree/kernel/processor/ProcessorInformation.h"
30#include "pedigree/kernel/time/Time.h"
31#include "pedigree/kernel/utilities/RequestQueue.h"
32#include "pedigree/kernel/utilities/assert.h"
33#include "pedigree/kernel/utilities/new"
37static_assert(__atomic_always_lock_free(
sizeof(
size_t),
nullptr),
38 "RequestQueue preallocated-publication words must be lock-free");
40 "RequestQueue preallocated-publication pointers must be lock-free");
50#if defined(PEDIGREE_BUILDUTILS)
69 m_ReleasingToken(releasingToken),
82 m_StateLevel = __atomic_load_n(&m_Thread->
m_nStateLevel, __ATOMIC_ACQUIRE);
86 m_Thread->armAtomicStateCleanup(m_Cleanup, abandon,
this);
100 m_Thread->disarmAtomicStateCleanup(m_Cleanup);
106 if (!thread || !queue) {
113 if (scope->m_Queue == queue) {
116 scope = scope->m_Previous;
121 static bool allowsReleaseHandoff(
const RequestQueue* queue,
124 if (!thread || !queue || !token) {
131 if (scope->m_Queue == queue && scope->m_ReleasingToken == token) {
134 scope = scope->m_Previous;
143 static void abandon(
void* context) {
154 if (m_StateLevel >= MAX_NESTED_EVENTS ||
156 FATAL_NOLOCK(
"RequestQueue callback scope stack was corrupted.");
177RequestQueue::PreallocatedRequest::PreallocatedRequest() : PreallocatedRequest(nullptr, nullptr) {}
179RequestQueue::PreallocatedRequest::PreallocatedRequest(ReleaseCallback releaseCallback,
180 void* releaseContext)
181 : m_Request(0, true, 0, 0, 0, 0, 0, 0, 0, 0, this),
184 m_ReleaseCallback(releaseCallback),
185 m_ReleaseContext(releaseContext) {}
187RequestQueue::PreallocatedRequest::~PreallocatedRequest() {
188 if (!isAvailable()) {
189 FATAL(
"Destroying a published RequestQueue preallocated token.");
193bool RequestQueue::PreallocatedRequest::isAvailable()
const {
194 return static_cast<size_t>(
m_State) == Idle && !m_ReleaseDepth;
200 m_State(static_cast<size_t>(LifecycleState::Stopped)),
211 m_pOverrunTimer(nullptr),
213#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
214 m_AfterPreallocatedAdmissionHook(nullptr),
215 m_AfterPreallocatedAdmissionContext(nullptr),
216 m_AfterIntakeExchangeHook(nullptr),
217 m_AfterIntakeExchangeContext(nullptr),
218 m_WorkerTransientRetries(0),
219 m_GuardedTransientRetries(0),
220 m_PublisherDrainRetries(0),
222#if PEDIGREE_CONCURRENCY_SMOKE_TESTS
223 m_AfterPreallocatedClaimHook(nullptr),
224 m_AfterPreallocatedClaimContext(nullptr),
231 m_Name(name.cstr(), name.length()) {
232 for (
size_t i = 0; i < REQUEST_QUEUE_NUM_PRIORITIES; ++i) {
234 m_pRequestQueueTail[i] =
nullptr;
238 m_OverrunChecker.queue =
this;
242RequestQueue::~RequestQueue() {
247 const LifecycleState state =
static_cast<LifecycleState
>(
m_State.value());
248 active = (state != LifecycleState::Stopped && state != LifecycleState::Destroyed) ||
253 FATAL(
"RequestQueue '" << m_Name
254 <<
"' reached its base destructor while active; "
255 "the most-derived destructor must call "
263 FATAL(
"Initialising a permanently destroyed RequestQueue");
271 const LifecycleState state =
static_cast<LifecycleState
>(
m_State.value());
272 if (state == LifecycleState::Accepting) {
278 if (state == LifecycleState::Destroyed) {
281 if (state == LifecycleState::Stopping || m_pThread ||
m_bWorkerReady.value()) {
282 ERROR(
"RequestQueue '" << m_Name <<
"' cannot start while stopping");
289 const bool explicitPlacement = workerPlacement(placement);
291 true,
true, explicitPlacement ? &placement :
nullptr);
292 worker->setName(
"RequestQueue worker");
296 assert(
static_cast<LifecycleState
>(
m_State.value()) == LifecycleState::Stopped);
300 m_OverrunChecker.resetBaselineLocked();
308 FATAL(
"RequestQueue '" << m_Name <<
"' could not start its worker");
319 if (
static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Stopped ||
321 FATAL(
"RequestQueue '" << m_Name <<
"' lost its worker during startup");
323 const WaitQueue::WakeReason reason = guard.waitForCompletion(
333 assert(
static_cast<LifecycleState
>(
m_State.value()) == LifecycleState::Stopped);
343 m_State =
static_cast<size_t>(LifecycleState::Accepting);
357bool RequestQueue::stopWorker() {
359 bool hadWorker =
false;
364 closePreallocatedAdmission();
365 if (
static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Destroyed) {
366 m_State =
static_cast<size_t>(LifecycleState::Stopped);
372 ERROR(
"RequestQueue '" << m_Name <<
"' worker cannot halt itself");
376 if (
static_cast<LifecycleState
>(
m_State.value()) == LifecycleState::Accepting) {
379 closePreallocatedAdmission();
380 m_State =
static_cast<size_t>(LifecycleState::Stopping);
389 waitForPreallocatedPublishers();
398 waitForPreallocatedPublishers();
400 if (!
worker->joinForCompletion()) {
401 ERROR(
"RequestQueue '" << m_Name <<
"' could not join its worker");
413 m_State =
static_cast<size_t>(LifecycleState::Stopped);
420void RequestQueue::closePreallocatedAdmission() {
424void RequestQueue::waitForPreallocatedPublishers() {
426#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
427 m_PublisherDrainRetries += 1;
437 FATAL(
"RequestQueue '" << m_Name <<
"' destroy re-entered a queue callback");
440 FATAL(
"RequestQueue '" << m_Name <<
"' destroy re-entered a lifecycle callback");
445 FATAL(
"RequestQueue '" << m_Name
446 <<
"' could not satisfy destroy's worker "
450 if (m_pOverrunTimer) {
451 if (!m_pOverrunTimer->unregisterHandler(&m_OverrunChecker)) {
452 FATAL(
"RequestQueue '" << m_Name <<
"' could not drain its timer callback");
454 m_pOverrunTimer =
nullptr;
461 for (
size_t priority = 0; priority < REQUEST_QUEUE_NUM_PRIORITIES; ++priority) {
463 FATAL(
"RequestQueue '" << m_Name
464 <<
"' retained an incomplete MPSC "
465 "publication after producer drain");
470 m_pRequestQueueTail[priority] =
nullptr;
473 Request* next = request->m_Next;
474 request->m_Next = cancelled;
480 m_nAsyncRequests = 0;
485 Request* next = cancelled->m_Next;
486 cancelled->m_Next =
nullptr;
489 releaseRequest(cancelled);
495 closePreallocatedAdmission();
496 m_State =
static_cast<size_t>(LifecycleState::Destroyed);
502 uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8) {
503 return addRequest(priority, RequestQueue::Block, p1, p2, p3, p4, p5, p6, p7, p8);
507 uint64_t p2, uint64_t p3, uint64_t p4, uint64_t p5, uint64_t p6,
508 uint64_t p7, uint64_t p8) {
514 Request* candidate =
new Request(priority,
false, p1, p2, p3, p4, p5, p6, p7, p8);
516 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
517 ERROR(
"RequestQueue '" << m_Name <<
"' rejected invalid priority " << priority);
523 bool rejected =
false;
524 bool executeInline =
false;
529 if (
static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Accepting) {
534 executeInline =
true;
536#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
537 m_GuardedTransientRetries += 1;
543 if (action != NewRequest) {
548 if (action == ReturnImmediately) {
579 return waitForRequest(request);
581 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
582 ERROR(
"RequestQueue '" << m_Name <<
"' rejected invalid priority " << priority);
583 Request* candidate =
new Request(priority,
false, p1, p2, p3, p4, p5, p6, p7, p8);
592 uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7,
594 return addAsyncRequestInternal(priority, p1, p2, p3, p4, p5, p6, p7, p8);
598 uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7,
601 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES ||
605 static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Accepting) {
609 Request* request =
new Request(priority,
true, p1, p2, p3, p4, p5, p6, p7, p8);
615 if (admission & PublicationClosed) {
621 m_nAsyncRequests += 1;
641 PreallocatedRequest& token,
size_t priority, uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4,
642 uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8) {
643 return publishPreallocatedRequest(token, PreallocatedRequest::Idle, priority, p1, p2, p3, p4, p5,
648 PreallocatedRequest& token,
size_t priority, uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4,
649 uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8) {
650 return publishPreallocatedRequest(token, PreallocatedRequest::Releasing, priority, p1, p2, p3, p4,
654RequestQueue::PreallocatedPublishResult RequestQueue::publishPreallocatedRequest(
655 PreallocatedRequest& token, PreallocatedRequest::State availableState,
size_t priority,
656 uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7,
658 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
659 return PreallocatedPublishResult::InvalidPriority;
663 (availableState != PreallocatedRequest::Releasing ||
664 !RequestQueueCallbackScope::allowsReleaseHandoff(
this, &token))) {
665 return PreallocatedPublishResult::QueueStopped;
675#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
676 if (m_AfterPreallocatedAdmissionHook) {
677 m_AfterPreallocatedAdmissionHook(m_AfterPreallocatedAdmissionContext);
680 if (admission & PublicationClosed) {
682 return PreallocatedPublishResult::QueueStopped;
685 if (!token.m_State.compareAndSwap(availableState, PreallocatedRequest::Claimed)) {
687 return PreallocatedPublishResult::TokenBusy;
689#if PEDIGREE_CONCURRENCY_SMOKE_TESTS
690 if (m_AfterPreallocatedClaimHook) {
691 m_AfterPreallocatedClaimHook(m_AfterPreallocatedClaimContext);
695 if (!token.m_State.compareAndSwap(availableState, PreallocatedRequest::Claimed)) {
696 return PreallocatedPublishResult::TokenBusy;
700 Request* request = &token.m_Request;
709 request->m_ReturnValue = 0;
711 request->m_References = 1;
713 request->m_Next =
nullptr;
714 request->m_Intake.next =
nullptr;
715 request->m_Priority = priority;
716 request->m_Asynchronous =
true;
717 request->m_Rejected =
false;
718 request->m_Completed =
false;
721 token.m_State = PreallocatedRequest::Published;
724 return PreallocatedPublishResult::Accepted;
726 token.m_State = PreallocatedRequest::Published;
728 m_nAsyncRequests += 1;
732 return PreallocatedPublishResult::Accepted;
736uint64_t RequestQueue::addAsyncRequestInternal(
size_t priority, uint64_t p1, uint64_t p2,
737 uint64_t p3, uint64_t p4, uint64_t p5, uint64_t p6,
738 uint64_t p7, uint64_t p8) {
740 Request* request =
new Request(priority,
true, p1, p2, p3, p4, p5, p6, p7, p8);
741 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
742 ERROR(
"RequestQueue '" << m_Name <<
"' rejected invalid priority " << priority);
754 Request* request =
new Request(priority,
true, p1, p2, p3, p4, p5, p6, p7, p8);
756 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
757 ERROR(
"RequestQueue '" << m_Name <<
"' rejected invalid priority " << priority);
762 bool rejected =
false;
763 bool overloaded =
false;
768 if (
static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Accepting) {
771#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
772 m_GuardedTransientRetries += 1;
781 m_nAsyncRequests += 1;
794 ERROR(
"RequestQueue: '" << m_Name <<
"' is not keeping up with async requests");
795 ERROR(
" -> priority=" << priority <<
", p1=" <<
Hex << p1 <<
", p2=" << p2 <<
", p3=" << p3
797 ERROR(
" -> p5=" <<
Hex << p5 <<
", p6=" << p6 <<
", p7=" << p7 <<
", p8=" << p8);
831 if (!m_pOverrunTimer) {
833 if (timer && timer->registerHandler(&m_OverrunChecker)) {
834 m_pOverrunTimer = timer;
836 !timer->
armHandler(&m_OverrunChecker, Time::getTicks() + Time::Multiplier::Second)) {
837 FATAL(
"RequestQueue could not arm its overrun checker");
848 return static_cast<LifecycleState
>(
m_State.value());
852 return RequestQueueCallbackScope::contains(
this);
863 return m_pThread != current;
876 if (request.isAvailable())
883 return request.isAvailable();
899 if (
static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Accepting) {
900 ERROR(
"RequestQueue '" << m_Name <<
"' cannot drain while it is stopping");
903 if (m_pThread == current) {
904 ERROR(
"RequestQueue '" << m_Name <<
"' worker cannot drain itself");
908 const WaitQueue::WakeReason reason = guard.waitForCompletion(
909 WaitQueue::Channel(
this, 1), Thread::CallbackDrain,
reinterpret_cast<uintptr_t
>(
this));
919 return queue->
work();
923#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
926 const TestAccess::Publication publication =
927 TestAccess::beginPush(lane.m_Queue, request->m_Intake);
928 if (m_AfterIntakeExchangeHook) {
929 m_AfterIntakeExchangeHook(m_AfterIntakeExchangeContext);
931 TestAccess::finishPush(lane.m_Queue, publication);
933 m_IntakeLanes[request->m_Priority].m_Queue.push(request->m_Intake);
945 assert(priority < REQUEST_QUEUE_NUM_PRIORITIES);
946 using PopResult = IntrusiveMpscQueue<IntakeNode, &IntakeNode::next>::PopResult;
950 const PopResult result =
m_IntakeLanes[priority].m_Queue.pop(node);
951 if (result == PopResult::Empty) {
954 if (result == PopResult::Transient) {
958 if (!node || !node->owner)
960 Request* request = node->owner;
961 assert(request->m_Priority == priority);
962 assert(!request->m_Next);
963 if (m_pRequestQueueTail[priority]) {
964 m_pRequestQueueTail[priority]->m_Next = request;
968 m_pRequestQueueTail[priority] = request;
974 using PopResult = IntrusiveMpscQueue<IntakeNode, &IntakeNode::next>::PopResult;
976 for (
size_t priority = 0; priority < REQUEST_QUEUE_NUM_PRIORITIES; ++priority) {
981 m_pRequestQueueTail[priority] =
nullptr;
983 request->m_Next =
nullptr;
985 return NextRequestResult::Item;
989 const PopResult result =
m_IntakeLanes[priority].m_Queue.pop(node);
990 if (result == PopResult::Transient) {
991#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
992 m_WorkerTransientRetries += 1;
996 return NextRequestResult::Retry;
998 if (result == PopResult::Item) {
999 if (!node || !node->owner)
1000 return NextRequestResult::Retry;
1001 request = node->owner;
1002 assert(request->m_Priority == priority);
1004 return NextRequestResult::Item;
1008 return NextRequestResult::Empty;
1020 if (!queued->m_pPreallocatedOwner &&
compareRequests(*queued, request)) {
1023 queued = queued->m_Next;
1031 auto guard = request->m_Completion.acquire();
1032 assert(!request->m_Completed);
1033 request->m_ReturnValue = returnValue;
1034 request->m_Rejected = rejected;
1035 request->m_Completed =
true;
1038 request->m_ReturnValue = returnValue;
1039 request->m_Rejected = rejected;
1040 request->m_Completed =
true;
1058 assert(&owner->m_Request == request);
1059 assert(
static_cast<size_t>(owner->m_State) == PreallocatedRequest::Published);
1060 owner->m_ReleaseDepth += 1;
1061 owner->m_State = PreallocatedRequest::Releasing;
1062 if (owner->m_ReleaseCallback) {
1064 owner->m_ReleaseCallback(owner->m_ReleaseContext);
1068 const size_t state = owner->m_State;
1069 if (state == PreallocatedRequest::Releasing) {
1070 if (owner->m_State.compareAndSwap(PreallocatedRequest::Releasing,
1071 PreallocatedRequest::Idle)) {
1076 if (state == PreallocatedRequest::Claimed) {
1085 assert(state == PreallocatedRequest::Published || state == PreallocatedRequest::Idle);
1093 owner->m_ReleaseDepth -= 1;
1101 request->m_References += 1;
1104void RequestQueue::releaseRequest(Request* request) {
1105 if (request->m_pPreallocatedOwner) {
1110 assert(
static_cast<size_t>(request->m_References));
1111 if ((request->m_References -= 1) == 0) {
1116uint64_t RequestQueue::waitForRequest(Request* request) {
1117 uint64_t result = 0;
1120 auto guard = request->m_Completion.acquire();
1121 if (request->m_Completed) {
1122 if (!request->m_Rejected) {
1123 result = request->m_ReturnValue;
1131 WaitQueue::WakeReason reason = guard.waitForCompletion(
WaitQueue::Channel(), Thread::CondWait,
1132 reinterpret_cast<uintptr_t
>(request));
1136 releaseRequest(request);
1150 static_cast<LifecycleState
>(
m_State.value()) != LifecycleState::Stopped ||
1152 FATAL(
"RequestQueue '" << m_Name <<
"' worker entered with invalid state");
1160 NextRequestResult next = NextRequestResult::Empty;
1161 bool accepting =
false;
1165 const LifecycleState state =
static_cast<LifecycleState
>(
m_State.value());
1166 if (state == LifecycleState::Stopping) {
1171 if (state == LifecycleState::Stopped) {
1175 }
else if (state == LifecycleState::Accepting) {
1178 if (next == NextRequestResult::Item) {
1190 if (next != NextRequestResult::Item) {
1194 if (!accepting || next == NextRequestResult::Retry) {
1199 auto waitGuard = m_WorkerWaiters.acquire();
1201 static_cast<LifecycleState
>(
m_State.value()) == LifecycleState::Accepting) {
1203 reinterpret_cast<uintptr_t
>(
this));
1211 uint64_t result =
executeRequest(request->p1, request->p2, request->p3, request->p4,
1212 request->p5, request->p6, request->p7, request->p8);
1220 if (request->m_Asynchronous) {
1221 assert(m_nAsyncRequests.value());
1222 m_nAsyncRequests -= 1;
1230 releaseRequest(request);
1251void RequestQueue::RequestQueueOverrunChecker::resetBaselineLocked() {
1252 m_LastQueueSize = 0;
1254 m_HasBacklogBaseline =
false;
1257RequestQueue::OverrunStatus RequestQueue::RequestQueueOverrunChecker::sample(
size_t& lastSize,
1258 size_t& currentSize) {
1259 auto guard = queue->m_RequestQueueWaiters.acquire();
1260 lastSize = m_LastQueueSize;
1261 const size_t total = queue->m_nTotalRequests.value();
1262 const size_t active = queue->m_nActiveRequests.value();
1263 assert(total >= active);
1264 currentSize = total - active;
1266 const size_t progress = queue->m_WorkerProgressGeneration;
1267 if (
static_cast<LifecycleState
>(queue->m_State.value()) != LifecycleState::Accepting ||
1269 resetBaselineLocked();
1270 return OverrunStatus::Clear;
1273 OverrunStatus status = OverrunStatus::Armed;
1274 if (m_HasBacklogBaseline) {
1275 if (progress == m_LastProgressGeneration && currentSize >= m_LastQueueSize) {
1276 status = OverrunStatus::Stalled;
1277 }
else if (progress != m_LastProgressGeneration && currentSize > m_LastQueueSize) {
1278 status = OverrunStatus::Overloaded;
1282 m_LastQueueSize = currentSize;
1283 m_LastProgressGeneration = progress;
1284 m_HasBacklogBaseline =
true;
1290 Timer* source = queue->m_pOverrunTimer;
1291 if (m_Tick < Time::Multiplier::Second) {
1293 !source->
armHandler(
this, Time::getTicks() + Time::Multiplier::Second - m_Tick)) {
1294 FATAL(
"RequestQueue could not rearm its overrun checker");
1298 m_Tick %= Time::Multiplier::Second;
1300 !source->
armHandler(
this, Time::getTicks() + Time::Multiplier::Second - m_Tick)) {
1301 FATAL(
"RequestQueue could not rearm its overrun checker");
1304 size_t lastSize = 0;
1305 size_t currentSize = 0;
1306 const OverrunStatus status = sample(lastSize, currentSize);
1307 if (status == OverrunStatus::Stalled) {
1311 WARNING(
"RequestQueue '" << queue->m_Name
1312 <<
"' completed no request during the watchdog interval with "
1313 << currentSize <<
" queued requests.");
1314 }
else if (status == OverrunStatus::Overloaded) {
1315 WARNING(
"RequestQueue '" << queue->m_Name <<
"' backlog grew from " << lastSize <<
" to "
1316 << currentSize <<
" despite worker progress.");
1321#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
1322RequestQueue::OverrunStatus RequestQueue::sampleOverrunForTest() {
1323 size_t lastSize = 0;
1324 size_t currentSize = 0;
1325 return m_OverrunChecker.sample(lastSize, currentSize);
virtual Timer * getTimer()=0
bool isOwnedByCurrentThread() const
void ringIrqWorkDoorbell()
void unregisterWorkerWake(SchedulerWorkerWake &worker)
static bool getInterrupts()
static ProcessorInformation & information()
static ExecutionContext executionContext()
static void setInterrupts(bool bEnable)
virtual void timer(uint64_t delta)
Atomic< PerProcessorScheduler * > m_pWorkerScheduler
void publishRequest(Request *request)
Request * m_pRequestQueue[REQUEST_QUEUE_NUM_PRIORITIES]
static void retainRequest(Request *request)
virtual bool compareRequests(const Request &a, const Request &b)
static void completeRequest(Request *request, uint64_t returnValue, bool rejected)
MUST_USE_RESULT bool waitForPreallocated(PreallocatedRequest &request)
size_t m_WorkerProgressGeneration
WaitQueue m_RequestQueueWaiters
Atomic< size_t > m_bWorkerActive
bool canWaitForCompletion()
MUST_USE_RESULT PreallocatedPublishResult republishPreallocatedWhileReleasing(PreallocatedRequest &request, size_t priority, uint64_t p1=0, uint64_t p2=0, uint64_t p3=0, uint64_t p4=0, uint64_t p5=0, uint64_t p6=0, uint64_t p7=0, uint64_t p8=0)
LifecycleState getLifecycleState()
void releasePreallocatedRequest(Request *request)
RequestQueue(const String &name)
IntakeLane m_IntakeLanes[REQUEST_QUEUE_NUM_PRIORITIES]
MUST_USE_RESULT bool publishAsyncRequest(size_t priority, uint64_t p1=0, uint64_t p2=0, uint64_t p3=0, uint64_t p4=0, uint64_t p5=0, uint64_t p6=0, uint64_t p7=0, uint64_t p8=0)
Request * findDuplicate(const Request &request)
Atomic< size_t > m_bWorkerReady
size_t m_nMaxAsyncRequests
Atomic< size_t > m_nTotalRequests
MUST_USE_RESULT bool halt()
virtual void initialise()
MUST_USE_RESULT bool resume()
void invokeCancelRequest(const Request &request)
MUST_USE_RESULT PreallocatedPublishResult publishPreallocated(PreallocatedRequest &request, size_t priority, uint64_t p1=0, uint64_t p2=0, uint64_t p3=0, uint64_t p4=0, uint64_t p5=0, uint64_t p6=0, uint64_t p7=0, uint64_t p8=0)
virtual void cancelRequest(const Request &request)
virtual uint64_t executeRequest(uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8)=0
Atomic< size_t > m_nActiveRequests
MUST_USE_RESULT uint64_t addRequest(size_t priority, uint64_t p1=0, uint64_t p2=0, uint64_t p3=0, uint64_t p4=0, uint64_t p5=0, uint64_t p6=0, uint64_t p7=0, uint64_t p8=0)
Request * m_pActiveRequest
static int trampoline(void *p)
NextRequestResult getNextRequest(Request *&request)
Atomic< size_t > m_PublicationState
uint64_t addAsyncRequest(size_t priority, uint64_t p1=0, uint64_t p2=0, uint64_t p3=0, uint64_t p4=0, uint64_t p5=0, uint64_t p6=0, uint64_t p7=0, uint64_t p8=0)
bool drainIntakeLocked(size_t priority)
void discardRequest(Request *request)
bool callbackActiveOnCurrentThread() const
static Scheduler & instance()
virtual bool supportsDeadlines() const
virtual bool armHandler(TimerHandler *, uint64_t)
MUST_USE_RESULT WakeReason wait(const Channel &channel=Channel(), size_t debugState=0, uintptr_t debugAddress=0, StackDiscardCleanup onStackDiscard=nullptr, void *stackDiscardContext=nullptr)
MUST_USE_RESULT WakeReason waitForCompletion(const Channel &channel=Channel(), size_t debugState=0, uintptr_t debugAddress=0)
RequestQueueCallbackScope * m_pRequestQueueCallback