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"
23#include "modules/Module.h"
25extern bool runNetworkFilterConcurrencyRegressions();
26extern bool runAnonymousMemoryRegionRegression();
27extern bool runSlamAllocatorConcurrencyRegression();
28extern bool runTextIoFlipLifetimeRegression();
29extern bool runTlbShootdownConcurrencyRegression();
30extern bool runVfsCallbackLifetimeRegressions();
31extern bool runRcuConcurrencyRegression();
32extern bool runSyscallLifetimeRegression();
38 :
RequestQueue(MakeConstantString(
"QEMU release handoff")),
40 allowReleaseReturn(0),
43 token(released, this) {}
45 ~HandoffQueue()
override {
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");
59 uint64_t executeRequest(uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t,
69 PreallocatedRequest token;
73 ClaimPause() : entered(0), release(0), claimCalls(0) {}
80struct PublishContext {
81 PublishContext(HandoffQueue* requestQueue)
82 : queue(requestQueue), result(
RequestQueue::PreallocatedPublishResult::QueueStopped) {}
85 RequestQueue::PreallocatedPublishResult result;
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");
100 BlockingConsumer() : entered(0), release(0), returned(0) {}
102 ~BlockingConsumer()
override {
119 void consume(uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t, uint64_t,
122 if (!release.acquireForCompletion()) {
123 FATAL(
"QEMU ProducerConsumer callback release was interrupted");
129struct ConsumerDestroyContext {
130 explicit ConsumerDestroyContext(BlockingConsumer* producerConsumer)
131 : consumer(producerConsumer), finished(0) {}
133 BlockingConsumer* consumer;
137struct TerminalResourceProbe {
140 : barrier(barrier), releases(releases), releasesBeforeDrain(releasesBeforeDrain) {}
147struct TerminalResourceProbeReleaser {
148 static void release(TerminalResourceProbe* resource) {
149 *resource->releases += 1;
150 if (!resource->barrier->isClosedAndDrained()) {
151 *resource->releasesBeforeDrain += 1;
158struct TerminalUnwindContext {
159 explicit TerminalUnwindContext(TerminalResourceOwner&& resource)
161 resource(pedigree_std::move(resource)),
168 TerminalResourceOwner resource;
175class TerminalUnwindCanary {
177 explicit TerminalUnwindCanary(
Atomic<size_t>& destructed) : m_Destructed(destructed) {}
179 ~TerminalUnwindCanary() {
187int waitForTerminalRequest(
void* context) {
188 TerminalUnwindContext* terminal =
reinterpret_cast<TerminalUnwindContext*
>(context);
190 if (!terminal->resource->barrier->tryAcquire(admission)) {
191 FATAL(
"QEMU terminal request could not admit its worker");
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;
201int destroyConsumer(
void* context) {
202 ConsumerDestroyContext* destroy =
reinterpret_cast<ConsumerDestroyContext*
>(context);
203 destroy->consumer->shutdown();
204 destroy->finished += 1;
208int publishHandoff(
void* context) {
209 PublishContext* publication =
reinterpret_cast<PublishContext*
>(context);
210 publication->result =
211 publication->queue->republishPreallocatedWhileReleasing(publication->queue->token, 0, 2);
215struct ReciprocalSyscallContext;
219 ReciprocalSyscallHandler(ReciprocalSyscallContext& context,
bool first)
220 : m_Context(context), m_First(first) {}
222 uintptr_t
syscall(SyscallState&)
override;
225 ReciprocalSyscallContext& m_Context;
229struct ReciprocalSyscallContext {
230 ReciprocalSyscallContext()
231 : first(*this, true),
232 second(*this, false),
237 firstProcessor(static_cast<size_t>(-1)),
238 secondProcessor(static_cast<size_t>(-1)) {}
240 ReciprocalSyscallHandler first;
241 ReciprocalSyscallHandler second;
252uintptr_t ReciprocalSyscallHandler::syscall(SyscallState&) {
258 m_Context.entered += 1;
259 for (
size_t attempt = 0; attempt < 65536 && m_Context.entered !=
static_cast<size_t>(2);
263 if (m_Context.entered !=
static_cast<size_t>(2)) {
264 m_Context.failures += 1;
269 m_First ? m_Context.secondRegistration : m_Context.firstRegistration;
270 if (!peer.reset() && peer) {
271 m_Context.rejections += 1;
273 m_Context.failures += 1;
275 m_Context.resetReturns += 1;
276 for (
size_t attempt = 0; attempt < 65536 && m_Context.resetReturns !=
static_cast<size_t>(2);
280 if (m_Context.resetReturns !=
static_cast<size_t>(2)) {
281 m_Context.failures += 1;
283 return m_First ? 0x71 : 0x72;
286struct ReciprocalSyscallInvocation {
291int invokeReciprocalSyscall(
void* context) {
292 ReciprocalSyscallInvocation* invocation =
reinterpret_cast<ReciprocalSyscallInvocation*
>(context);
298void testSyscallReciprocalUnregister() {
299 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN syscall-reciprocal-unregister-smp");
302 ReciprocalSyscallContext context;
305 if (context.firstRegistration) {
306 context.firstRegistration.reset();
308 FATAL(
"QEMU syscall reciprocal handlers could not register");
311 ReciprocalSyscallInvocation first = {TUI, 0};
312 ReciprocalSyscallInvocation second = {native, 0};
315 new Thread(process, invokeReciprocalSyscall, &first,
nullptr,
false,
false,
true);
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");
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");
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");
343void testProducerConsumerTeardown() {
344 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN producerconsumer-teardown-completion");
346 BlockingConsumer consumer;
347 if (!consumer.start()) {
348 FATAL(
"QEMU ProducerConsumer worker did not start");
351 if (!consumer.entered.acquireForCompletion()) {
352 FATAL(
"QEMU ProducerConsumer callback did not enter");
355 ConsumerDestroyContext destroy(&consumer);
357 &destroy,
nullptr,
false,
true);
358 destroyer->setName(
"QEMU ProducerConsumer destroyer");
360 bool joinPublished =
false;
361 for (
size_t attempt = 0; attempt < 4096; ++attempt) {
362 uintptr_t address = 0;
363 if (destroyer->
getStatus() == Thread::Sleeping &&
365 joinPublished =
true;
370 if (!joinPublished) {
371 FATAL(
"QEMU ProducerConsumer destroyer did not join its worker");
375 bool teardownAbandoned =
false;
376 for (
size_t attempt = 0; attempt < 4096; ++attempt) {
377 if (destroyer->
getStatus() == Thread::AwaitingJoin) {
378 teardownAbandoned =
true;
383 if (teardownAbandoned) {
384 FATAL(
"QEMU ProducerConsumer teardown abandoned its worker join");
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");
393 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS producerconsumer-teardown-completion");
396void testTerminalRequestStackUnwind() {
397 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN terminal-request-stack-unwind");
402 TerminalResourceProbe resource(&barrier, &releases, &releasesBeforeDrain);
403 TerminalUnwindContext context(TerminalResourceOwner::adopt(&resource));
405 FATAL(
"QEMU terminal request could not admit its drain probe");
408 &context,
nullptr,
false,
true);
409 worker->setName(
"QEMU terminal request stack unwind");
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;
424 FATAL(
"QEMU closed operation barrier admitted late work");
427 if (barrier.isClosedAndDrained()) {
428 FATAL(
"QEMU operation barrier drained before its worker finished");
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");
444 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS terminal-request-stack-unwind");
447void testPinnedLinkerUnloadRejection() {
448 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN linker-pinned-unload-rejection");
450 static char linkerName[] =
"linker";
454 FATAL(
"QEMU linker module did not reject unload while preserving its image");
457 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS linker-pinned-unload-rejection");
460void testFilesystemModuleUnloadRejection() {
463 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN posix-runtime-unload-rejection");
464 static char posixName[] =
"posix";
467 FATAL(
"QEMU POSIX module did not reject runtime unload while preserving its image");
469 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS posix-runtime-unload-rejection");
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");
479 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS rawfs-runtime-unload-rejection");
481 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN ramfs-pinned-unload-rejection");
482 static char ramfsName[] =
"ramfs";
485 FATAL(
"QEMU RamFs module did not reject unload while preserving its image");
487 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS ramfs-pinned-unload-rejection");
491 if (!runSyscallLifetimeRegression()) {
492 FATAL(
"QEMU syscall handler lifetime regression failed");
494 if (!runRcuConcurrencyRegression()) {
495 FATAL(
"QEMU RCU publication and reclamation regression failed");
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");
504 testProducerConsumerTeardown();
505 testTerminalRequestStackUnwind();
507 if (!runAnonymousMemoryRegionRegression()) {
508 FATAL(
"QEMU anonymous MemoryRegion regression failed");
510 if (!runTlbShootdownConcurrencyRegression()) {
511 FATAL(
"QEMU TLB shootdown regression failed");
513 if (!runSlamAllocatorConcurrencyRegression()) {
514 FATAL(
"QEMU SLAM allocator concurrency regression failed");
516 if (!runTextIoFlipLifetimeRegression()) {
517 FATAL(
"QEMU TextIO flip-worker lifetime regression failed");
519 if (!runVfsCallbackLifetimeRegressions()) {
520 FATAL(
"QEMU VFS callback lifetime regression failed");
523 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN requestqueue-release-handoff");
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");
535 queue.setAfterPreallocatedClaimHookForTest(pauseClaimedPublication, &pause);
537 PublishContext publication(&queue);
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");
545 queue.allowReleaseReturn.release();
547 if (!queue.addAsyncRequest(0, 3)) {
548 FATAL(
"QEMU RequestQueue handoff progress probe was rejected");
551 bool independentProgress =
false;
552 for (
size_t attempt = 0; attempt < 4096; ++attempt) {
553 if (queue.executions.value() >= 2) {
554 independentProgress =
true;
560 pause.release.release();
563 const bool drained = queue.drain();
564 queue.setAfterPreallocatedClaimHookForTest(
nullptr,
nullptr);
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");
574 NOTICE(
"QEMU-CONCURRENCY-TEST: PASS requestqueue-release-handoff");
575 NOTICE(
"QEMU-CONCURRENCY-TEST: BEGIN shutdown-worker-drain-smp");
582MODULE_INFO(
"concurrency-smoke", &entry, &exit,
"console",
"vfs",
"rawfs",
"network-stack");
bool unloadModule(const char *name, bool silent=false, bool progress=true)
bool moduleIsLoaded(char *name)
static KernelElf & instance()
MUST_USE_RESULT bool tryEnter()
static Scheduler & instance()
virtual uintptr_t syscall(SyscallState &State)=0
virtual bool registerSyscallHandler(Service_t Service, SyscallHandler *pHandler, Registration ®istration, FastEntry entry=nullptr)=0
static EXPORTED_PUBLIC SyscallManager & instance()
void setUnwindState(UnwindType ut)
@ TerminateThread
Exit only this thread during Process exit.
DebugState getDebugState(uintptr_t &address)