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"
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"
34constexpr size_t ReturnTimeout = 2 * Time::Multiplier::Second;
35constexpr size_t SourceDescriptor = 90;
45void streamSignalHandler(
size_t) {
47 const int descriptor = g_SignalCloseDescriptor;
48 if (descriptor >= 0 && g_SignalCloseDescriptor.compareAndSwap(descriptor, -1)) {
49 g_SignalCloseResult = posix_close(descriptor);
53void closeFromStreamControlLock() {
54 streamSignalHandler(SIGUSR1);
57void rejectLateCloneAfterHandlerClose() {
58 if (!g_LateCloneSource) {
64 g_LateCloneRejected += 1;
69void closeFromEndpointLockOrReadinessLease() {
70 streamSignalHandler(SIGUSR1);
71 rejectLateCloneAfterHandlerClose();
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);
83 struct cmsghdr alignment;
84 uint8_t bytes[CMSG_SPACE(
sizeof(
int))];
89 uint8_t payload[1024];
92 ControlBuffer control;
95enum class IoOperation {
101 IoContext(
int descriptor, uint8_t* payload,
size_t length, IoOperation operation)
102 : descriptor(descriptor),
105 operation(operation),
114 IoOperation operation;
121int ioWorker(
void* parameter) {
122 IoContext* context =
reinterpret_cast<IoContext*
>(parameter);
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);
130 context->result = posix_send(context->descriptor, context->payload, context->length, 0);
132 context->error = thread->
getErrno();
133 context->returned += 1;
137bool waitUntilBlocked(
Thread* thread, IoContext& context,
139 const Time::Timestamp deadline = Time::getTicks() + ReturnTimeout;
140 while (Time::getTicks() < deadline) {
142 uintptr_t address = 0;
143 if (context.entered == 1 && !context.returned && thread->
getWaitDebugInfo(wait) && wait.queue &&
144 wait.queued && thread->
getDebugState(address) == expectedState) {
147 if (thread->
getStatus() == Thread::AwaitingJoin) {
155bool waitUntilReturned(IoContext& context) {
156 const Time::Timestamp deadline = Time::getTicks() + ReturnTimeout;
157 while (!context.returned && Time::getTicks() < deadline) {
160 return context.returned == 1;
163bool sendCaughtSignal(
Thread* thread) {
165 ~0UL, 0,
true,
true);
173bool closeSocket(
int& descriptor) {
174 if (descriptor < 0) {
177 const bool closed = posix_close(descriptor) == 0;
182bool closePair(StreamFixture& fixture) {
183 const bool first = closeSocket(fixture.sockets[0]);
184 const bool second = closeSocket(fixture.sockets[1]);
185 return first && second;
188bool createPair(StreamFixture& fixture) {
189 fixture.sockets[0] = fixture.sockets[1] = -1;
190 return posix_socketpair(AF_UNIX, SOCK_STREAM, 0, fixture.sockets) == 0;
195 source->
fd = SourceDescriptor;
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);
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));
219bool signalHandlerCloseWhileControlLocked(
PosixSubsystem* subsystem, StreamFixture& fixture,
220 IoOperation operation) {
221 if (!createPair(fixture)) {
226 fixture.payload[0] =
'c';
227 prepareControlMessage(fixture,
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);
237 prepareControlMessage(fixture,
false);
240 const size_t operationSide = operation == IoOperation::Receive ? 1 : 0;
241 const size_t closeSide = 1;
243 g_SignalCloseResult = -2;
244 g_SignalCloseDescriptor = fixture.sockets[closeSide];
245 setUnixStreamControlLockHookForTest(closeFromStreamControlLock);
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);
254 const bool handlerClosed = g_SignalCloseResult == 0;
256 fixture.sockets[closeSide] = -1;
258 g_SignalCloseDescriptor = -1;
260 if (operation == IoOperation::Send) {
261 passed = result == -1 && error == Error::BrokenPipe && posix_close(SourceDescriptor) == 0 &&
262 source->descriptorOwnerCount() == 0 && passed;
264 passed = result == 0 && source->descriptorOwnerCount() == 0 && passed;
266 passed = handlerClosed && g_SignalCalls == 1 && SocketRights::inFlightForTest() == 0 && passed;
267 return closePair(fixture) && passed;
270bool lateCloneIsRejected(StreamFixture& fixture) {
271 if (!createPair(fixture)) {
276 if (!acquireDescriptor(fixture.sockets[0], retained)) {
281 const bool closed = closeSocket(fixture.sockets[0]);
286 return rejected && closePair(fixture);
289bool handlerCloseWhileEndpointMutationLocked(StreamFixture& fixture) {
290 if (!createPair(fixture)) {
295 if (!acquireDescriptor(fixture.sockets[0], retained)) {
302 g_SignalCloseResult = -2;
303 g_LateCloneRejected = 0;
304 g_LateCloneSource = &*retained;
305 g_SignalCloseDescriptor = fixture.sockets[0];
306 setUnixEndpointMutationLockHookForTest(closeFromEndpointLockOrReadinessLease);
309 const bool created = network->create();
310 const int error = thread->
getErrno();
311 setUnixEndpointMutationLockHookForTest(
nullptr);
313 const bool handlerClosed = g_SignalCloseResult == 0;
315 fixture.sockets[0] = -1;
317 g_SignalCloseDescriptor = -1;
318 g_LateCloneSource =
nullptr;
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) &&
326 return closePair(fixture) && passed;
329bool handlerCloseWhileEndpointReadinessActive(StreamFixture& fixture) {
330 if (!createPair(fixture)) {
335 if (!acquireDescriptor(fixture.sockets[0], retained)) {
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);
350 const bool handlerClosed = g_SignalCloseResult == 0;
352 fixture.sockets[0] = -1;
354 g_SignalCloseDescriptor = -1;
355 g_LateCloneSource =
nullptr;
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) &&
361 return closePair(fixture) && passed;
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) {
374 if (!acquireDescriptor(fixture.sockets[0], firstDescriptor) ||
375 !acquireDescriptor(fixture.sockets[1], secondDescriptor)) {
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);
394 const bool firstClosed = g_SignalCloseResult == 0;
395 const bool secondClosed = g_SecondSignalCloseResult == 0;
397 fixture.sockets[0] = -1;
400 fixture.sockets[1] = -1;
402 g_SignalCloseDescriptor = -1;
403 g_SecondSignalCloseDescriptor = -1;
404 g_LateCloneSource =
nullptr;
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;
415bool setNonblocking(
int descriptor,
bool nonblocking) {
416 const int flags = posix_fcntl(descriptor, F_GETFL,
nullptr);
420 const int replacement = nonblocking ? flags | O_NONBLOCK : flags & ~O_NONBLOCK;
421 return posix_fcntl(descriptor, F_SETFL,
reinterpret_cast<void*
>(replacement)) == 0;
424bool fillSendQueue(StreamFixture& fixture) {
425 if (!setNonblocking(fixture.sockets[0],
true)) {
429 ByteSet(fixture.payload, 0x51,
sizeof(fixture.payload));
434 const ssize_t result =
435 posix_send(fixture.sockets[0], fixture.payload,
sizeof(fixture.payload), 0);
437 written +=
static_cast<size_t>(result);
440 if (result != -1 || thread->
getErrno() != Error::NoMoreProcesses) {
446 return written == MAX_UNIX_STREAM_QUEUE && setNonblocking(fixture.sockets[0],
false);
449bool signalInterruptsEmptyReceive(
Process* process, StreamFixture& fixture) {
450 if (!createPair(fixture)) {
455 IoContext context(fixture.sockets[1], fixture.payload, 1, IoOperation::Receive);
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);
463 if (started && !returned) {
465 if (!context.returned) {
467 waitUntilReturned(context);
470 const bool joined = started &&
worker->joinForCompletion();
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;
479bool signalInterruptsFullSend(
Process* process, StreamFixture& fixture) {
480 if (!createPair(fixture) || !fillSendQueue(fixture)) {
486 IoContext context(fixture.sockets[0], fixture.payload, 1, IoOperation::Send);
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);
494 if (started && !returned) {
496 if (!context.returned) {
498 waitUntilReturned(context);
501 const bool joined = started &&
worker->joinForCompletion();
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;
510bool signalInterruptsSerializationWait(
Process* process, StreamFixture& fixture,
511 IoOperation operation) {
512 if (!createPair(fixture)) {
515 if (operation == IoOperation::Send && !fillSendQueue(fixture)) {
520 const size_t operationSide = operation == IoOperation::Receive ? 1 : 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");
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;
538 const bool closed = closePair(fixture);
539 const bool holderReturned = holderStarted && waitUntilReturned(holderContext);
540 if (waiterStarted && !waiterContext.returned) {
542 waitUntilReturned(waiterContext);
544 if (holderStarted && !holderContext.returned) {
545 sendCaughtSignal(holder);
546 waitUntilReturned(holderContext);
549 const bool waiterJoined = waiterStarted &&
waiter->joinForCompletion();
551 if (!waiterStarted) {
554 if (!holderStarted) {
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;
568bool partialSendWinsSignal(
Process* process, StreamFixture& fixture) {
569 if (!createPair(fixture) || !fillSendQueue(fixture)) {
574 if (posix_recv(fixture.sockets[1], fixture.payload, 1, 0) != 1) {
578 fixture.payload[0] =
'x';
579 fixture.payload[1] =
'y';
582 IoContext context(fixture.sockets[0], fixture.payload, 2, IoOperation::Send);
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);
590 if (started && !returned) {
592 if (!context.returned) {
594 waitUntilReturned(context);
597 const bool joined = started &&
worker->joinForCompletion();
601 const bool closed = closePair(fixture);
602 return started && blocked && signalled && returned && joined && context.result == 1 &&
603 context.error == 0 && g_SignalCalls == 1 && closed;
606bool closeWakesBlockedIo(
Process* process, StreamFixture& fixture, IoOperation operation,
608 if (!createPair(fixture)) {
611 if (operation == IoOperation::Send && !fillSendQueue(fixture)) {
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);
620 if (operation == IoOperation::Receive) {
621 worker->setName(
"hosted AF_UNIX close-woken receiver");
623 worker->setName(
"hosted AF_UNIX close-woken sender");
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);
630 if (started && !returned) {
632 if (!context.returned) {
634 waitUntilReturned(context);
637 const bool joined = started &&
worker->joinForCompletion();
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;
649 explicit SuiteContext(
Process* process) : process(process), completed(0), passed(false) {}
656int suiteWorker(
void* parameter) {
657 SuiteContext* context =
reinterpret_cast<SuiteContext*
>(parameter);
659 uintptr_t mappingAddress = 0;
660 if (!context->process->allocateUserRange(Process::UserRegion::Normal, pageSize, mappingAddress)) {
661 context->completed += 1;
665 uintptr_t mappedAddress = mappingAddress;
667 mappedAddress, pageSize, MemoryMappedObject::Read | MemoryMappedObject::Write);
668 if (!mapping || mappedAddress != mappingAddress ||
sizeof(StreamFixture) > pageSize) {
672 context->process->freeUserRange(Process::UserRegion::Normal, mappingAddress, pageSize);
673 context->completed += 1;
677 StreamFixture* fixture =
reinterpret_cast<StreamFixture*
>(mappedAddress);
678 ByteSet(fixture, 0,
sizeof(*fixture));
680 bool passed = lateCloneIsRejected(*fixture);
681 passed = handlerCloseWhileEndpointMutationLocked(*fixture) && passed;
682 passed = handlerCloseWhileEndpointReadinessActive(*fixture) && passed;
683 passed = handlerClosesBothEndpointsWhilePairLocksOwned(*fixture) && 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;
690 signalInterruptsSerializationWait(context->process, *fixture, IoOperation::Receive) && 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;
700 context->process->freeUserRange(Process::UserRegion::Normal, mappingAddress, pageSize);
701 context->passed = passed;
702 context->completed += 1;
703 return passed ? 0 : 1;
707bool runHostedUnixStreamInterruptionRegressions(
Process* kernelProcess) {
710 VFS::HostedRootViewScope fixture;
712 if (!fixture.open(filesystem)) {
717 new PosixProcess(kernelProcess,
true, Process::FilesystemContextMode::Deferred);
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();
729 bool passed = started && joined && context.completed == 1 && context.passed;
732 const bool rootRestored = fixture.close();
734 FATAL(
"Hosted filesystem fixture retained owners after teardown");
741 "HOSTED-SYSCALL-TEST: FAIL unix-stream-interruption: "
742 "EINTR, partial-progress, or close wake semantics regressed");
746 NOTICE(
"HOSTED-SYSCALL-TEST: PASS unix-stream-interruption");
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()
static constexpr size_t getPageSize() PURE
void addFileDescriptor(size_t fd, FileDescriptor *pFd)
static ProcessorInformation & information()
static Scheduler & instance()
void setErrno(size_t err)
bool getWaitDebugInfo(WaitDebugInfo &info)
DebugState getDebugState(uintptr_t &address)
bool sendEvent(Event *pEvent)
bool pairWith(UnixSocketSyscalls *other)
virtual ReadyMask queryReady(bool reading, bool writing)
Filesystem * getRootFilesystem() const