The Pedigree Project 0.1
RequestQueue.cc
1/*
2 * Copyright (c) 2008-2014, Pedigree Developers
3 *
4 * Please see the CONTRIB file in the root of the source tree for a full
5 * list of contributors.
6 *
7 * Permission to use, copy, modify, and distribute this software for any
8 * purpose with or without fee is hereby granted, provided that the above
9 * copyright notice and this permission notice appear in all copies.
10 *
11 * THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
12 * WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
13 * MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
14 * ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
15 * WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
16 * ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
17 * OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
18 */
19
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"
34
35class Process;
36
37static_assert(__atomic_always_lock_free(sizeof(size_t), nullptr),
38 "RequestQueue preallocated-publication words must be lock-free");
39static_assert(__atomic_always_lock_free(sizeof(PerProcessorScheduler*), nullptr),
40 "RequestQueue preallocated-publication pointers must be lock-free");
41
50#if defined(PEDIGREE_BUILDUTILS)
52 public:
54
55 static bool contains(const RequestQueue*) {
56 return false;
57 }
58
59 static bool allowsReleaseHandoff(const RequestQueue*, const RequestQueue::PreallocatedRequest*) {
60 return false;
61 }
62};
63#else
65 public:
67 RequestQueue::PreallocatedRequest* releasingToken = nullptr)
68 : m_Queue(queue),
69 m_ReleasingToken(releasingToken),
70 m_Thread(nullptr),
71 m_StateLevel(0),
72 m_Previous(nullptr),
73 m_Cleanup(),
74 m_Active(false) {
75 m_Thread = Processor::information().getCurrentThread();
76 if (!m_Thread) {
77 return;
78 }
79
80 const bool interruptsWereEnabled = Processor::getInterrupts();
82 m_StateLevel = __atomic_load_n(&m_Thread->m_nStateLevel, __ATOMIC_ACQUIRE);
83 m_Previous = m_Thread->m_StateLevels[m_StateLevel].m_pRequestQueueCallback;
84
85 // Publish cleanup before this callback becomes visible to nested work.
86 m_Thread->armAtomicStateCleanup(m_Cleanup, abandon, this);
87 m_Active = true;
88 m_Thread->m_StateLevels[m_StateLevel].m_pRequestQueueCallback = this;
89 Processor::setInterrupts(interruptsWereEnabled);
90 }
91
93 if (!m_Active) {
94 return;
95 }
96
97 const bool interruptsWereEnabled = Processor::getInterrupts();
99 restore();
100 m_Thread->disarmAtomicStateCleanup(m_Cleanup);
101 Processor::setInterrupts(interruptsWereEnabled);
102 }
103
104 static bool contains(const RequestQueue* queue) {
105 Thread* thread = Processor::information().getCurrentThread();
106 if (!thread || !queue) {
107 return false;
108 }
109
111 thread->m_StateLevels[thread->m_nStateLevel].m_pRequestQueueCallback;
112 while (scope) {
113 if (scope->m_Queue == queue) {
114 return true;
115 }
116 scope = scope->m_Previous;
117 }
118 return false;
119 }
120
121 static bool allowsReleaseHandoff(const RequestQueue* queue,
123 Thread* thread = Processor::information().getCurrentThread();
124 if (!thread || !queue || !token) {
125 return false;
126 }
127
129 thread->m_StateLevels[thread->m_nStateLevel].m_pRequestQueueCallback;
130 while (scope) {
131 if (scope->m_Queue == queue && scope->m_ReleasingToken == token) {
132 return true;
133 }
134 scope = scope->m_Previous;
135 }
136 return false;
137 }
138
139 private:
141 RequestQueueCallbackScope& operator=(const RequestQueueCallbackScope&) = delete;
142
143 static void abandon(void* context) {
144 RequestQueueCallbackScope* scope = reinterpret_cast<RequestQueueCallbackScope*>(context);
145 if (scope) {
146 scope->restore();
147 }
148 }
149
150 void restore() {
151 if (!m_Active) {
152 return;
153 }
154 if (m_StateLevel >= MAX_NESTED_EVENTS ||
155 m_Thread->m_StateLevels[m_StateLevel].m_pRequestQueueCallback != this) {
156 FATAL_NOLOCK("RequestQueue callback scope stack was corrupted.");
157 return;
158 }
159
160 const bool interruptsWereEnabled = Processor::getInterrupts();
162 m_Thread->m_StateLevels[m_StateLevel].m_pRequestQueueCallback = m_Previous;
163 m_Active = false;
164 Processor::setInterrupts(interruptsWereEnabled);
165 }
166
167 RequestQueue* m_Queue;
168 RequestQueue::PreallocatedRequest* m_ReleasingToken;
169 Thread* m_Thread;
170 size_t m_StateLevel;
171 RequestQueueCallbackScope* m_Previous;
172 AtomicStateCleanupRecord m_Cleanup;
173 bool m_Active;
174};
175#endif
176
177RequestQueue::PreallocatedRequest::PreallocatedRequest() : PreallocatedRequest(nullptr, nullptr) {}
178
179RequestQueue::PreallocatedRequest::PreallocatedRequest(ReleaseCallback releaseCallback,
180 void* releaseContext)
181 : m_Request(0, true, 0, 0, 0, 0, 0, 0, 0, 0, this),
182 m_State(Idle),
183 m_ReleaseDepth(0),
184 m_ReleaseCallback(releaseCallback),
185 m_ReleaseContext(releaseContext) {}
186
187RequestQueue::PreallocatedRequest::~PreallocatedRequest() {
188 if (!isAvailable()) {
189 FATAL("Destroying a published RequestQueue preallocated token.");
190 }
191}
192
193bool RequestQueue::PreallocatedRequest::isAvailable() const {
194 return static_cast<size_t>(m_State) == Idle && !m_ReleaseDepth;
195}
196
198 : m_IntakeLanes(),
199 m_pActiveRequest(nullptr),
200 m_State(static_cast<size_t>(LifecycleState::Stopped)),
201#if THREADS
204 m_WorkerWaiters(),
205 m_pThread(nullptr),
206 m_pWorkerScheduler(nullptr),
207 m_WorkerWake(),
211 m_pOverrunTimer(nullptr),
212 m_PublicationState(PublicationClosed),
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),
221#endif
222#if PEDIGREE_CONCURRENCY_SMOKE_TESTS
223 m_AfterPreallocatedClaimHook(nullptr),
224 m_AfterPreallocatedClaimContext(nullptr),
225#endif
226#endif
228 m_nAsyncRequests(0),
231 m_Name(name.cstr(), name.length()) {
232 for (size_t i = 0; i < REQUEST_QUEUE_NUM_PRIORITIES; ++i) {
233 m_pRequestQueue[i] = nullptr;
234 m_pRequestQueueTail[i] = nullptr;
235 }
236
237#if THREADS
238 m_OverrunChecker.queue = this;
239#endif
240}
241
242RequestQueue::~RequestQueue() {
243#if THREADS
244 bool active = false;
245 {
246 auto guard = m_RequestQueueWaiters.acquire();
247 const LifecycleState state = static_cast<LifecycleState>(m_State.value());
248 active = (state != LifecycleState::Stopped && state != LifecycleState::Destroyed) ||
249 m_pThread || m_pOverrunTimer || m_nTotalRequests.value() ||
250 m_PublicationState.value() != PublicationClosed;
251 }
252 if (active) {
253 FATAL("RequestQueue '" << m_Name
254 << "' reached its base destructor while active; "
255 "the most-derived destructor must call "
256 "destroy().");
257 }
258#endif
259}
260
262 if (!resume()) {
263 FATAL("Initialising a permanently destroyed RequestQueue");
264 }
265}
266
267#if THREADS
269 {
270 auto guard = m_RequestQueueWaiters.acquire();
271 const LifecycleState state = static_cast<LifecycleState>(m_State.value());
272 if (state == LifecycleState::Accepting) {
273 assert(m_pThread && m_bWorkerReady.value());
274 assert(!(m_PublicationState.value() & PublicationClosed));
275 return true;
276 }
277
278 if (state == LifecycleState::Destroyed) {
279 return false;
280 }
281 if (state == LifecycleState::Stopping || m_pThread || m_bWorkerReady.value()) {
282 ERROR("RequestQueue '" << m_Name << "' cannot start while stopping");
283 return false;
284 }
285 }
286
287 Process* process = Scheduler::instance().getKernelProcess();
288 ThreadPlacement placement;
289 const bool explicitPlacement = workerPlacement(placement);
290 Thread* worker = new Thread(process, &trampoline, reinterpret_cast<void*>(this), nullptr, false,
291 true, true, explicitPlacement ? &placement : nullptr);
292 worker->setName("RequestQueue worker");
293
294 {
295 auto guard = m_RequestQueueWaiters.acquire();
296 assert(static_cast<LifecycleState>(m_State.value()) == LifecycleState::Stopped);
297 assert(!m_pThread);
298 assert(!m_bWorkerReady.value());
299 assert(m_PublicationState.value() & PublicationClosed);
300 m_OverrunChecker.resetBaselineLocked();
301 m_pThread = worker;
302 m_pWorkerScheduler = worker->getScheduler();
303 m_pWorkerScheduler.value()->registerWorkerWake(m_WorkerWake, m_WorkerWaiters);
304 }
305
306 // The delayed worker cannot observe partially published queue state.
307 if (!worker->start()) {
308 FATAL("RequestQueue '" << m_Name << "' could not start its worker");
309 }
310
311 // A terminal request can retire a delayed Thread before its entry point
312 // runs. Do not publish a usable queue until work() has installed the
313 // queue-owned lifetime deferral.
314 while (true) {
315 auto guard = m_RequestQueueWaiters.acquire();
316 if (m_bWorkerReady.value()) {
317 break;
318 }
319 if (static_cast<LifecycleState>(m_State.value()) != LifecycleState::Stopped ||
320 m_pThread != worker) {
321 FATAL("RequestQueue '" << m_Name << "' lost its worker during startup");
322 }
323 const WaitQueue::WakeReason reason = guard.waitForCompletion(
324 WaitQueue::Channel(this, 2), Thread::CondWait, reinterpret_cast<uintptr_t>(this));
325 (void)reason;
326 }
327
328 while (true) {
329 bool opened = false;
330 {
331 auto guard = m_RequestQueueWaiters.acquire();
332 assert(m_pThread == worker && m_bWorkerReady.value());
333 assert(static_cast<LifecycleState>(m_State.value()) == LifecycleState::Stopped);
334
335 // A rejected hard producer briefly contributes to the closed
336 // gate's low-bit count. Reopen only the exact closed-and-drained
337 // state so its eventual decrement can never underflow a new
338 // publication lifetime.
339 opened = m_PublicationState.compareAndSwap(PublicationClosed, 0);
340 if (opened) {
341 // This is the final acceptance point. The worker, scheduler, and
342 // lock-free work publication are already established.
343 m_State = static_cast<size_t>(LifecycleState::Accepting);
344 }
345 }
346
347 if (opened) {
348 break;
349 }
351 }
352
353 m_pWorkerScheduler.value()->ringIrqWorkDoorbell(m_WorkerWake);
354 return true;
355}
356
357bool RequestQueue::stopWorker() {
358 Thread* worker = nullptr;
359 bool hadWorker = false;
360 {
361 auto guard = m_RequestQueueWaiters.acquire();
362 worker = m_pThread;
363 if (!worker) {
364 closePreallocatedAdmission();
365 if (static_cast<LifecycleState>(m_State.value()) != LifecycleState::Destroyed) {
366 m_State = static_cast<size_t>(LifecycleState::Stopped);
367 }
368 m_bWorkerReady = 0;
369 m_bWorkerActive = 0;
370 m_pWorkerScheduler = nullptr;
371 } else if (worker == Processor::information().getCurrentThread()) {
372 ERROR("RequestQueue '" << m_Name << "' worker cannot halt itself");
373 return false;
374 } else {
375 hadWorker = true;
376 if (static_cast<LifecycleState>(m_State.value()) == LifecycleState::Accepting) {
377 // Close before publishing Stopping, so every producer that can
378 // still observe Accepting is either rejected or counted below.
379 closePreallocatedAdmission();
380 m_State = static_cast<size_t>(LifecycleState::Stopping);
381 }
382 }
383 }
384
385 if (!hadWorker) {
386 // Even a closed queue admits rejected publishers into the low-bit
387 // lifetime count long enough to observe the closed bit. Do not let a
388 // never-started or already-stopped queue outrun one of those callers.
389 waitForPreallocatedPublishers();
390 return true;
391 }
392
393 PerProcessorScheduler* scheduler = m_pWorkerScheduler.value();
394 if (scheduler) {
395 scheduler->ringIrqWorkDoorbell(m_WorkerWake);
396 }
397
398 waitForPreallocatedPublishers();
399
400 if (!worker->joinForCompletion()) {
401 ERROR("RequestQueue '" << m_Name << "' could not join its worker");
402 return false;
403 }
404
405 if (scheduler) {
406 scheduler->unregisterWorkerWake(m_WorkerWake);
407 }
408
409 {
410 auto guard = m_RequestQueueWaiters.acquire();
411 m_pThread = nullptr;
412 m_pWorkerScheduler = nullptr;
413 m_State = static_cast<size_t>(LifecycleState::Stopped);
414 m_bWorkerReady = 0;
415 m_bWorkerActive = 0;
416 }
417 return true;
418}
419
420void RequestQueue::closePreallocatedAdmission() {
421 m_PublicationState |= PublicationClosed;
422}
423
424void RequestQueue::waitForPreallocatedPublishers() {
425 while (m_PublicationState.value() & PublicationCountMask) {
426#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
427 m_PublisherDrainRetries += 1;
428#endif
430 }
431}
432#endif
433
435#if THREADS
437 FATAL("RequestQueue '" << m_Name << "' destroy re-entered a queue callback");
438 }
440 FATAL("RequestQueue '" << m_Name << "' destroy re-entered a lifecycle callback");
441 }
442 TerminationDeferral terminationDeferral;
443 LockGuard<Mutex> lifecycleGuard(m_LifecycleMutex);
444 if (!stopWorker()) {
445 FATAL("RequestQueue '" << m_Name
446 << "' could not satisfy destroy's worker "
447 "drain contract");
448 }
449
450 if (m_pOverrunTimer) {
451 if (!m_pOverrunTimer->unregisterHandler(&m_OverrunChecker)) {
452 FATAL("RequestQueue '" << m_Name << "' could not drain its timer callback");
453 }
454 m_pOverrunTimer = nullptr;
455 }
456
457 Request* cancelled = nullptr;
458 {
459 auto guard = m_RequestQueueWaiters.acquire();
460 assert(!m_pActiveRequest);
461 for (size_t priority = 0; priority < REQUEST_QUEUE_NUM_PRIORITIES; ++priority) {
462 if (!drainIntakeLocked(priority)) {
463 FATAL("RequestQueue '" << m_Name
464 << "' retained an incomplete MPSC "
465 "publication after producer drain");
466 }
467
468 Request* request = m_pRequestQueue[priority];
469 m_pRequestQueue[priority] = nullptr;
470 m_pRequestQueueTail[priority] = nullptr;
471
472 while (request) {
473 Request* next = request->m_Next;
474 request->m_Next = cancelled;
475 cancelled = request;
476 request = next;
477 }
478 }
480 m_nAsyncRequests = 0;
482 }
483
484 while (cancelled) {
485 Request* next = cancelled->m_Next;
486 cancelled->m_Next = nullptr;
487 invokeCancelRequest(*cancelled);
488 completeRequest(cancelled, 0, true);
489 releaseRequest(cancelled);
490 cancelled = next;
491 }
492
493 {
494 auto guard = m_RequestQueueWaiters.acquire();
495 closePreallocatedAdmission();
496 m_State = static_cast<size_t>(LifecycleState::Destroyed);
497 }
498#endif
499}
500
501uint64_t RequestQueue::addRequest(size_t priority, uint64_t p1, uint64_t p2, uint64_t p3,
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);
504}
505
506uint64_t RequestQueue::addRequest(size_t priority, ActionOnDuplicate action, uint64_t p1,
507 uint64_t p2, uint64_t p3, uint64_t p4, uint64_t p5, uint64_t p6,
508 uint64_t p7, uint64_t p8) {
509#if THREADS
511 return 0;
512 }
513
514 Request* candidate = new Request(priority, false, p1, p2, p3, p4, p5, p6, p7, p8);
515
516 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
517 ERROR("RequestQueue '" << m_Name << "' rejected invalid priority " << priority);
518 discardRequest(candidate);
519 return 0;
520 }
521
522 Request* request = nullptr;
523 bool rejected = false;
524 bool executeInline = false;
525 while (true) {
526 bool retry = false;
527 {
528 auto guard = m_RequestQueueWaiters.acquire();
529 if (static_cast<LifecycleState>(m_State.value()) != LifecycleState::Accepting) {
530 rejected = true;
531 } else if (m_pThread == Processor::information().getCurrentThread()) {
532 // A worker cannot wait for itself. Execute nested synchronous
533 // work inline after dropping the queue guard.
534 executeInline = true;
535 } else if (action != NewRequest && !drainIntakeLocked(priority)) {
536#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
537 m_GuardedTransientRetries += 1;
538#endif
539 // An allocation-backed predecessor can be hidden behind an
540 // unfinished token link, so do not deduplicate around it.
541 retry = true;
542 } else {
543 if (action != NewRequest) {
544 request = findDuplicate(*candidate);
545 }
546
547 if (request) {
548 if (action == ReturnImmediately) {
549 rejected = true;
550 } else {
551 retainRequest(request);
552 }
553 } else {
554 m_nTotalRequests += 1;
555 publishRequest(candidate);
556 request = candidate;
557 candidate = nullptr;
558 }
559 }
560 }
561
562 if (!retry) {
563 break;
564 }
566 }
567
568 if (executeInline) {
569 delete candidate;
570 return executeRequest(p1, p2, p3, p4, p5, p6, p7, p8);
571 }
572 if (candidate) {
573 discardRequest(candidate);
574 }
575 if (rejected) {
576 return 0;
577 }
578 assert(request);
579 return waitForRequest(request);
580#else
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);
584 discardRequest(candidate);
585 return 0;
586 }
587 return executeRequest(p1, p2, p3, p4, p5, p6, p7, p8);
588#endif
589}
590
591uint64_t RequestQueue::addAsyncRequest(size_t priority, uint64_t p1, uint64_t p2, uint64_t p3,
592 uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7,
593 uint64_t p8) {
594 return addAsyncRequestInternal(priority, p1, p2, p3, p4, p5, p6, p7, p8);
595}
596
597bool RequestQueue::publishAsyncRequest(size_t priority, uint64_t p1, uint64_t p2, uint64_t p3,
598 uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7,
599 uint64_t p8) {
600#if THREADS
601 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES ||
602 Processor::executionContext() != ExecutionContext::WaitableThread ||
605 static_cast<LifecycleState>(m_State.value()) != LifecycleState::Accepting) {
606 return false;
607 }
608
609 Request* request = new Request(priority, true, p1, p2, p3, p4, p5, p6, p7, p8);
610 if (!request) {
611 return false;
612 }
613
614 const size_t admission = (m_PublicationState += 1);
615 if (admission & PublicationClosed) {
617 delete request;
618 return false;
619 }
620
621 m_nAsyncRequests += 1;
622 m_nTotalRequests += 1;
623 publishRequest(request);
625 return true;
626#else
627 (void)priority;
628 (void)p1;
629 (void)p2;
630 (void)p3;
631 (void)p4;
632 (void)p5;
633 (void)p6;
634 (void)p7;
635 (void)p8;
636 return false;
637#endif
638}
639
640RequestQueue::PreallocatedPublishResult RequestQueue::publishPreallocated(
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,
644 p6, p7, p8);
645}
646
647RequestQueue::PreallocatedPublishResult RequestQueue::republishPreallocatedWhileReleasing(
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,
651 p5, p6, p7, p8);
652}
653
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,
657 uint64_t p8) {
658 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
659 return PreallocatedPublishResult::InvalidPriority;
660 }
661
663 (availableState != PreallocatedRequest::Releasing ||
664 !RequestQueueCallbackScope::allowsReleaseHandoff(this, &token))) {
665 return PreallocatedPublishResult::QueueStopped;
666 }
667
668 // The sole callback exception is its own Releasing token: this path only
669 // hands that token back through the lock-free intake, performs no
670 // allocation or lifecycle transition, and the closed admission gate below
671 // still rejects it during shutdown.
672
673#if THREADS
674 const size_t admission = (m_PublicationState += 1);
675#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
676 if (m_AfterPreallocatedAdmissionHook) {
677 m_AfterPreallocatedAdmissionHook(m_AfterPreallocatedAdmissionContext);
678 }
679#endif
680 if (admission & PublicationClosed) {
682 return PreallocatedPublishResult::QueueStopped;
683 }
684
685 if (!token.m_State.compareAndSwap(availableState, PreallocatedRequest::Claimed)) {
687 return PreallocatedPublishResult::TokenBusy;
688 }
689#if PEDIGREE_CONCURRENCY_SMOKE_TESTS
690 if (m_AfterPreallocatedClaimHook) {
691 m_AfterPreallocatedClaimHook(m_AfterPreallocatedClaimContext);
692 }
693#endif
694#else
695 if (!token.m_State.compareAndSwap(availableState, PreallocatedRequest::Claimed)) {
696 return PreallocatedPublishResult::TokenBusy;
697 }
698#endif
699
700 Request* request = &token.m_Request;
701 request->p1 = p1;
702 request->p2 = p2;
703 request->p3 = p3;
704 request->p4 = p4;
705 request->p5 = p5;
706 request->p6 = p6;
707 request->p7 = p7;
708 request->p8 = p8;
709 request->m_ReturnValue = 0;
710#if THREADS
711 request->m_References = 1;
712#endif
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;
719
720#if !THREADS
721 token.m_State = PreallocatedRequest::Published;
722 executeRequest(p1, p2, p3, p4, p5, p6, p7, p8);
724 return PreallocatedPublishResult::Accepted;
725#else
726 token.m_State = PreallocatedRequest::Published;
727 // Readiness must be visible before the node can become consumable.
728 m_nAsyncRequests += 1;
729 m_nTotalRequests += 1;
730 publishRequest(request);
732 return PreallocatedPublishResult::Accepted;
733#endif
734}
735
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) {
739#if !THREADS
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);
743 discardRequest(request);
744 return 0;
745 }
746 executeRequest(p1, p2, p3, p4, p5, p6, p7, p8);
747 delete request;
748 return 1;
749#else
751 return 0;
752 }
753
754 Request* request = new Request(priority, true, p1, p2, p3, p4, p5, p6, p7, p8);
755
756 if (priority >= REQUEST_QUEUE_NUM_PRIORITIES) {
757 ERROR("RequestQueue '" << m_Name << "' rejected invalid priority " << priority);
758 discardRequest(request);
759 return 0;
760 }
761
762 bool rejected = false;
763 bool overloaded = false;
764 while (true) {
765 bool retry = false;
766 {
767 auto guard = m_RequestQueueWaiters.acquire();
768 if (static_cast<LifecycleState>(m_State.value()) != LifecycleState::Accepting) {
769 rejected = true;
770 } else if (!drainIntakeLocked(priority)) {
771#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
772 m_GuardedTransientRetries += 1;
773#endif
774 retry = true;
775 } else if (findDuplicate(*request)) {
776 rejected = true;
777 } else if (m_nAsyncRequests.value() >= m_nMaxAsyncRequests) {
778 rejected = true;
779 overloaded = true;
780 } else {
781 m_nAsyncRequests += 1;
782 m_nTotalRequests += 1;
783 publishRequest(request);
784 }
785 }
786
787 if (!retry) {
788 break;
789 }
791 }
792
793 if (overloaded) {
794 ERROR("RequestQueue: '" << m_Name << "' is not keeping up with async requests");
795 ERROR(" -> priority=" << priority << ", p1=" << Hex << p1 << ", p2=" << p2 << ", p3=" << p3
796 << ", p4=" << p4);
797 ERROR(" -> p5=" << Hex << p5 << ", p6=" << p6 << ", p7=" << p7 << ", p8=" << p8);
798 }
799 if (rejected) {
800 discardRequest(request);
801 return 0;
802 }
803 return 1;
804#endif
805}
806
808#if THREADS
810 return false;
811 }
812 TerminationDeferral terminationDeferral;
813 LockGuard<Mutex> lifecycleGuard(m_LifecycleMutex);
814 return stopWorker();
815#else
816 return true;
817#endif
818}
819
821#if THREADS
823 return false;
824 }
825 TerminationDeferral terminationDeferral;
826 LockGuard<Mutex> lifecycleGuard(m_LifecycleMutex);
827 if (!startWorker()) {
828 return false;
829 }
830
831 if (!m_pOverrunTimer) {
832 Timer* timer = Machine::instance().getTimer();
833 if (timer && timer->registerHandler(&m_OverrunChecker)) {
834 m_pOverrunTimer = timer;
835 if (timer->supportsDeadlines() &&
836 !timer->armHandler(&m_OverrunChecker, Time::getTicks() + Time::Multiplier::Second)) {
837 FATAL("RequestQueue could not arm its overrun checker");
838 }
839 }
840 }
841 return true;
842#else
843 return true;
844#endif
845}
846
847RequestQueue::LifecycleState RequestQueue::getLifecycleState() {
848 return static_cast<LifecycleState>(m_State.value());
849}
850
852 return RequestQueueCallbackScope::contains(this);
853}
854
857 return false;
858#if THREADS
859 Thread* current = Processor::information().getCurrentThread();
861 return false;
862 auto guard = m_RequestQueueWaiters.acquire();
863 return m_pThread != current;
864#else
865 return true;
866#endif
867}
868
871 return false;
872 TerminationDeferral lifetime;
873#if THREADS
874 while (true) {
875 auto guard = m_RequestQueueWaiters.acquire();
876 if (request.isAvailable())
877 return true;
878 const auto reason =
879 guard.waitForCompletion(WaitQueue::Channel(&request), Thread::CallbackDrain);
880 (void)reason;
881 }
882#else
883 return request.isAvailable();
884#endif
885}
886
888#if THREADS
890 return false;
891 }
892 TerminationDeferral terminationDeferral;
893 Thread* current = Processor::information().getCurrentThread();
894 while (true) {
895 auto guard = m_RequestQueueWaiters.acquire();
896 if (!m_nTotalRequests.value()) {
897 return true;
898 }
899 if (static_cast<LifecycleState>(m_State.value()) != LifecycleState::Accepting) {
900 ERROR("RequestQueue '" << m_Name << "' cannot drain while it is stopping");
901 return false;
902 }
903 if (m_pThread == current) {
904 ERROR("RequestQueue '" << m_Name << "' worker cannot drain itself");
905 return false;
906 }
907
908 const WaitQueue::WakeReason reason = guard.waitForCompletion(
909 WaitQueue::Channel(this, 1), Thread::CallbackDrain, reinterpret_cast<uintptr_t>(this));
910 (void)reason;
911 }
912#else
913 return true;
914#endif
915}
916
918 RequestQueue* queue = reinterpret_cast<RequestQueue*>(p);
919 return queue->work();
920}
921
923#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
925 IntakeLane& lane = m_IntakeLanes[request->m_Priority];
926 const TestAccess::Publication publication =
927 TestAccess::beginPush(lane.m_Queue, request->m_Intake);
928 if (m_AfterIntakeExchangeHook) {
929 m_AfterIntakeExchangeHook(m_AfterIntakeExchangeContext);
930 }
931 TestAccess::finishPush(lane.m_Queue, publication);
932#else
933 m_IntakeLanes[request->m_Priority].m_Queue.push(request->m_Intake);
934#endif
935
936#if THREADS
937 PerProcessorScheduler* scheduler = m_pWorkerScheduler.value();
938 if (scheduler) {
939 scheduler->ringIrqWorkDoorbell(m_WorkerWake);
940 }
941#endif
942}
943
944bool RequestQueue::drainIntakeLocked(size_t priority) {
945 assert(priority < REQUEST_QUEUE_NUM_PRIORITIES);
946 using PopResult = IntrusiveMpscQueue<IntakeNode, &IntakeNode::next>::PopResult;
947
948 while (true) {
949 IntakeNode* node = nullptr;
950 const PopResult result = m_IntakeLanes[priority].m_Queue.pop(node);
951 if (result == PopResult::Empty) {
952 return true;
953 }
954 if (result == PopResult::Transient) {
955 return false;
956 }
957
958 if (!node || !node->owner)
959 return false;
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;
965 } else {
966 m_pRequestQueue[priority] = request;
967 }
968 m_pRequestQueueTail[priority] = request;
969 }
970}
971
972RequestQueue::NextRequestResult RequestQueue::getNextRequest(Request*& out) {
973 out = nullptr;
974 using PopResult = IntrusiveMpscQueue<IntakeNode, &IntakeNode::next>::PopResult;
975
976 for (size_t priority = 0; priority < REQUEST_QUEUE_NUM_PRIORITIES; ++priority) {
977 Request* request = m_pRequestQueue[priority];
978 if (request) {
979 m_pRequestQueue[priority] = request->m_Next;
980 if (!m_pRequestQueue[priority]) {
981 m_pRequestQueueTail[priority] = nullptr;
982 }
983 request->m_Next = nullptr;
984 out = request;
985 return NextRequestResult::Item;
986 }
987
988 IntakeNode* node = nullptr;
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;
993#endif
994 // An unlinked request at this priority must not be overtaken by
995 // work from a lower-priority lane.
996 return NextRequestResult::Retry;
997 }
998 if (result == PopResult::Item) {
999 if (!node || !node->owner)
1000 return NextRequestResult::Retry;
1001 request = node->owner;
1002 assert(request->m_Priority == priority);
1003 out = request;
1004 return NextRequestResult::Item;
1005 }
1006 }
1007
1008 return NextRequestResult::Empty;
1009}
1010
1012 if (m_pActiveRequest && !m_pActiveRequest->m_pPreallocatedOwner &&
1013 m_pActiveRequest->m_Priority == request.m_Priority &&
1014 compareRequests(*m_pActiveRequest, request)) {
1015 return m_pActiveRequest;
1016 }
1017
1018 Request* queued = m_pRequestQueue[request.m_Priority];
1019 while (queued) {
1020 if (!queued->m_pPreallocatedOwner && compareRequests(*queued, request)) {
1021 return queued;
1022 }
1023 queued = queued->m_Next;
1024 }
1025
1026 return nullptr;
1027}
1028
1029void RequestQueue::completeRequest(Request* request, uint64_t returnValue, bool rejected) {
1030#if THREADS
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;
1036 guard.wakeAll();
1037#else
1038 request->m_ReturnValue = returnValue;
1039 request->m_Rejected = rejected;
1040 request->m_Completed = true;
1041#endif
1042}
1043
1045 invokeCancelRequest(*request);
1046 delete request;
1047}
1048
1050 RequestQueueCallbackScope callback(this);
1051 cancelRequest(request);
1052}
1053
1055 assert(request);
1056 PreallocatedRequest* owner = request->m_pPreallocatedOwner;
1057 assert(owner);
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) {
1063 RequestQueueCallbackScope callback(this, owner);
1064 owner->m_ReleaseCallback(owner->m_ReleaseContext);
1065 }
1066
1067 while (true) {
1068 const size_t state = owner->m_State;
1069 if (state == PreallocatedRequest::Releasing) {
1070 if (owner->m_State.compareAndSwap(PreallocatedRequest::Releasing,
1071 PreallocatedRequest::Idle)) {
1072 break;
1073 }
1074 continue;
1075 }
1076 if (state == PreallocatedRequest::Claimed) {
1077 // The claimant now owns every remaining token transition. Waiting
1078 // for it here would make unrelated queue progress depend on the
1079 // scheduling latency of a producer on another CPU.
1080 break;
1081 }
1082
1083 // Published is an asynchronous republication. Idle is a nested inline
1084 // republication when threading is disabled.
1085 assert(state == PreallocatedRequest::Published || state == PreallocatedRequest::Idle);
1086 break;
1087 }
1088#if THREADS
1089 // Availability can immediately release token storage. All later accesses
1090 // belong to the retained queue, including the notification wait guard.
1091 auto guard = m_RequestQueueWaiters.acquire();
1092#endif
1093 owner->m_ReleaseDepth -= 1;
1094#if THREADS
1095 guard.wakeAll(WaitQueue::WakeReason::Signalled, WaitQueue::Channel(owner));
1096#endif
1097}
1098
1099#if THREADS
1101 request->m_References += 1;
1102}
1103
1104void RequestQueue::releaseRequest(Request* request) {
1105 if (request->m_pPreallocatedOwner) {
1107 return;
1108 }
1109
1110 assert(static_cast<size_t>(request->m_References));
1111 if ((request->m_References -= 1) == 0) {
1112 delete request;
1113 }
1114}
1115
1116uint64_t RequestQueue::waitForRequest(Request* request) {
1117 uint64_t result = 0;
1118
1119 while (true) {
1120 auto guard = request->m_Completion.acquire();
1121 if (request->m_Completed) {
1122 if (!request->m_Rejected) {
1123 result = request->m_ReturnValue;
1124 }
1125 break;
1126 }
1127
1128 // A synchronous request transfers payload lifetime to the queue until
1129 // execution completes. Signals and terminal teardown may wake this
1130 // thread, but neither can make that completion contract optional.
1131 WaitQueue::WakeReason reason = guard.waitForCompletion(WaitQueue::Channel(), Thread::CondWait,
1132 reinterpret_cast<uintptr_t>(request));
1133 (void)reason;
1134 }
1135
1136 releaseRequest(request);
1137 return result;
1138}
1139#endif
1140
1142#if THREADS
1143 // The queue, not an unrelated terminal request, owns worker retirement.
1144 // This prevents an idle death from leaving Accepting with no worker and
1145 // prevents active executeRequest state from being abandoned.
1146 TerminationDeferral workerLifetime;
1147 {
1148 auto guard = m_RequestQueueWaiters.acquire();
1149 if (m_pThread != Processor::information().getCurrentThread() ||
1150 static_cast<LifecycleState>(m_State.value()) != LifecycleState::Stopped ||
1151 m_bWorkerReady.value()) {
1152 FATAL("RequestQueue '" << m_Name << "' worker entered with invalid state");
1153 }
1154 m_bWorkerReady = 1;
1155 guard.wakeAll(WaitQueue::WakeReason::Signalled, WaitQueue::Channel(this, 2));
1156 }
1157
1158 while (true) {
1159 Request* request = nullptr;
1160 NextRequestResult next = NextRequestResult::Empty;
1161 bool accepting = false;
1162 m_bWorkerActive = 1;
1163 {
1164 auto guard = m_RequestQueueWaiters.acquire();
1165 const LifecycleState state = static_cast<LifecycleState>(m_State.value());
1166 if (state == LifecycleState::Stopping) {
1167 m_bWorkerReady = 0;
1168 m_bWorkerActive = 0;
1169 return 0;
1170 }
1171 if (state == LifecycleState::Stopped) {
1172 // startWorker() has not yet reached its final acceptance
1173 // point. Stay runnable long enough to publish readiness.
1174 m_bWorkerActive = 0;
1175 } else if (state == LifecycleState::Accepting) {
1176 accepting = true;
1177 next = getNextRequest(request);
1178 if (next == NextRequestResult::Item) {
1179 assert(request);
1180 assert(!m_pActiveRequest);
1181 m_pActiveRequest = request;
1184 } else {
1185 m_bWorkerActive = 0;
1186 }
1187 }
1188 }
1189
1190 if (next != NextRequestResult::Item) {
1191 // A transient MPSC link must be retried without sleeping. An empty
1192 // accepting queue can sleep because every producer publishes its work
1193 // count before ringing this worker's scheduler wake edge.
1194 if (!accepting || next == NextRequestResult::Retry) {
1196 continue;
1197 }
1198
1199 auto waitGuard = m_WorkerWaiters.acquire();
1200 if (!m_nTotalRequests.value() &&
1201 static_cast<LifecycleState>(m_State.value()) == LifecycleState::Accepting) {
1202 const WaitQueue::WakeReason reason = waitGuard.wait(WaitQueue::Channel(), Thread::CondWait,
1203 reinterpret_cast<uintptr_t>(this));
1204 (void)reason;
1205 }
1206 continue;
1207 }
1208
1209 assert(request);
1210
1211 uint64_t result = executeRequest(request->p1, request->p2, request->p3, request->p4,
1212 request->p5, request->p6, request->p7, request->p8);
1213 completeRequest(request, result, false);
1214
1215 {
1216 auto guard = m_RequestQueueWaiters.acquire();
1217 assert(m_pActiveRequest == request);
1218 m_pActiveRequest = nullptr;
1219 assert(m_nTotalRequests.value());
1220 if (request->m_Asynchronous) {
1221 assert(m_nAsyncRequests.value());
1222 m_nAsyncRequests -= 1;
1223 }
1224 }
1225
1226 // Drop the queue's ownership after removing the request from every
1227 // location discoverable by duplicate detection. A preallocated-token
1228 // release callback may republish dependent work, so it runs before
1229 // the final drain predicate is observed.
1230 releaseRequest(request);
1231
1232 {
1233 auto guard = m_RequestQueueWaiters.acquire();
1234 assert(m_nTotalRequests.value());
1235 const size_t remaining = (m_nTotalRequests -= 1);
1237 if (!remaining) {
1238 guard.wakeAll(WaitQueue::WakeReason::Signalled, WaitQueue::Channel(this, 1));
1239 }
1240 }
1241
1242 m_bWorkerActive = 0;
1244 }
1245#else
1246 return 0;
1247#endif
1248}
1249
1250#if THREADS
1251void RequestQueue::RequestQueueOverrunChecker::resetBaselineLocked() {
1252 m_LastQueueSize = 0;
1253 m_LastProgressGeneration = queue->m_WorkerProgressGeneration;
1254 m_HasBacklogBaseline = false;
1255}
1256
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;
1265
1266 const size_t progress = queue->m_WorkerProgressGeneration;
1267 if (static_cast<LifecycleState>(queue->m_State.value()) != LifecycleState::Accepting ||
1268 !currentSize) {
1269 resetBaselineLocked();
1270 return OverrunStatus::Clear;
1271 }
1272
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;
1279 }
1280 }
1281
1282 m_LastQueueSize = currentSize;
1283 m_LastProgressGeneration = progress;
1284 m_HasBacklogBaseline = true;
1285 return status;
1286}
1287
1289 m_Tick += delta;
1290 Timer* source = queue->m_pOverrunTimer;
1291 if (m_Tick < Time::Multiplier::Second) {
1292 if (source && source->supportsDeadlines() &&
1293 !source->armHandler(this, Time::getTicks() + Time::Multiplier::Second - m_Tick)) {
1294 FATAL("RequestQueue could not rearm its overrun checker");
1295 }
1296 return;
1297 }
1298 m_Tick %= Time::Multiplier::Second;
1299 if (source && source->supportsDeadlines() &&
1300 !source->armHandler(this, Time::getTicks() + Time::Multiplier::Second - m_Tick)) {
1301 FATAL("RequestQueue could not rearm its overrun checker");
1302 }
1303
1304 size_t lastSize = 0;
1305 size_t currentSize = 0;
1306 const OverrunStatus status = sample(lastSize, currentSize);
1307 if (status == OverrunStatus::Stalled) {
1308 // A checked writeback can wait on another queue or a device command whose
1309 // deadline exceeds this sampling interval. Lack of completion is diagnostic;
1310 // the operation's own timeout decides whether its backend has failed.
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.");
1317 }
1318}
1319#endif
1320
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);
1326}
1327#endif
virtual Timer * getTimer()=0
bool isOwnedByCurrentThread() const
Definition Mutex.cc:30
void unregisterWorkerWake(SchedulerWorkerWake &worker)
static bool getInterrupts()
static ProcessorInformation & information()
static ExecutionContext executionContext()
Definition Processor.cc:109
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)
Atomic< size_t > m_State
LifecycleState getLifecycleState()
void releasePreallocatedRequest(Request *request)
Mutex m_LifecycleMutex
RequestQueue(const String &name)
virtual void destroy()
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()
Definition Scheduler.h:96
void yield()
Definition Scheduler.cc:236
size_t m_nStateLevel
Definition Thread.h:1177
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)
Definition WaitQueue.cc:104
MUST_USE_RESULT WakeReason waitForCompletion(const Channel &channel=Channel(), size_t debugState=0, uintptr_t debugAddress=0)
Definition WaitQueue.cc:117
@ Hex
Definition Log.h:124
RequestQueueCallbackScope * m_pRequestQueueCallback
Definition Thread.h:1155