The Pedigree Project 0.1
modules/system/concurrency-smoke/main.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/Atomic.h"
9#include "pedigree/kernel/Log.h"
10#include "pedigree/kernel/linker/KernelElf.h"
11#include "pedigree/kernel/process/OperationBarrier.h"
12#include "pedigree/kernel/process/Scheduler.h"
13#include "pedigree/kernel/process/Semaphore.h"
14#include "pedigree/kernel/process/Thread.h"
15#include "pedigree/kernel/processor/Processor.h"
16#include "pedigree/kernel/processor/SyscallHandler.h"
17#include "pedigree/kernel/processor/SyscallManager.h"
18#include "pedigree/kernel/processor/Syscalls.h"
19#include "pedigree/kernel/utilities/ProducerConsumer.h"
20#include "pedigree/kernel/utilities/RequestQueue.h"
21#include "pedigree/kernel/utilities/UniqueResource.h"
22
23#include "modules/Module.h"
24
25extern bool runNetworkFilterConcurrencyRegressions();
26extern bool runAnonymousMemoryRegionRegression();
27extern bool runSlamAllocatorConcurrencyRegression();
28extern bool runTextIoFlipLifetimeRegression();
29extern bool runTlbShootdownConcurrencyRegression();
30extern bool runVfsCallbackLifetimeRegressions();
31extern bool runRcuConcurrencyRegression();
32extern bool runSyscallLifetimeRegression();
33
34namespace {
35class HandoffQueue : public RequestQueue {
36 public:
37 HandoffQueue()
38 : RequestQueue(MakeConstantString("QEMU release handoff")),
39 releaseEntered(0),
40 allowReleaseReturn(0),
41 releaseCalls(0),
42 executions(0),
43 token(released, this) {}
44
45 ~HandoffQueue() override {
46 destroy();
47 }
48
49 static void released(void* context) {
50 HandoffQueue* queue = reinterpret_cast<HandoffQueue*>(context);
51 if ((queue->releaseCalls += 1) == 1) {
52 queue->releaseEntered.release();
53 if (!queue->allowReleaseReturn.acquireForCompletion()) {
54 FATAL("QEMU RequestQueue release callback was interrupted");
55 }
56 }
57 }
58
59 uint64_t executeRequest(uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t,
60 uint64_t) override {
61 executions += 1;
62 return 42;
63 }
64
65 Semaphore releaseEntered;
66 Semaphore allowReleaseReturn;
67 Atomic<size_t> releaseCalls;
68 Atomic<size_t> executions;
69 PreallocatedRequest token;
70};
71
72struct ClaimPause {
73 ClaimPause() : entered(0), release(0), claimCalls(0) {}
74
75 Semaphore entered;
76 Semaphore release;
77 Atomic<size_t> claimCalls;
78};
79
80struct PublishContext {
81 PublishContext(HandoffQueue* requestQueue)
82 : queue(requestQueue), result(RequestQueue::PreallocatedPublishResult::QueueStopped) {}
83
84 HandoffQueue* queue;
85 RequestQueue::PreallocatedPublishResult result;
86};
87
88void pauseClaimedPublication(void* context) {
89 ClaimPause* pause = reinterpret_cast<ClaimPause*>(context);
90 if ((pause->claimCalls += 1) == 1) {
91 pause->entered.release();
92 if (!pause->release.acquireForCompletion()) {
93 FATAL("QEMU RequestQueue publisher pause was interrupted");
94 }
95 }
96}
97
98class BlockingConsumer : public ProducerConsumer {
99 public:
100 BlockingConsumer() : entered(0), release(0), returned(0) {}
101
102 ~BlockingConsumer() override {
103 shutdown();
104 }
105
106 bool start() {
107 return initialise();
108 }
109
110 void shutdown() {
111 destroy();
112 }
113
114 Semaphore entered;
115 Semaphore release;
116 Atomic<size_t> returned;
117
118 private:
119 void consume(uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t,
120 uint64_t) override {
121 entered.release();
122 if (!release.acquireForCompletion()) {
123 FATAL("QEMU ProducerConsumer callback release was interrupted");
124 }
125 returned += 1;
126 }
127};
128
129struct ConsumerDestroyContext {
130 explicit ConsumerDestroyContext(BlockingConsumer* producerConsumer)
131 : consumer(producerConsumer), finished(0) {}
132
133 BlockingConsumer* consumer;
134 Atomic<size_t> finished;
135};
136
137struct TerminalResourceProbe {
138 TerminalResourceProbe(OperationBarrier* barrier, Atomic<size_t>* releases,
139 Atomic<size_t>* releasesBeforeDrain)
140 : barrier(barrier), releases(releases), releasesBeforeDrain(releasesBeforeDrain) {}
141
142 OperationBarrier* barrier;
143 Atomic<size_t>* releases;
144 Atomic<size_t>* releasesBeforeDrain;
145};
146
147struct TerminalResourceProbeReleaser {
148 static void release(TerminalResourceProbe* resource) {
149 *resource->releases += 1;
150 if (!resource->barrier->isClosedAndDrained()) {
151 *resource->releasesBeforeDrain += 1;
152 }
153 }
154};
155
157
158struct TerminalUnwindContext {
159 explicit TerminalUnwindContext(TerminalResourceOwner&& resource)
160 : wait(0),
161 resource(pedigree_std::move(resource)),
162 entered(0),
163 returned(0),
164 interrupted(0),
165 destructed(0) {}
166
167 Semaphore wait;
168 TerminalResourceOwner resource;
169 Atomic<size_t> entered;
170 Atomic<size_t> returned;
171 Atomic<size_t> interrupted;
172 Atomic<size_t> destructed;
173};
174
175class TerminalUnwindCanary {
176 public:
177 explicit TerminalUnwindCanary(Atomic<size_t>& destructed) : m_Destructed(destructed) {}
178
179 ~TerminalUnwindCanary() {
180 m_Destructed += 1;
181 }
182
183 private:
184 Atomic<size_t>& m_Destructed;
185};
186
187int waitForTerminalRequest(void* context) {
188 TerminalUnwindContext* terminal = reinterpret_cast<TerminalUnwindContext*>(context);
189 OperationBarrier::Lease admission;
190 if (!terminal->resource->barrier->tryAcquire(admission)) {
191 FATAL("QEMU terminal request could not admit its worker");
192 }
193 TerminalUnwindCanary canary(terminal->destructed);
194 TerminalResourceOwner resource = pedigree_std::move(terminal->resource);
195 terminal->entered += 1;
196 terminal->interrupted = terminal->wait.acquire() ? 0 : 1;
197 terminal->returned += 1;
198 return terminal->interrupted ? 0 : 1;
199}
200
201int destroyConsumer(void* context) {
202 ConsumerDestroyContext* destroy = reinterpret_cast<ConsumerDestroyContext*>(context);
203 destroy->consumer->shutdown();
204 destroy->finished += 1;
205 return 0;
206}
207
208int publishHandoff(void* context) {
209 PublishContext* publication = reinterpret_cast<PublishContext*>(context);
210 publication->result =
211 publication->queue->republishPreallocatedWhileReleasing(publication->queue->token, 0, 2);
212 return 0;
213}
214
215struct ReciprocalSyscallContext;
216
217class ReciprocalSyscallHandler : public SyscallHandler {
218 public:
219 ReciprocalSyscallHandler(ReciprocalSyscallContext& context, bool first)
220 : m_Context(context), m_First(first) {}
221
222 uintptr_t syscall(SyscallState&) override;
223
224 private:
225 ReciprocalSyscallContext& m_Context;
226 bool m_First;
227};
228
229struct ReciprocalSyscallContext {
230 ReciprocalSyscallContext()
231 : first(*this, true),
232 second(*this, false),
233 entered(0),
234 rejections(0),
235 resetReturns(0),
236 failures(0),
237 firstProcessor(static_cast<size_t>(-1)),
238 secondProcessor(static_cast<size_t>(-1)) {}
239
240 ReciprocalSyscallHandler first;
241 ReciprocalSyscallHandler second;
242 SyscallManager::Registration firstRegistration;
243 SyscallManager::Registration secondRegistration;
244 Atomic<size_t> entered;
245 Atomic<size_t> rejections;
246 Atomic<size_t> resetReturns;
247 Atomic<size_t> failures;
248 Atomic<size_t> firstProcessor;
249 Atomic<size_t> secondProcessor;
250};
251
252uintptr_t ReciprocalSyscallHandler::syscall(SyscallState&) {
253 if (m_First) {
254 m_Context.firstProcessor = Processor::id();
255 } else {
256 m_Context.secondProcessor = Processor::id();
257 }
258 m_Context.entered += 1;
259 for (size_t attempt = 0; attempt < 65536 && m_Context.entered != static_cast<size_t>(2);
260 ++attempt) {
262 }
263 if (m_Context.entered != static_cast<size_t>(2)) {
264 m_Context.failures += 1;
265 return 0;
266 }
267
269 m_First ? m_Context.secondRegistration : m_Context.firstRegistration;
270 if (!peer.reset() && peer) {
271 m_Context.rejections += 1;
272 } else {
273 m_Context.failures += 1;
274 }
275 m_Context.resetReturns += 1;
276 for (size_t attempt = 0; attempt < 65536 && m_Context.resetReturns != static_cast<size_t>(2);
277 ++attempt) {
279 }
280 if (m_Context.resetReturns != static_cast<size_t>(2)) {
281 m_Context.failures += 1;
282 }
283 return m_First ? 0x71 : 0x72;
284}
285
286struct ReciprocalSyscallInvocation {
287 Service_t service;
288 uintptr_t result;
289};
290
291int invokeReciprocalSyscall(void* context) {
292 ReciprocalSyscallInvocation* invocation = reinterpret_cast<ReciprocalSyscallInvocation*>(context);
293 return SyscallManager::instance().dispatchHandlerForTest(invocation->service, invocation->result)
294 ? 0
295 : 1;
296}
297
298void testSyscallReciprocalUnregister() {
299 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN syscall-reciprocal-unregister-smp");
300
302 ReciprocalSyscallContext context;
303 if (!manager.registerSyscallHandler(TUI, &context.first, context.firstRegistration) ||
304 !manager.registerSyscallHandler(native, &context.second, context.secondRegistration)) {
305 if (context.firstRegistration) {
306 context.firstRegistration.reset();
307 }
308 FATAL("QEMU syscall reciprocal handlers could not register");
309 }
310
311 ReciprocalSyscallInvocation first = {TUI, 0};
312 ReciprocalSyscallInvocation second = {native, 0};
313 Process* process = Scheduler::instance().getKernelProcess();
314 Thread* firstThread =
315 new Thread(process, invokeReciprocalSyscall, &first, nullptr, false, false, true);
316 Thread* secondThread =
317 new Thread(process, invokeReciprocalSyscall, &second, nullptr, false, false, true);
318 firstThread->setName("QEMU syscall reciprocal callback A");
319 secondThread->setName("QEMU syscall reciprocal callback B");
320 if (!firstThread->start() || !secondThread->start()) {
321 FATAL("QEMU syscall reciprocal workers did not start");
322 }
323
324 const bool firstJoined = firstThread->joinForCompletion();
325 const bool secondJoined = secondThread->joinForCompletion();
326 const bool registrationsPreserved = context.firstRegistration && context.secondRegistration;
327 const bool firstRetired = context.firstRegistration.reset();
328 const bool secondRetired = context.secondRegistration.reset();
329 if (!firstJoined || !secondJoined || first.result != 0x71 || second.result != 0x72 ||
330 context.entered != static_cast<size_t>(2) || context.rejections != static_cast<size_t>(2) ||
331 context.resetReturns != static_cast<size_t>(2) || context.failures ||
332 context.firstProcessor == context.secondProcessor || !registrationsPreserved ||
333 !firstRetired || !secondRetired) {
334 FATAL("QEMU syscall reciprocal unregister did not preserve external cleanup ownership");
335 }
336
337 NOTICE("QEMU-CONCURRENCY-TEST: syscall reciprocal cpus="
338 << Dec << static_cast<size_t>(context.firstProcessor) << "/"
339 << static_cast<size_t>(context.secondProcessor));
340 NOTICE("QEMU-CONCURRENCY-TEST: PASS syscall-reciprocal-unregister-smp");
341}
342
343void testProducerConsumerTeardown() {
344 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN producerconsumer-teardown-completion");
345
346 BlockingConsumer consumer;
347 if (!consumer.start()) {
348 FATAL("QEMU ProducerConsumer worker did not start");
349 }
350 consumer.produce(1);
351 if (!consumer.entered.acquireForCompletion()) {
352 FATAL("QEMU ProducerConsumer callback did not enter");
353 }
354
355 ConsumerDestroyContext destroy(&consumer);
356 Thread* destroyer = new Thread(Scheduler::instance().getKernelProcess(), destroyConsumer,
357 &destroy, nullptr, false, true);
358 destroyer->setName("QEMU ProducerConsumer destroyer");
359
360 bool joinPublished = false;
361 for (size_t attempt = 0; attempt < 4096; ++attempt) {
362 uintptr_t address = 0;
363 if (destroyer->getStatus() == Thread::Sleeping &&
364 destroyer->getDebugState(address) == Thread::Joining) {
365 joinPublished = true;
366 break;
367 }
369 }
370 if (!joinPublished) {
371 FATAL("QEMU ProducerConsumer destroyer did not join its worker");
372 }
373
375 bool teardownAbandoned = false;
376 for (size_t attempt = 0; attempt < 4096; ++attempt) {
377 if (destroyer->getStatus() == Thread::AwaitingJoin) {
378 teardownAbandoned = true;
379 break;
380 }
382 }
383 if (teardownAbandoned) {
384 FATAL("QEMU ProducerConsumer teardown abandoned its worker join");
385 }
386
387 consumer.release.release();
388 if (!destroyer->joinForCompletion() || destroy.finished.value() != static_cast<size_t>(1) ||
389 consumer.returned.value() != static_cast<size_t>(1)) {
390 FATAL("QEMU ProducerConsumer teardown did not drain its worker");
391 }
392
393 NOTICE("QEMU-CONCURRENCY-TEST: PASS producerconsumer-teardown-completion");
394}
395
396void testTerminalRequestStackUnwind() {
397 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN terminal-request-stack-unwind");
398
399 OperationBarrier barrier;
400 Atomic<size_t> releases(0);
401 Atomic<size_t> releasesBeforeDrain(0);
402 TerminalResourceProbe resource(&barrier, &releases, &releasesBeforeDrain);
403 TerminalUnwindContext context(TerminalResourceOwner::adopt(&resource));
404 if (!barrier.tryEnter()) {
405 FATAL("QEMU terminal request could not admit its drain probe");
406 }
407 Thread* worker = new Thread(Scheduler::instance().getKernelProcess(), waitForTerminalRequest,
408 &context, nullptr, false, true);
409 worker->setName("QEMU terminal request stack unwind");
410
411 bool waitPublished = false;
412 for (size_t attempt = 0; worker && attempt < 4096; ++attempt) {
413 uintptr_t address = 0;
414 if (worker->getStatus() == Thread::Sleeping &&
415 worker->getDebugState(address) == Thread::SemWait) {
416 waitPublished = true;
417 break;
418 }
420 }
421
422 barrier.close();
423 if (barrier.tryEnter()) {
424 FATAL("QEMU closed operation barrier admitted late work");
425 }
426 barrier.leave();
427 if (barrier.isClosedAndDrained()) {
428 FATAL("QEMU operation barrier drained before its worker finished");
429 }
430 if (worker) {
431 worker->setUnwindState(Thread::TerminateThread);
432 }
433 barrier.wait();
434 const bool joined = worker->joinForCompletion();
435 if (!joined || !waitPublished || context.entered.value() != static_cast<size_t>(1) ||
436 context.interrupted.value() != static_cast<size_t>(1) ||
437 context.returned.value() != static_cast<size_t>(1) ||
438 context.destructed.value() != static_cast<size_t>(1) ||
439 releases.value() != static_cast<size_t>(1) ||
440 releasesBeforeDrain.value() != static_cast<size_t>(1) || !barrier.isClosedAndDrained()) {
441 FATAL("QEMU terminal request did not unwind the worker's C++ stack");
442 }
443
444 NOTICE("QEMU-CONCURRENCY-TEST: PASS terminal-request-stack-unwind");
445}
446
447void testPinnedLinkerUnloadRejection() {
448 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN linker-pinned-unload-rejection");
449
450 static char linkerName[] = "linker";
451 KernelElf& kernelElf = KernelElf::instance();
452 if (!kernelElf.moduleIsLoaded(linkerName) || kernelElf.unloadModule(linkerName, true, false) ||
453 !kernelElf.moduleIsLoaded(linkerName)) {
454 FATAL("QEMU linker module did not reject unload while preserving its image");
455 }
456
457 NOTICE("QEMU-CONCURRENCY-TEST: PASS linker-pinned-unload-rejection");
458}
459
460void testFilesystemModuleUnloadRejection() {
461 KernelElf& kernelElf = KernelElf::instance();
462
463 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN posix-runtime-unload-rejection");
464 static char posixName[] = "posix";
465 if (!kernelElf.moduleIsLoaded(posixName) || kernelElf.unloadModule(posixName, true, false) ||
466 !kernelElf.moduleIsLoaded(posixName)) {
467 FATAL("QEMU POSIX module did not reject runtime unload while preserving its image");
468 }
469 NOTICE("QEMU-CONCURRENCY-TEST: PASS posix-runtime-unload-rejection");
470
471 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN rawfs-runtime-unload-rejection");
472 static char rawfsName[] = "rawfs";
473 const bool rawfsLoadedBefore = kernelElf.moduleIsLoaded(rawfsName);
474 const bool rawfsUnloadAccepted = kernelElf.unloadModule(rawfsName, true, false);
475 const bool rawfsLoadedAfter = kernelElf.moduleIsLoaded(rawfsName);
476 if (!rawfsLoadedBefore || rawfsUnloadAccepted || !rawfsLoadedAfter) {
477 FATAL("QEMU rawfs module did not reject runtime unload while preserving its image");
478 }
479 NOTICE("QEMU-CONCURRENCY-TEST: PASS rawfs-runtime-unload-rejection");
480
481 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN ramfs-pinned-unload-rejection");
482 static char ramfsName[] = "ramfs";
483 if (!kernelElf.moduleIsLoaded(ramfsName) || kernelElf.unloadModule(ramfsName, true, false) ||
484 !kernelElf.moduleIsLoaded(ramfsName)) {
485 FATAL("QEMU RamFs module did not reject unload while preserving its image");
486 }
487 NOTICE("QEMU-CONCURRENCY-TEST: PASS ramfs-pinned-unload-rejection");
488}
489
490bool entry() {
491 if (!runSyscallLifetimeRegression()) {
492 FATAL("QEMU syscall handler lifetime regression failed");
493 }
494 if (!runRcuConcurrencyRegression()) {
495 FATAL("QEMU RCU publication and reclamation regression failed");
496 }
497 testPinnedLinkerUnloadRejection();
498 testFilesystemModuleUnloadRejection();
499 testSyscallReciprocalUnregister();
500 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN network-filter-reciprocal-removal-smp");
501 if (!runNetworkFilterConcurrencyRegressions()) {
502 FATAL("QEMU NetworkFilter reciprocal-removal regression failed");
503 }
504 testProducerConsumerTeardown();
505 testTerminalRequestStackUnwind();
506
507 if (!runAnonymousMemoryRegionRegression()) {
508 FATAL("QEMU anonymous MemoryRegion regression failed");
509 }
510 if (!runTlbShootdownConcurrencyRegression()) {
511 FATAL("QEMU TLB shootdown regression failed");
512 }
513 if (!runSlamAllocatorConcurrencyRegression()) {
514 FATAL("QEMU SLAM allocator concurrency regression failed");
515 }
516 if (!runTextIoFlipLifetimeRegression()) {
517 FATAL("QEMU TextIO flip-worker lifetime regression failed");
518 }
519 if (!runVfsCallbackLifetimeRegressions()) {
520 FATAL("QEMU VFS callback lifetime regression failed");
521 }
522
523 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN requestqueue-release-handoff");
524
525 HandoffQueue queue;
526 ClaimPause pause;
527 queue.initialise();
528
529 const RequestQueue::PreallocatedPublishResult initial =
530 queue.publishPreallocated(queue.token, 0, 1);
531 if (initial != RequestQueue::PreallocatedPublishResult::Accepted ||
532 !queue.releaseEntered.acquireForCompletion()) {
533 FATAL("QEMU RequestQueue handoff setup failed");
534 }
535 queue.setAfterPreallocatedClaimHookForTest(pauseClaimedPublication, &pause);
536
537 PublishContext publication(&queue);
538 Thread* publisher = new Thread(Scheduler::instance().getKernelProcess(), publishHandoff,
539 &publication, nullptr, false, true, true);
540 publisher->setName("QEMU RequestQueue handoff publisher");
541 if (!publisher->start() || !pause.entered.acquireForCompletion()) {
542 FATAL("QEMU RequestQueue handoff publisher did not claim its token");
543 }
544
545 queue.allowReleaseReturn.release();
546
547 if (!queue.addAsyncRequest(0, 3)) {
548 FATAL("QEMU RequestQueue handoff progress probe was rejected");
549 }
550
551 bool independentProgress = false;
552 for (size_t attempt = 0; attempt < 4096; ++attempt) {
553 if (queue.executions.value() >= 2) {
554 independentProgress = true;
555 break;
556 }
558 }
559
560 pause.release.release();
561
562 const bool joined = publisher->joinForCompletion();
563 const bool drained = queue.drain();
564 queue.setAfterPreallocatedClaimHookForTest(nullptr, nullptr);
565 queue.destroy();
566
567 if (!independentProgress || !drained || !joined ||
568 publication.result != RequestQueue::PreallocatedPublishResult::Accepted ||
569 queue.executions.value() != 3 || queue.releaseCalls.value() != 2 ||
570 !queue.token.isAvailable()) {
571 FATAL("QEMU RequestQueue handoff blocked independent queue progress");
572 }
573
574 NOTICE("QEMU-CONCURRENCY-TEST: PASS requestqueue-release-handoff");
575 NOTICE("QEMU-CONCURRENCY-TEST: BEGIN shutdown-worker-drain-smp");
576 return true;
577}
578
579void exit() {}
580} // namespace
581
582MODULE_INFO("concurrency-smoke", &entry, &exit, "console", "vfs", "rawfs", "network-stack");
bool unloadModule(const char *name, bool silent=false, bool progress=true)
Definition KernelElf.cc:929
bool moduleIsLoaded(char *name)
static KernelElf & instance()
Definition KernelElf.h:135
MUST_USE_RESULT bool tryEnter()
static ProcessorId id()
virtual void destroy()
static Scheduler & instance()
Definition Scheduler.h:96
void yield()
Definition Scheduler.cc:236
void release(size_t n=1)
Definition Semaphore.cc:549
virtual uintptr_t syscall(SyscallState &State)=0
virtual bool registerSyscallHandler(Service_t Service, SyscallHandler *pHandler, Registration &registration, FastEntry entry=nullptr)=0
static EXPORTED_PUBLIC SyscallManager & instance()
void setUnwindState(UnwindType ut)
Definition Thread.cc:3580
@ TerminateThread
Exit only this thread during Process exit.
Definition Thread.h:519
bool joinForCompletion()
Definition Thread.cc:2750
Status getStatus() const
Definition Thread.h:433
DebugState getDebugState(uintptr_t &address)
Definition Thread.h:574
bool start()
Definition Thread.cc:751
@ Dec
Definition Log.h:126