The Pedigree Project 0.1
unix-stream-interruption-regressions.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/errors.h"
11#include "pedigree/kernel/process/Process.h"
12#include "pedigree/kernel/process/Scheduler.h"
13#include "pedigree/kernel/process/SignalEvent.h"
14#include "pedigree/kernel/process/Thread.h"
15#include "pedigree/kernel/processor/PhysicalMemoryManager.h"
16#include "pedigree/kernel/processor/Processor.h"
17#include "pedigree/kernel/time/Time.h"
18#include "pedigree/kernel/utilities/lib.h"
19
20#include <fcntl.h>
21#include <signal.h>
22
23#include "modules/subsys/posix/FileDescriptor.h"
24#include "modules/subsys/posix/PosixProcess.h"
25#include "modules/subsys/posix/PosixSubsystem.h"
26#include "modules/subsys/posix/UnixFilesystem.h"
27#include "modules/subsys/posix/file-syscalls.h"
28#include "modules/subsys/posix/net-syscalls.h"
30#include "modules/system/vfs/VFS.h"
31#include <sys/socket.h>
32
33namespace {
34constexpr size_t ReturnTimeout = 2 * Time::Multiplier::Second;
35constexpr size_t SourceDescriptor = 90;
36
37Atomic<size_t> g_SignalCalls(0);
38Atomic<int> g_SignalCloseDescriptor(-1);
39Atomic<int> g_SignalCloseResult(-2);
40Atomic<int> g_SecondSignalCloseDescriptor(-1);
41Atomic<int> g_SecondSignalCloseResult(-2);
42Atomic<size_t> g_LateCloneRejected(0);
43FileDescriptor* g_LateCloneSource = nullptr;
44
45void streamSignalHandler(size_t) {
46 g_SignalCalls += 1;
47 const int descriptor = g_SignalCloseDescriptor;
48 if (descriptor >= 0 && g_SignalCloseDescriptor.compareAndSwap(descriptor, -1)) {
49 g_SignalCloseResult = posix_close(descriptor);
50 }
51}
52
53void closeFromStreamControlLock() {
54 streamSignalHandler(SIGUSR1);
55}
56
57void rejectLateCloneAfterHandlerClose() {
58 if (!g_LateCloneSource) {
59 return;
60 }
61
62 FileDescriptor* lateClone = new FileDescriptor(*g_LateCloneSource);
63 if (lateClone->networkImpl && !lateClone->networkPublished()) {
64 g_LateCloneRejected += 1;
65 }
66 delete lateClone;
67}
68
69void closeFromEndpointLockOrReadinessLease() {
70 streamSignalHandler(SIGUSR1);
71 rejectLateCloneAfterHandlerClose();
72}
73
74void closePairFromEndpointMutationLocks() {
75 closeFromEndpointLockOrReadinessLease();
76 const int descriptor = g_SecondSignalCloseDescriptor;
77 if (descriptor >= 0 && g_SecondSignalCloseDescriptor.compareAndSwap(descriptor, -1)) {
78 g_SecondSignalCloseResult = posix_close(descriptor);
79 }
80}
81
82union ControlBuffer {
83 struct cmsghdr alignment;
84 uint8_t bytes[CMSG_SPACE(sizeof(int))];
85};
86
87struct StreamFixture {
88 int sockets[2];
89 uint8_t payload[1024];
90 struct iovec vector;
91 struct msghdr message;
92 ControlBuffer control;
93};
94
95enum class IoOperation {
96 Receive,
97 Send,
98};
99
100struct IoContext {
101 IoContext(int descriptor, uint8_t* payload, size_t length, IoOperation operation)
102 : descriptor(descriptor),
103 payload(payload),
104 length(length),
105 operation(operation),
106 entered(0),
107 returned(0),
108 result(-2),
109 error(0) {}
110
111 int descriptor;
112 uint8_t* payload;
113 size_t length;
114 IoOperation operation;
115 Atomic<size_t> entered;
116 Atomic<size_t> returned;
117 ssize_t result;
118 int error;
119};
120
121int ioWorker(void* parameter) {
122 IoContext* context = reinterpret_cast<IoContext*>(parameter);
123 Thread* thread = Processor::information().getCurrentThread();
124 thread->setErrno(0);
125 thread->clearInterruption();
126 context->entered += 1;
127 if (context->operation == IoOperation::Receive) {
128 context->result = posix_recv(context->descriptor, context->payload, context->length, 0);
129 } else {
130 context->result = posix_send(context->descriptor, context->payload, context->length, 0);
131 }
132 context->error = thread->getErrno();
133 context->returned += 1;
134 return 0;
135}
136
137bool waitUntilBlocked(Thread* thread, IoContext& context,
138 Thread::DebugState expectedState = Thread::CondWait) {
139 const Time::Timestamp deadline = Time::getTicks() + ReturnTimeout;
140 while (Time::getTicks() < deadline) {
141 Thread::WaitDebugInfo wait = {};
142 uintptr_t address = 0;
143 if (context.entered == 1 && !context.returned && thread->getWaitDebugInfo(wait) && wait.queue &&
144 wait.queued && thread->getDebugState(address) == expectedState) {
145 return true;
146 }
147 if (thread->getStatus() == Thread::AwaitingJoin) {
148 return false;
149 }
151 }
152 return false;
153}
154
155bool waitUntilReturned(IoContext& context) {
156 const Time::Timestamp deadline = Time::getTicks() + ReturnTimeout;
157 while (!context.returned && Time::getTicks() < deadline) {
159 }
160 return context.returned == 1;
161}
162
163bool sendCaughtSignal(Thread* thread) {
164 SignalEvent* signal = new SignalEvent(reinterpret_cast<uintptr_t>(&streamSignalHandler), SIGUSR1,
165 ~0UL, 0, true, true);
166 if (thread->sendEvent(signal)) {
167 return true;
168 }
169 delete signal;
170 return false;
171}
172
173bool closeSocket(int& descriptor) {
174 if (descriptor < 0) {
175 return true;
176 }
177 const bool closed = posix_close(descriptor) == 0;
178 descriptor = -1;
179 return closed;
180}
181
182bool closePair(StreamFixture& fixture) {
183 const bool first = closeSocket(fixture.sockets[0]);
184 const bool second = closeSocket(fixture.sockets[1]);
185 return first && second;
186}
187
188bool createPair(StreamFixture& fixture) {
189 fixture.sockets[0] = fixture.sockets[1] = -1;
190 return posix_socketpair(AF_UNIX, SOCK_STREAM, 0, fixture.sockets) == 0;
191}
192
193FileDescriptor::OpenFileDescriptionLease addSourceDescriptor(PosixSubsystem* subsystem) {
194 FileDescriptor* source = new FileDescriptor;
195 source->fd = SourceDescriptor;
197 subsystem->addFileDescriptor(SourceDescriptor, source);
198 return description;
199}
200
201void prepareControlMessage(StreamFixture& fixture, bool sending) {
202 ByteSet(&fixture.control, 0, sizeof(fixture.control));
203 fixture.vector = {fixture.payload, 1};
204 fixture.message = {};
205 fixture.message.msg_iov = &fixture.vector;
206 fixture.message.msg_iovlen = 1;
207 fixture.message.msg_control = fixture.control.bytes;
208 fixture.message.msg_controllen = sizeof(fixture.control.bytes);
209 if (sending) {
210 struct cmsghdr* header = CMSG_FIRSTHDR(&fixture.message);
211 header->cmsg_len = CMSG_LEN(sizeof(int));
212 header->cmsg_level = SOL_SOCKET;
213 header->cmsg_type = SCM_RIGHTS;
214 const int source = SourceDescriptor;
215 MemoryCopy(CMSG_DATA(header), &source, sizeof(source));
216 }
217}
218
219bool signalHandlerCloseWhileControlLocked(PosixSubsystem* subsystem, StreamFixture& fixture,
220 IoOperation operation) {
221 if (!createPair(fixture)) {
222 return false;
223 }
224
225 FileDescriptor::OpenFileDescriptionLease source = addSourceDescriptor(subsystem);
226 fixture.payload[0] = 'c';
227 prepareControlMessage(fixture, true);
228 bool passed = true;
229 if (operation == IoOperation::Receive) {
230 if (posix_sendmsg(fixture.sockets[0], &fixture.message, 0) != 1 ||
231 posix_close(SourceDescriptor) != 0 || source->descriptorOwnerCount() != 1 ||
232 SocketRights::inFlightForTest() != 1) {
233 posix_close(SourceDescriptor);
234 closePair(fixture);
235 return false;
236 }
237 prepareControlMessage(fixture, false);
238 }
239
240 const size_t operationSide = operation == IoOperation::Receive ? 1 : 0;
241 const size_t closeSide = 1;
242 g_SignalCalls = 0;
243 g_SignalCloseResult = -2;
244 g_SignalCloseDescriptor = fixture.sockets[closeSide];
245 setUnixStreamControlLockHookForTest(closeFromStreamControlLock);
246 Thread* thread = Processor::information().getCurrentThread();
247 thread->setErrno(0);
248 const ssize_t result = operation == IoOperation::Receive
249 ? posix_recvmsg(fixture.sockets[operationSide], &fixture.message, 0)
250 : posix_sendmsg(fixture.sockets[operationSide], &fixture.message, 0);
251 const int error = thread->getErrno();
252 setUnixStreamControlLockHookForTest(nullptr);
253
254 const bool handlerClosed = g_SignalCloseResult == 0;
255 if (handlerClosed) {
256 fixture.sockets[closeSide] = -1;
257 }
258 g_SignalCloseDescriptor = -1;
259
260 if (operation == IoOperation::Send) {
261 passed = result == -1 && error == Error::BrokenPipe && posix_close(SourceDescriptor) == 0 &&
262 source->descriptorOwnerCount() == 0 && passed;
263 } else {
264 passed = result == 0 && source->descriptorOwnerCount() == 0 && passed;
265 }
266 passed = handlerClosed && g_SignalCalls == 1 && SocketRights::inFlightForTest() == 0 && passed;
267 return closePair(fixture) && passed;
268}
269
270bool lateCloneIsRejected(StreamFixture& fixture) {
271 if (!createPair(fixture)) {
272 return false;
273 }
274
275 DescriptorLease retained;
276 if (!acquireDescriptor(fixture.sockets[0], retained)) {
277 closePair(fixture);
278 return false;
279 }
280
281 const bool closed = closeSocket(fixture.sockets[0]);
282 FileDescriptor* lateClone = closed ? new FileDescriptor(*retained) : nullptr;
283 const bool rejected = lateClone && retained->networkImpl && !lateClone->networkPublished();
284 delete lateClone;
285 retained.reset();
286 return rejected && closePair(fixture);
287}
288
289bool handlerCloseWhileEndpointMutationLocked(StreamFixture& fixture) {
290 if (!createPair(fixture)) {
291 return false;
292 }
293
294 DescriptorLease retained;
295 if (!acquireDescriptor(fixture.sockets[0], retained)) {
296 closePair(fixture);
297 return false;
298 }
300
301 g_SignalCalls = 0;
302 g_SignalCloseResult = -2;
303 g_LateCloneRejected = 0;
304 g_LateCloneSource = &*retained;
305 g_SignalCloseDescriptor = fixture.sockets[0];
306 setUnixEndpointMutationLockHookForTest(closeFromEndpointLockOrReadinessLease);
307 Thread* thread = Processor::information().getCurrentThread();
308 thread->setErrno(0);
309 const bool created = network->create();
310 const int error = thread->getErrno();
311 setUnixEndpointMutationLockHookForTest(nullptr);
312
313 const bool handlerClosed = g_SignalCloseResult == 0;
314 if (handlerClosed) {
315 fixture.sockets[0] = -1;
316 }
317 g_SignalCloseDescriptor = -1;
318 g_LateCloneSource = nullptr;
319
320 const ReadyMask ready = network->queryReady(true, true);
321 const ssize_t peerResult = posix_recv(fixture.sockets[1], fixture.payload, 1, 0);
322 const bool passed = !created && error == Error::BadFileDescriptor && handlerClosed &&
323 g_SignalCalls == 1 && g_LateCloneRejected == 1 &&
324 (ready & (ReadyInvalid | ReadyHangup)) == (ReadyInvalid | ReadyHangup) &&
325 peerResult == 0;
326 return closePair(fixture) && passed;
327}
328
329bool handlerCloseWhileEndpointReadinessActive(StreamFixture& fixture) {
330 if (!createPair(fixture)) {
331 return false;
332 }
333
334 DescriptorLease retained;
335 if (!acquireDescriptor(fixture.sockets[0], retained)) {
336 closePair(fixture);
337 return false;
338 }
340
341 g_SignalCalls = 0;
342 g_SignalCloseResult = -2;
343 g_LateCloneRejected = 0;
344 g_LateCloneSource = &*retained;
345 g_SignalCloseDescriptor = fixture.sockets[0];
346 setUnixEndpointReadinessLeaseHookForTest(closeFromEndpointLockOrReadinessLease);
347 const ReadyMask ready = network->queryReady(true, true);
348 setUnixEndpointReadinessLeaseHookForTest(nullptr);
349
350 const bool handlerClosed = g_SignalCloseResult == 0;
351 if (handlerClosed) {
352 fixture.sockets[0] = -1;
353 }
354 g_SignalCloseDescriptor = -1;
355 g_LateCloneSource = nullptr;
356
357 const ssize_t peerResult = posix_recv(fixture.sockets[1], fixture.payload, 1, 0);
358 const bool passed = handlerClosed && g_SignalCalls == 1 && g_LateCloneRejected == 1 &&
359 (ready & (ReadyInvalid | ReadyHangup)) == (ReadyInvalid | ReadyHangup) &&
360 peerResult == 0;
361 return closePair(fixture) && passed;
362}
363
364bool handlerClosesBothEndpointsWhilePairLocksOwned(StreamFixture& fixture) {
365 fixture.sockets[0] = posix_socket(AF_UNIX, SOCK_STREAM, 0);
366 fixture.sockets[1] = posix_socket(AF_UNIX, SOCK_STREAM, 0);
367 if (fixture.sockets[0] < 0 || fixture.sockets[1] < 0) {
368 closePair(fixture);
369 return false;
370 }
371
372 DescriptorLease firstDescriptor;
373 DescriptorLease secondDescriptor;
374 if (!acquireDescriptor(fixture.sockets[0], firstDescriptor) ||
375 !acquireDescriptor(fixture.sockets[1], secondDescriptor)) {
376 closePair(fixture);
377 return false;
378 }
379 UnixSocketSyscalls* first = static_cast<UnixSocketSyscalls*>(firstDescriptor->networkImpl.get());
380 UnixSocketSyscalls* second =
381 static_cast<UnixSocketSyscalls*>(secondDescriptor->networkImpl.get());
382
383 g_SignalCalls = 0;
384 g_SignalCloseResult = -2;
385 g_SecondSignalCloseResult = -2;
386 g_LateCloneRejected = 0;
387 g_LateCloneSource = &*firstDescriptor;
388 g_SignalCloseDescriptor = fixture.sockets[0];
389 g_SecondSignalCloseDescriptor = fixture.sockets[1];
390 setUnixEndpointMutationLockHookForTest(closePairFromEndpointMutationLocks);
391 const bool paired = first->pairWith(second);
392 setUnixEndpointMutationLockHookForTest(nullptr);
393
394 const bool firstClosed = g_SignalCloseResult == 0;
395 const bool secondClosed = g_SecondSignalCloseResult == 0;
396 if (firstClosed) {
397 fixture.sockets[0] = -1;
398 }
399 if (secondClosed) {
400 fixture.sockets[1] = -1;
401 }
402 g_SignalCloseDescriptor = -1;
403 g_SecondSignalCloseDescriptor = -1;
404 g_LateCloneSource = nullptr;
405
406 const ReadyMask firstReady = first->queryReady(true, true);
407 const ReadyMask secondReady = second->queryReady(true, true);
408 const bool passed = !paired && firstClosed && secondClosed && g_SignalCalls == 1 &&
409 g_LateCloneRejected == 1 &&
410 (firstReady & (ReadyInvalid | ReadyHangup)) == (ReadyInvalid | ReadyHangup) &&
411 (secondReady & (ReadyInvalid | ReadyHangup)) == (ReadyInvalid | ReadyHangup);
412 return closePair(fixture) && passed;
413}
414
415bool setNonblocking(int descriptor, bool nonblocking) {
416 const int flags = posix_fcntl(descriptor, F_GETFL, nullptr);
417 if (flags < 0) {
418 return false;
419 }
420 const int replacement = nonblocking ? flags | O_NONBLOCK : flags & ~O_NONBLOCK;
421 return posix_fcntl(descriptor, F_SETFL, reinterpret_cast<void*>(replacement)) == 0;
422}
423
424bool fillSendQueue(StreamFixture& fixture) {
425 if (!setNonblocking(fixture.sockets[0], true)) {
426 return false;
427 }
428
429 ByteSet(fixture.payload, 0x51, sizeof(fixture.payload));
430 size_t written = 0;
431 while (true) {
432 Thread* thread = Processor::information().getCurrentThread();
433 thread->setErrno(0);
434 const ssize_t result =
435 posix_send(fixture.sockets[0], fixture.payload, sizeof(fixture.payload), 0);
436 if (result > 0) {
437 written += static_cast<size_t>(result);
438 continue;
439 }
440 if (result != -1 || thread->getErrno() != Error::NoMoreProcesses) {
441 return false;
442 }
443 break;
444 }
445
446 return written == MAX_UNIX_STREAM_QUEUE && setNonblocking(fixture.sockets[0], false);
447}
448
449bool signalInterruptsEmptyReceive(Process* process, StreamFixture& fixture) {
450 if (!createPair(fixture)) {
451 return false;
452 }
453
454 g_SignalCalls = 0;
455 IoContext context(fixture.sockets[1], fixture.payload, 1, IoOperation::Receive);
456 Thread* worker = new Thread(process, ioWorker, &context, nullptr, false, true, true);
457 worker->setName("hosted AF_UNIX interrupted receiver");
458 const bool started = worker->start();
459 const bool blocked = started && waitUntilBlocked(worker, context);
460 const bool signalled = blocked && sendCaughtSignal(worker);
461 const bool returned = signalled && waitUntilReturned(context);
462
463 if (started && !returned) {
464 closePair(fixture);
465 if (!context.returned) {
466 sendCaughtSignal(worker);
467 waitUntilReturned(context);
468 }
469 }
470 const bool joined = started && worker->joinForCompletion();
471 if (!started) {
472 delete worker;
473 }
474 const bool closed = closePair(fixture);
475 return started && blocked && signalled && returned && joined && context.result == -1 &&
476 context.error == Error::Interrupted && g_SignalCalls == 1 && closed;
477}
478
479bool signalInterruptsFullSend(Process* process, StreamFixture& fixture) {
480 if (!createPair(fixture) || !fillSendQueue(fixture)) {
481 closePair(fixture);
482 return false;
483 }
484
485 g_SignalCalls = 0;
486 IoContext context(fixture.sockets[0], fixture.payload, 1, IoOperation::Send);
487 Thread* worker = new Thread(process, ioWorker, &context, nullptr, false, true, true);
488 worker->setName("hosted AF_UNIX interrupted sender");
489 const bool started = worker->start();
490 const bool blocked = started && waitUntilBlocked(worker, context);
491 const bool signalled = blocked && sendCaughtSignal(worker);
492 const bool returned = signalled && waitUntilReturned(context);
493
494 if (started && !returned) {
495 closePair(fixture);
496 if (!context.returned) {
497 sendCaughtSignal(worker);
498 waitUntilReturned(context);
499 }
500 }
501 const bool joined = started && worker->joinForCompletion();
502 if (!started) {
503 delete worker;
504 }
505 const bool closed = closePair(fixture);
506 return started && blocked && signalled && returned && joined && context.result == -1 &&
507 context.error == Error::Interrupted && g_SignalCalls == 1 && closed;
508}
509
510bool signalInterruptsSerializationWait(Process* process, StreamFixture& fixture,
511 IoOperation operation) {
512 if (!createPair(fixture)) {
513 return false;
514 }
515 if (operation == IoOperation::Send && !fillSendQueue(fixture)) {
516 closePair(fixture);
517 return false;
518 }
519
520 const size_t operationSide = operation == IoOperation::Receive ? 1 : 0;
521 g_SignalCalls = 0;
522 IoContext holderContext(fixture.sockets[operationSide], fixture.payload, 1, operation);
523 IoContext waiterContext(fixture.sockets[operationSide], fixture.payload, 1, operation);
524 Thread* holder = new Thread(process, ioWorker, &holderContext, nullptr, false, true, true);
525 Thread* waiter = new Thread(process, ioWorker, &waiterContext, nullptr, false, true, true);
526 holder->setName("hosted AF_UNIX serialization holder");
527 waiter->setName("hosted AF_UNIX serialization waiter");
528
529 const bool holderStarted = holder->start();
530 const bool holderBlocked = holderStarted && waitUntilBlocked(holder, holderContext);
531 const bool waiterStarted = holderBlocked && waiter->start();
532 const bool waiterBlocked =
533 waiterStarted && waitUntilBlocked(waiter, waiterContext, Thread::SemWait);
534 const bool signalled = waiterBlocked && sendCaughtSignal(waiter);
535 const bool waiterReturned = signalled && waitUntilReturned(waiterContext);
536 const bool holderStayedBlocked = !holderContext.returned;
537
538 const bool closed = closePair(fixture);
539 const bool holderReturned = holderStarted && waitUntilReturned(holderContext);
540 if (waiterStarted && !waiterContext.returned) {
541 sendCaughtSignal(waiter);
542 waitUntilReturned(waiterContext);
543 }
544 if (holderStarted && !holderContext.returned) {
545 sendCaughtSignal(holder);
546 waitUntilReturned(holderContext);
547 }
548
549 const bool waiterJoined = waiterStarted && waiter->joinForCompletion();
550 const bool holderJoined = holderStarted && holder->joinForCompletion();
551 if (!waiterStarted) {
552 delete waiter;
553 }
554 if (!holderStarted) {
555 delete holder;
556 }
557
558 const bool holderResult =
559 operation == IoOperation::Receive
560 ? holderContext.result == 0
561 : holderContext.result == -1 && holderContext.error == Error::BrokenPipe;
562 return holderStarted && holderBlocked && waiterStarted && waiterBlocked && signalled &&
563 waiterReturned && holderStayedBlocked && closed && holderReturned && waiterJoined &&
564 holderJoined && waiterContext.result == -1 && waiterContext.error == Error::Interrupted &&
565 holderResult && g_SignalCalls == 1;
566}
567
568bool partialSendWinsSignal(Process* process, StreamFixture& fixture) {
569 if (!createPair(fixture) || !fillSendQueue(fixture)) {
570 closePair(fixture);
571 return false;
572 }
573
574 if (posix_recv(fixture.sockets[1], fixture.payload, 1, 0) != 1) {
575 closePair(fixture);
576 return false;
577 }
578 fixture.payload[0] = 'x';
579 fixture.payload[1] = 'y';
580
581 g_SignalCalls = 0;
582 IoContext context(fixture.sockets[0], fixture.payload, 2, IoOperation::Send);
583 Thread* worker = new Thread(process, ioWorker, &context, nullptr, false, true, true);
584 worker->setName("hosted AF_UNIX partial interrupted sender");
585 const bool started = worker->start();
586 const bool blocked = started && waitUntilBlocked(worker, context);
587 const bool signalled = blocked && sendCaughtSignal(worker);
588 const bool returned = signalled && waitUntilReturned(context);
589
590 if (started && !returned) {
591 closePair(fixture);
592 if (!context.returned) {
593 sendCaughtSignal(worker);
594 waitUntilReturned(context);
595 }
596 }
597 const bool joined = started && worker->joinForCompletion();
598 if (!started) {
599 delete worker;
600 }
601 const bool closed = closePair(fixture);
602 return started && blocked && signalled && returned && joined && context.result == 1 &&
603 context.error == 0 && g_SignalCalls == 1 && closed;
604}
605
606bool closeWakesBlockedIo(Process* process, StreamFixture& fixture, IoOperation operation,
607 bool closeLocal) {
608 if (!createPair(fixture)) {
609 return false;
610 }
611 if (operation == IoOperation::Send && !fillSendQueue(fixture)) {
612 closePair(fixture);
613 return false;
614 }
615
616 const size_t operationSide = operation == IoOperation::Receive ? 1 : 0;
617 const size_t closeSide = closeLocal ? operationSide : 1 - operationSide;
618 IoContext context(fixture.sockets[operationSide], fixture.payload, 1, operation);
619 Thread* worker = new Thread(process, ioWorker, &context, nullptr, false, true, true);
620 if (operation == IoOperation::Receive) {
621 worker->setName("hosted AF_UNIX close-woken receiver");
622 } else {
623 worker->setName("hosted AF_UNIX close-woken sender");
624 }
625 const bool started = worker->start();
626 const bool blocked = started && waitUntilBlocked(worker, context);
627 const bool triggeringClose = blocked && closeSocket(fixture.sockets[closeSide]);
628 const bool returned = triggeringClose && waitUntilReturned(context);
629
630 if (started && !returned) {
631 closePair(fixture);
632 if (!context.returned) {
633 sendCaughtSignal(worker);
634 waitUntilReturned(context);
635 }
636 }
637 const bool joined = started && worker->joinForCompletion();
638 if (!started) {
639 delete worker;
640 }
641 const bool closed = closePair(fixture);
642 const bool result = operation == IoOperation::Receive
643 ? context.result == 0
644 : context.result == -1 && context.error == Error::BrokenPipe;
645 return started && blocked && triggeringClose && returned && joined && result && closed;
646}
647
648struct SuiteContext {
649 explicit SuiteContext(Process* process) : process(process), completed(0), passed(false) {}
650
651 Process* process;
652 Atomic<size_t> completed;
653 bool passed;
654};
655
656int suiteWorker(void* parameter) {
657 SuiteContext* context = reinterpret_cast<SuiteContext*>(parameter);
658 const size_t pageSize = PhysicalMemoryManager::getPageSize();
659 uintptr_t mappingAddress = 0;
660 if (!context->process->allocateUserRange(Process::UserRegion::Normal, pageSize, mappingAddress)) {
661 context->completed += 1;
662 return 1;
663 }
664
665 uintptr_t mappedAddress = mappingAddress;
667 mappedAddress, pageSize, MemoryMappedObject::Read | MemoryMappedObject::Write);
668 if (!mapping || mappedAddress != mappingAddress || sizeof(StreamFixture) > pageSize) {
669 if (mapping) {
670 MemoryMapManager::instance().remove(mappedAddress, pageSize);
671 }
672 context->process->freeUserRange(Process::UserRegion::Normal, mappingAddress, pageSize);
673 context->completed += 1;
674 return 1;
675 }
676
677 StreamFixture* fixture = reinterpret_cast<StreamFixture*>(mappedAddress);
678 ByteSet(fixture, 0, sizeof(*fixture));
679 PosixSubsystem* subsystem = static_cast<PosixSubsystem*>(context->process->getSubsystem());
680 bool passed = lateCloneIsRejected(*fixture);
681 passed = handlerCloseWhileEndpointMutationLocked(*fixture) && passed;
682 passed = handlerCloseWhileEndpointReadinessActive(*fixture) && passed;
683 passed = handlerClosesBothEndpointsWhilePairLocksOwned(*fixture) && passed;
684 passed =
685 signalHandlerCloseWhileControlLocked(subsystem, *fixture, IoOperation::Receive) && passed;
686 passed = signalHandlerCloseWhileControlLocked(subsystem, *fixture, IoOperation::Send) && passed;
687 passed = signalInterruptsEmptyReceive(context->process, *fixture) && passed;
688 passed = signalInterruptsFullSend(context->process, *fixture) && passed;
689 passed =
690 signalInterruptsSerializationWait(context->process, *fixture, IoOperation::Receive) && passed;
691 passed =
692 signalInterruptsSerializationWait(context->process, *fixture, IoOperation::Send) && passed;
693 passed = partialSendWinsSignal(context->process, *fixture) && passed;
694 passed = closeWakesBlockedIo(context->process, *fixture, IoOperation::Receive, false) && passed;
695 passed = closeWakesBlockedIo(context->process, *fixture, IoOperation::Receive, true) && passed;
696 passed = closeWakesBlockedIo(context->process, *fixture, IoOperation::Send, false) && passed;
697 passed = closeWakesBlockedIo(context->process, *fixture, IoOperation::Send, true) && passed;
698
699 MemoryMapManager::instance().remove(mappedAddress, pageSize);
700 context->process->freeUserRange(Process::UserRegion::Normal, mappingAddress, pageSize);
701 context->passed = passed;
702 context->completed += 1;
703 return passed ? 0 : 1;
704}
705} // namespace
706
707bool runHostedUnixStreamInterruptionRegressions(Process* kernelProcess) {
709 auto* priorView = VFS::instance().mountView();
710 VFS::HostedRootViewScope fixture;
711 UnixFilesystem* filesystem = new UnixFilesystem;
712 if (!fixture.open(filesystem)) {
713 delete filesystem;
714 return false;
715 }
716 Process* process =
717 new PosixProcess(kernelProcess, true, Process::FilesystemContextMode::Deferred);
718 process->setSubsystem(new PosixSubsystem);
719 const bool contextInstalled = fixture.installContext(*process);
720 SuiteContext context(process);
721 Thread* worker = new Thread(process, suiteWorker, &context, nullptr, false, true, true);
722 worker->setName("hosted AF_UNIX stream interruption suite");
723 const bool started = contextInstalled && worker->start();
724 const bool joined = started && worker->joinForCompletion();
725 if (!started) {
726 delete worker;
727 }
728
729 bool passed = started && joined && context.completed == 1 && context.passed;
730 delete process;
731
732 const bool rootRestored = fixture.close();
733 if (!rootRestored)
734 FATAL("Hosted filesystem fixture retained owners after teardown");
735 passed = rootRestored && VFS::instance().getRootFilesystem() == priorRoot &&
736 VFS::instance().mountView() == priorView && passed;
737 delete filesystem;
738
739 if (!passed) {
740 ERROR(
741 "HOSTED-SYSCALL-TEST: FAIL unix-stream-interruption: "
742 "EINTR, partial-progress, or close wake semantics regressed");
743 return false;
744 }
745
746 NOTICE("HOSTED-SYSCALL-TEST: PASS unix-stream-interruption");
747 return true;
748}
Memory-mapped file interface.
OpenFileDescriptionLease acquireOpenFileDescription() const
SharedPointer< NetworkSyscalls > networkImpl
Network syscall implementation for this descriptor (if it's a socket).
size_t fd
Descriptor number.
bool networkPublished() const
MemoryMappedObject * mapAnon(uintptr_t &address, size_t length, MemoryMappedObject::Permissions perms)
size_t remove(uintptr_t base, size_t length)
static MemoryMapManager & instance()
void addFileDescriptor(size_t fd, FileDescriptor *pFd)
static ProcessorInformation & information()
static Scheduler & instance()
Definition Scheduler.h:96
void yield()
Definition Scheduler.cc:226
T * get() const
void setErrno(size_t err)
Definition Thread.h:478
DebugState
Definition Thread.h:176
bool getWaitDebugInfo(WaitDebugInfo &info)
Definition Thread.cc:3184
size_t getErrno()
Definition Thread.h:473
bool joinForCompletion()
Definition Thread.cc:2771
Status getStatus() const
Definition Thread.h:431
DebugState getDebugState(uintptr_t &address)
Definition Thread.h:570
bool start()
Definition Thread.cc:794
bool sendEvent(Event *pEvent)
Definition Thread.cc:1158
bool pairWith(UnixSocketSyscalls *other)
virtual ReadyMask queryReady(bool reading, bool writing)
Filesystem * getRootFilesystem() const
Definition VFS.cc:631
static VFS & instance()
Definition VFS.cc:291