The Pedigree Project 0.1
WaitQueue.cc
1/*
2 * Copyright (c) 2026, Pedigree Developers
3 *
4 * Permission to use, copy, modify, and distribute this software for any
5 * purpose with or without fee is hereby granted.
6 */
7
8#include "pedigree/kernel/Log.h"
9#include "pedigree/kernel/process/Mutex.h"
10#include "pedigree/kernel/process/PerProcessorScheduler.h"
11#include "pedigree/kernel/process/Thread.h"
12#include "pedigree/kernel/process/WaitQueue.h"
13#include "pedigree/kernel/processor/Processor.h"
14#include "pedigree/kernel/processor/ProcessorInformation.h"
15#include "pedigree/kernel/utilities/Iterator.h"
16#include "pedigree/kernel/utilities/assert.h"
17
18#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
19WaitQueue::BeforeBlockHook WaitQueue::m_BeforeBlockHook = nullptr;
20#endif
21
22namespace {
23#if PEDIGREE_AFFINITY_TESTS
24Thread* g_ReadyPublicationTarget = nullptr;
25WaitQueue::ReadyPublicationHook g_ReadyPublicationHook = nullptr;
26#endif
27WaitQueue::WakeReason terminalWakeReason(Thread::UnwindType state) {
28 if (state == Thread::Exit) {
29 return WaitQueue::WakeReason::Unwinding;
30 }
31 if (state == Thread::TerminateThread) {
32 return WaitQueue::WakeReason::Terminating;
33 }
34 FATAL("Unknown terminal WaitQueue unwind state.");
35 return WaitQueue::WakeReason::Spurious;
36}
37} // namespace
38
39WaitQueue::Guard::Guard(WaitQueue& queue)
40 : m_Queue(&queue), m_OwnsLock(false), m_pFirstReady(nullptr), m_pLastReady(nullptr) {
41 if (!Processor::guardDeviceHardIrqOperation(DeviceHardIrqOperation::WaitQueueAccess)) {
42 return;
43 }
44
45 m_Queue->m_Lock.acquire();
46 m_OwnsLock = true;
47}
48
49WaitQueue::Guard::Guard(Guard&& other) noexcept
50 : m_Queue(other.m_Queue),
51 m_OwnsLock(other.m_OwnsLock),
52 m_pFirstReady(other.m_pFirstReady),
53 m_pLastReady(other.m_pLastReady) {
54 other.m_Queue = nullptr;
55 other.m_OwnsLock = false;
56 other.m_pFirstReady = nullptr;
57 other.m_pLastReady = nullptr;
58}
59
60WaitQueue::Guard::~Guard() {
61 release();
62}
63
64void WaitQueue::Guard::release() {
65 if (m_Queue && m_OwnsLock) {
66 m_Queue->clearWaitIntentIfEmpty();
67 m_Queue->m_Lock.release();
68 m_OwnsLock = false;
69
70 // Ready publication takes scheduler-owned locks. Keep that ordering
71 // outside the WaitQueue lock because scheduler lifecycle paths can
72 // themselves complete waits.
73 while (m_pFirstReady) {
74 Waiter* waiter = m_pFirstReady;
75 m_pFirstReady = waiter->notificationNext;
76 waiter->notificationNext = nullptr;
77 WaitQueue::publishReady(waiter);
78 }
79 m_pLastReady = nullptr;
80 }
81}
82
84 if (m_OwnsLock) {
85 // Paired with the producer's predicate publication and intent load.
86 // Queue membership alone is too late: the final predicate check must
87 // already be protected against a producer skipping the queue lock.
88 __atomic_store_n(&m_Queue->m_WaitIntent, true, __ATOMIC_SEQ_CST);
89 }
90}
91
92void WaitQueue::Guard::queueSchedulerNotification(Waiter* waiter) {
93 assert(waiter);
94 assert(!waiter->notificationNext);
95 if (m_pLastReady) {
96 m_pLastReady->notificationNext = waiter;
97 } else {
98 m_pFirstReady = waiter;
99 }
100 m_pLastReady = waiter;
101}
102
103WaitQueue::WakeReason WaitQueue::Guard::wait(const Channel& channel, size_t debugState,
104 uintptr_t debugAddress,
105 StackDiscardCleanup onStackDiscard,
106 void* stackDiscardContext) {
107 assert(m_Queue);
108 if (!m_OwnsLock) {
109 return WakeReason::Spurious;
110 }
111 assert(m_OwnsLock);
112 Thread::StackDiscardScope discardScope(onStackDiscard, stackDiscardContext);
113 return m_Queue->wait(*this, nullptr, channel, debugState, debugAddress, false, true);
114}
115
116WaitQueue::WakeReason WaitQueue::Guard::waitForCompletion(const Channel& channel, size_t debugState,
117 uintptr_t debugAddress) {
118 assert(m_Queue);
119 if (!m_OwnsLock) {
120 return WakeReason::Spurious;
121 }
122 assert(m_OwnsLock);
123 return m_Queue->wait(*this, nullptr, channel, debugState, debugAddress, true, true);
124}
125
126WaitQueue::WakeReason WaitQueue::Guard::waitWithoutEventDispatch(const Channel& channel,
127 size_t debugState,
128 uintptr_t debugAddress) {
129 assert(m_Queue);
130 if (!m_OwnsLock) {
131 return WakeReason::Spurious;
132 }
133 assert(m_OwnsLock);
134 return m_Queue->wait(*this, nullptr, channel, debugState, debugAddress, false, false);
135}
136
137WaitQueue::WakeReason WaitQueue::Guard::waitAndUnlock(Mutex& mutex, const Channel& channel,
138 size_t debugState, uintptr_t debugAddress,
139 StackDiscardCleanup onStackDiscard,
140 void* stackDiscardContext) {
141 assert(m_Queue);
142 if (!m_OwnsLock) {
143 return WakeReason::Spurious;
144 }
145 assert(m_OwnsLock);
146 Thread::StackDiscardScope discardScope(onStackDiscard, stackDiscardContext);
147 return m_Queue->wait(*this, &mutex, channel, debugState, debugAddress, false, true);
148}
149
151 const Channel& channel,
152 size_t debugState,
153 uintptr_t debugAddress) {
154 assert(m_Queue);
155 if (!m_OwnsLock) {
156 return WakeReason::Spurious;
157 }
158 assert(m_OwnsLock);
159 return m_Queue->wait(*this, &mutex, channel, debugState, debugAddress, true, true);
160}
161
162bool WaitQueue::Guard::wakeOne(WakeReason reason, const Channel& channel) {
163 assert(m_Queue);
164 if (!m_OwnsLock) {
165 return false;
166 }
167 assert(m_OwnsLock);
168 assert(reason != WakeReason::Waiting);
169 return m_Queue->wakeOneLocked(*this, reason, channel);
170}
171
172size_t WaitQueue::Guard::wakeAll(WakeReason reason, const Channel& channel) {
173 assert(m_Queue);
174 if (!m_OwnsLock) {
175 return 0;
176 }
177 assert(m_OwnsLock);
178 assert(reason != WakeReason::Waiting);
179 return m_Queue->wakeAllLocked(*this, reason, channel);
180}
181
182size_t WaitQueue::Guard::wakeAndRequeue(const Channel& source, size_t wakeCount,
183 const Channel& destination, size_t requeueCount) {
184 assert(m_Queue);
185 if (!m_OwnsLock) {
186 return 0;
187 }
188 assert(m_OwnsLock);
189 return m_Queue->wakeAndRequeueLocked(*this, source, wakeCount, destination, requeueCount);
190}
191
192WaitQueue::WaitQueue()
193 : m_Lock(false),
194 m_pFirstWaiter(nullptr),
195 m_pLastWaiter(nullptr),
196 m_WaiterCount(0),
197 m_WaitIntent(false) {}
198
199WaitQueue::~WaitQueue() {
200 if (waiterCount()) {
201 FATAL("Destroying a WaitQueue with live waiters.");
202 }
203}
204
205WaitQueue::WakeReason WaitQueue::wait(Guard& guard, Mutex* mutex, const Channel& channel,
206 size_t debugState, uintptr_t debugAddress, bool deferTerminal,
207 bool dispatchEvents) {
208 if (mutex && !mutex->isOwnedByCurrentThread()) {
209 FATAL(
210 "WaitQueue::waitAndUnlock requires current-thread mutex "
211 "ownership");
212 }
213
214 Thread* thread = Processor::information().getCurrentThread();
215 assert(thread);
216 const size_t stateLevel = thread->getStateLevel();
217 Waiter& waiter = thread->m_StateLevels[stateLevel].m_Waiter;
218
219 thread->m_Lock.acquire();
220 if (waiter.loadQueue()) {
221 FATAL("Thread attempted to enter two wait queues at once.");
222 }
223 thread->clearTerminalWaitCancelledBeforeBlockUnlocked(stateLevel);
224 const Thread::UnwindType unwindState = thread->getUnwindState();
225 if (unwindState != Thread::Continue && !deferTerminal) {
226 thread->m_Lock.release();
227 guard.release();
228 return terminalWakeReason(unwindState);
229 }
230 waiter.thread = thread;
231 waiter.scheduler = thread->m_pScheduler;
232 assert(waiter.scheduler);
233 waiter.storeChannel(channel);
234 waiter.stateLevel = stateLevel;
235 waiter.storeReason(WakeReason::Waiting);
236 waiter.setQueued(false);
237 waiter.notificationNext = nullptr;
238 waiter.previous = m_pLastWaiter;
239 waiter.next = nullptr;
240 thread->setDebugState(static_cast<Thread::DebugState>(debugState), debugAddress);
241
242 if (m_pLastWaiter) {
243 m_pLastWaiter->next = &waiter;
244 } else {
245 m_pFirstWaiter = &waiter;
246 }
247 m_pLastWaiter = &waiter;
248 ++m_WaiterCount;
249 waiter.setQueued(true);
250
251 // Publish only after every field and the queue membership are complete.
252 // Debugger snapshots never need to take the target thread's spinlock.
253 waiter.storeQueue(this);
254 thread->m_Lock.release();
255
256 if (mutex) {
257 mutex->release();
258 }
259
260 // The persistent waiter is visible before the queue lock is released.
261 // Any wake after this point changes waiter.reason, even if the thread has
262 // not yet committed the Sleeping state.
263 guard.release();
264
265#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
266 if (m_BeforeBlockHook) {
267 m_BeforeBlockHook(this, thread, channel, debugState);
268 }
269#endif
270
271 Processor::information().getScheduler().blockCurrent();
272
273 m_Lock.acquire();
274 thread->m_Lock.acquire();
275 thread->clearTerminalWaitCancelledBeforeBlockUnlocked(stateLevel);
276 if (waiter.loadQueue() == this) {
277 // Unpublish before the persistent record can be reused.
278 waiter.storeQueue(nullptr);
279 }
280 WakeReason reason = waiter.loadReason();
281 waiter.scheduler = nullptr;
282 thread->setDebugState(Thread::None, 0);
283 thread->m_Lock.release();
284 removeWaiterLocked(&waiter);
285 m_Lock.release();
286
287 if (reason == WakeReason::Waiting) {
288 reason = WakeReason::Spurious;
289 }
290
291 Thread::UnwindType terminalState = thread->getUnwindState();
292 if (terminalState != Thread::Continue && !deferTerminal) {
293 reason = terminalWakeReason(terminalState);
294 }
295
296 // Event handlers may themselves block. Dispatch only after removing the
297 // outer wait record so nested event state gets an independent wait.
298 // An ordinary wake and an event publication can race. Once the outer wait
299 // record is retired, dispatch any event which is now deliverable regardless
300 // of which wake reason won.
301 if (dispatchEvents && (terminalState == Thread::Continue || deferTerminal)) {
302 thread->m_StateLevels[stateLevel].m_bDispatchingWaitEvent = true;
303 Processor::information().getScheduler().checkEventState(0);
304 thread->m_StateLevels[stateLevel].m_bDispatchingWaitEvent = false;
305 }
306
307 // An event handler can itself request process termination. Do not enter a
308 // fresh external-mutex wait after that terminal decision.
309 terminalState = thread->getUnwindState();
310 if (terminalState != Thread::Continue && !deferTerminal) {
311 reason = terminalWakeReason(terminalState);
312 }
313
314 if (mutex) {
315 // Event delivery happens before reacquiring the external mutex, avoiding
316 // re-entry into a handler that needs the same lock. Reacquisition is
317 // itself an ownership barrier: an outer signal/timeout marker is
318 // retained, and terminal propagation happens only after mutex ownership
319 // has been restored.
320 const bool acquired = mutex->acquireForCompletion();
321 if (!acquired) {
322 FATAL("WaitQueue could not reacquire its caller mutex");
323 }
324
325 terminalState = thread->getUnwindState();
326 if (terminalState != Thread::Continue && !deferTerminal) {
327 reason = terminalWakeReason(terminalState);
328 }
329 }
330 return reason;
331}
332
333bool WaitQueue::wakeOne(WakeReason reason, const Channel& channel) {
334 Guard guard(*this);
335 return guard.wakeOne(reason, channel);
336}
337
338size_t WaitQueue::wakeAll(WakeReason reason, const Channel& channel) {
339 Guard guard(*this);
340 return guard.wakeAll(reason, channel);
341}
342
343size_t WaitQueue::wakeAllIfWaiting(WakeReason reason, const Channel& channel) {
344 if (!Processor::guardDeviceHardIrqOperation(DeviceHardIrqOperation::WaitQueueAccess)) {
345 return 0;
346 }
347 if (!__atomic_load_n(&m_WaitIntent, __ATOMIC_SEQ_CST)) {
348 return 0;
349 }
350 return wakeAll(reason, channel);
351}
352
353void WaitQueue::clearWaitIntentIfEmpty() {
354 // The queue lock excludes another waiter publishing intent until unlock.
355 // Clear on both abandoned enrollment and retirement of the last waiter.
356 if (!m_WaiterCount && __atomic_load_n(&m_WaitIntent, __ATOMIC_RELAXED)) {
357 __atomic_store_n(&m_WaitIntent, false, __ATOMIC_SEQ_CST);
358 }
359}
360
361bool WaitQueue::wakeOneLocked(Guard& guard, WakeReason reason, const Channel& channel) {
362 assert(reason == WakeReason::Signalled || reason == WakeReason::Event ||
363 reason == WakeReason::Spurious);
364 for (Waiter* waiter = m_pFirstWaiter; waiter; waiter = waiter->next) {
365 if (!(waiter->channel == channel)) {
366 continue;
367 }
368
369 if (completeWaiter(guard, waiter, reason)) {
370 return true;
371 }
372 }
373
374 return false;
375}
376
377size_t WaitQueue::wakeAllLocked(Guard& guard, WakeReason reason, const Channel& channel) {
378 assert(reason == WakeReason::Signalled || reason == WakeReason::Event ||
379 reason == WakeReason::Spurious);
380 size_t count = 0;
381 for (Waiter* waiter = m_pFirstWaiter; waiter; waiter = waiter->next) {
382 if (!(waiter->channel == channel)) {
383 continue;
384 }
385
386 if (completeWaiter(guard, waiter, reason)) {
387 ++count;
388 }
389 }
390
391 return count;
392}
393
394size_t WaitQueue::wakeAndRequeueLocked(Guard& guard, const Channel& source, size_t wakeCount,
395 const Channel& destination, size_t requeueCount) {
396 size_t woken = 0;
397 size_t requeued = 0;
398 for (Waiter* waiter = m_pFirstWaiter; waiter; waiter = waiter->next) {
399 if (!(waiter->channel == source) || waiter->loadReason() != WakeReason::Waiting) {
400 continue;
401 }
402
403 if (woken < wakeCount) {
404 if (completeWaiter(guard, waiter, WakeReason::Signalled)) {
405 ++woken;
406 }
407 continue;
408 }
409
410 if (requeued < requeueCount) {
411 Thread* thread = waiter->thread;
412 bool moved = false;
413 thread->m_Lock.acquire();
414 if (thread->m_StateLevels[waiter->stateLevel].m_Waiter.loadQueue() == this &&
415 waiter->loadReason() == WakeReason::Waiting && waiter->channel == source) {
416 // Queue-before-thread is the same lock order used by wake and
417 // cancellation. Diagnostics reject a channel being replaced.
418 waiter->storeChannel(destination);
419 moved = true;
420 }
421 thread->m_Lock.release();
422 if (moved) {
423 ++requeued;
424 }
425 }
426 }
427
428 return woken + requeued;
429}
430
431bool WaitQueue::completeWaiter(Guard& guard, Waiter* waiter, WakeReason reason) {
432 Thread* thread = waiter->thread;
433 bool becameReady = false;
434 bool completed = false;
435
436 thread->m_Lock.acquire();
437 if (thread->m_StateLevels[waiter->stateLevel].m_Waiter.loadQueue() == this &&
438 waiter->loadReason() == WakeReason::Waiting) {
439 waiter->storeReason(reason);
440 completed = true;
441 if (thread->m_Status == Thread::Sleeping) {
442 thread->m_Status = Thread::Ready;
443 __atomic_store_n(&thread->m_ReadyPublicationPending, true, __ATOMIC_RELEASE);
444 becameReady = true;
445 }
446 }
447 thread->m_Lock.release();
448
449 if (becameReady) {
450 guard.queueSchedulerNotification(waiter);
451 }
452 return completed;
453}
454
455void WaitQueue::publishReady(Waiter* waiter) {
456 assert(waiter);
457 Thread* thread = waiter->thread;
458 assert(thread);
459#if PEDIGREE_AFFINITY_TESTS
460 const auto hook = __atomic_load_n(&g_ReadyPublicationHook, __ATOMIC_ACQUIRE);
461 if (hook && thread == __atomic_load_n(&g_ReadyPublicationTarget, __ATOMIC_ACQUIRE))
462 hook(thread);
463#endif
464 thread->publishReadyNotification();
465}
466
467#if PEDIGREE_AFFINITY_TESTS
468void WaitQueue::setReadyPublicationHookForTest(Thread* target, ReadyPublicationHook hook) {
469 __atomic_store_n(&g_ReadyPublicationTarget, target, __ATOMIC_RELEASE);
470 __atomic_store_n(&g_ReadyPublicationHook, hook, __ATOMIC_RELEASE);
471}
472#endif
473
474void WaitQueue::removeWaiterLocked(Waiter* waiter) {
475 if (!waiter->isQueued()) {
476 return;
477 }
478
479 if (waiter->previous) {
480 waiter->previous->next = waiter->next;
481 } else {
482 assert(m_pFirstWaiter == waiter);
483 m_pFirstWaiter = waiter->next;
484 }
485
486 if (waiter->next) {
487 waiter->next->previous = waiter->previous;
488 } else {
489 assert(m_pLastWaiter == waiter);
490 m_pLastWaiter = waiter->previous;
491 }
492
493 assert(m_WaiterCount);
494 --m_WaiterCount;
495 clearWaitIntentIfEmpty();
496 waiter->previous = nullptr;
497 waiter->next = nullptr;
498 waiter->setQueued(false);
499}
500
501void WaitQueue::cancel(Waiter* waiter, WakeReason reason) {
502 Thread* thread = waiter->thread;
503 bool makeReady = false;
504 {
505 Guard guard(*this);
506 if (!guard.m_OwnsLock) {
507 return;
508 }
509 thread->m_Lock.acquire();
510 if (waiter->loadQueue() != this) {
511 thread->m_Lock.release();
512 return;
513 }
514
515 removeWaiterLocked(waiter);
516 waiter->storeReason(reason);
517 makeReady = thread->m_Status == Thread::Sleeping;
518
519 if (reason == WakeReason::Terminating) {
520 waiter->storeQueue(nullptr);
521 if (!makeReady && thread->m_Status == Thread::Running) {
522 // blockCurrent() has not committed Sleeping yet. Keep no
523 // linked waiter, but leave an exact one-shot handoff on the
524 // waiter's state level. A nested event may have interrupted
525 // that level between publication and blockCurrent().
527 }
528 }
529 thread->m_Lock.release();
530 }
531
532 bool becameReady = false;
533 if (makeReady) {
534 thread->m_Lock.acquire();
535 if (thread->m_Status == Thread::Sleeping) {
536 thread->m_Status = Thread::Ready;
537 __atomic_store_n(&thread->m_ReadyPublicationPending, true, __ATOMIC_RELEASE);
538 becameReady = true;
539 }
540 thread->m_Lock.release();
541 }
542
543 if (becameReady) {
544 publishReady(waiter);
545 }
546}
547
548size_t WaitQueue::waiterCount() {
549 Guard guard(*this);
550 if (!guard.m_OwnsLock) {
551 return 0;
552 }
553 return m_WaiterCount;
554}
555
556#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
557void WaitQueue::setBeforeBlockHook(BeforeBlockHook hook) {
558 m_BeforeBlockHook = hook;
559}
560#endif
Definition Mutex.h:56
static ProcessorInformation & information()
static bool guardDeviceHardIrqOperation(DeviceHardIrqOperation operation)
Definition Processor.h:563
void release()
Definition Spinlock.cc:161
bool acquire(bool recurse=false, bool safe=true)
Definition Spinlock.cc:35
UnwindType
Definition Thread.h:512
@ Continue
No unwind necessary, carry on as normal.
Definition Thread.h:513
@ TerminateThread
Exit only this thread during Process exit.
Definition Thread.h:515
@ Exit
Exit the owning process at the next safe boundary.
Definition Thread.h:514
void setDebugState(DebugState state, uintptr_t address)
Definition Thread.h:589
volatile Status m_Status
Definition Thread.h:1303
DebugState
Definition Thread.h:176
UnwindType getUnwindState()
Definition Thread.h:531
Spinlock m_Lock
Definition Thread.h:1257
void markTerminalWaitCancelledBeforeBlockUnlocked(size_t level)
Definition Thread.cc:3743
size_t getStateLevel() const
Definition Thread.h:314
void prepareToWait()
Definition WaitQueue.cc:83
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:103
MUST_USE_RESULT WakeReason waitForCompletion(const Channel &channel=Channel(), size_t debugState=0, uintptr_t debugAddress=0)
Definition WaitQueue.cc:116
size_t wakeAndRequeue(const Channel &source, size_t wakeCount, const Channel &destination, size_t requeueCount)
Definition WaitQueue.cc:182
MUST_USE_RESULT WakeReason waitWithoutEventDispatch(const Channel &channel, size_t debugState, uintptr_t debugAddress)
Definition WaitQueue.cc:126
MUST_USE_RESULT WakeReason waitAndUnlock(Mutex &mutex, const Channel &channel=Channel(), size_t debugState=0, uintptr_t debugAddress=0, StackDiscardCleanup onStackDiscard=nullptr, void *stackDiscardContext=nullptr)
Definition WaitQueue.cc:137
MUST_USE_RESULT WakeReason waitAndUnlockForCompletion(Mutex &mutex, const Channel &channel=Channel(), size_t debugState=0, uintptr_t debugAddress=0)
Definition WaitQueue.cc:150
size_t wakeAllIfWaiting(WakeReason reason=WakeReason::Signalled, const Channel &channel=Channel())
Definition WaitQueue.cc:343
WaitQueue::Waiter m_Waiter
Definition Thread.h:1166
bool m_bDispatchingWaitEvent
Definition Thread.h:1145