The Pedigree Project 0.1
unix-stream-interruption.c
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 <errno.h>
9#include <fcntl.h>
10#include <pthread.h>
11#include <sched.h>
12#include <signal.h>
13#include <stddef.h>
14#include <stdint.h>
15#include <stdio.h>
16#include <string.h>
17#include <unistd.h>
18
19#include <sys/socket.h>
20#include <sys/syscall.h>
21#include <sys/wait.h>
22
23extern void fail(void) __attribute__((noreturn));
24
25enum { wait_attempts = 100000 };
26
27enum io_operation {
28 receive_operation,
29 send_operation,
30};
31
32struct io_context {
33 int descriptor;
34 char* payload;
35 size_t length;
36 enum io_operation operation;
37 int use_message;
38 int rights_descriptor;
39 volatile int entered;
40 volatile int returned;
41 long tid;
42 ssize_t result;
43 int error;
44};
45
47 struct cmsghdr alignment;
48 unsigned char bytes[CMSG_SPACE(sizeof(int))];
49};
50
51static volatile sig_atomic_t signal_calls;
52static volatile sig_atomic_t close_from_signal = -1;
53static volatile sig_atomic_t signal_close_result = -2;
54
55static void signal_handler(int signal_number) {
56 if (signal_number == SIGUSR1) {
57 ++signal_calls;
58 const int descriptor = __atomic_exchange_n(&close_from_signal, -1, __ATOMIC_RELAXED);
59 if (descriptor >= 0) {
60 __atomic_store_n(&signal_close_result, close(descriptor), __ATOMIC_RELEASE);
61 }
62 }
63}
64
65static void* run_io(void* parameter) {
66 struct io_context* context = parameter;
67 context->tid = syscall(SYS_gettid);
68 errno = 0;
69 __atomic_store_n(&context->entered, 1, __ATOMIC_RELEASE);
70 if (context->use_message) {
71 struct iovec vector = {
72 .iov_base = context->payload,
73 .iov_len = context->length,
74 };
75 union control_buffer control = {0};
76 struct msghdr message = {
77 .msg_iov = &vector,
78 .msg_iovlen = 1,
79 .msg_control = control.bytes,
80 .msg_controllen = sizeof(control.bytes),
81 };
82 if (context->operation == send_operation) {
83 struct cmsghdr* header = CMSG_FIRSTHDR(&message);
84 header->cmsg_len = CMSG_LEN(sizeof(int));
85 header->cmsg_level = SOL_SOCKET;
86 header->cmsg_type = SCM_RIGHTS;
87 memcpy(CMSG_DATA(header), &context->rights_descriptor, sizeof(context->rights_descriptor));
88 context->result = sendmsg(context->descriptor, &message, MSG_NOSIGNAL);
89 } else {
90 context->result = recvmsg(context->descriptor, &message, 0);
91 }
92 } else if (context->operation == receive_operation) {
93 context->result = recv(context->descriptor, context->payload, context->length, 0);
94 } else {
95 context->result = send(context->descriptor, context->payload, context->length, MSG_NOSIGNAL);
96 }
97 context->error = errno;
98 __atomic_store_n(&context->returned, 1, __ATOMIC_RELEASE);
99 return 0;
100}
101
102static int wait_for_value(volatile int* value) {
103 for (size_t attempt = 0; attempt < wait_attempts; ++attempt) {
104 if (__atomic_load_n(value, __ATOMIC_ACQUIRE))
105 return 0;
106 sched_yield();
107 }
108 return -1;
109}
110
111static int install_signal_handler(void) {
112 struct sigaction action = {0};
113 signal_calls = 0;
114 __atomic_store_n(&close_from_signal, -1, __ATOMIC_RELAXED);
115 __atomic_store_n(&signal_close_result, -2, __ATOMIC_RELAXED);
116 action.sa_handler = signal_handler;
117 return sigemptyset(&action.sa_mask) || sigaction(SIGUSR1, &action, 0) ||
118 signal(SIGPIPE, SIG_IGN) == SIG_ERR;
119}
120
121static int interrupt_worker(struct io_context* context) {
122 for (size_t attempt = 0; attempt < 32; ++attempt) {
123 if (__atomic_load_n(&context->returned, __ATOMIC_ACQUIRE))
124 return 0;
125 if (syscall(SYS_tkill, context->tid, SIGUSR1))
126 return -1;
127 for (size_t pause = 0; pause < 32; ++pause) {
128 if (__atomic_load_n(&context->returned, __ATOMIC_ACQUIRE))
129 return 0;
130 sched_yield();
131 }
132 }
133 return -1;
134}
135
136static int fill_send_queue(int descriptor, const char* payload, size_t length) {
137 const int flags = fcntl(descriptor, F_GETFL);
138 if (flags < 0 || fcntl(descriptor, F_SETFL, flags | O_NONBLOCK))
139 return -1;
140
141 size_t written = 0;
142 while (1) {
143 errno = 0;
144 const ssize_t result = send(descriptor, payload, length, MSG_NOSIGNAL);
145 if (result > 0) {
146 written += (size_t)result;
147 continue;
148 }
149 if (result != -1 || (errno != EAGAIN && errno != EWOULDBLOCK))
150 return -1;
151 break;
152 }
153
154 return written && !fcntl(descriptor, F_SETFL, flags) ? 0 : -1;
155}
156
157static int interrupted_receive_child(void) {
158 alarm(10);
159 if (install_signal_handler())
160 return 10;
161
162 int sockets[2];
163 char payload = 0;
164 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets)) {
165 fprintf(stderr, "interrupted receive socketpair failed: errno=%d\n", errno);
166 return 11;
167 }
168
169 struct io_context context = {
170 .descriptor = sockets[1],
171 .payload = &payload,
172 .length = 1,
173 .operation = receive_operation,
174 .result = -2,
175 };
176 pthread_t worker;
177 if (pthread_create(&worker, 0, run_io, &context) || wait_for_value(&context.entered) ||
178 context.tid <= 0 || interrupt_worker(&context) || pthread_join(worker, 0))
179 return 12;
180
181 if (context.result != -1 || context.error != EINTR || signal_calls < 1)
182 return 13;
183 return close(sockets[0]) || close(sockets[1]) ? 14 : 0;
184}
185
186static int partial_send_child(void) {
187 alarm(10);
188 if (install_signal_handler())
189 return 20;
190
191 int sockets[2];
192 char payload[1024];
193 memset(payload, 'q', sizeof(payload));
194 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets) ||
195 fill_send_queue(sockets[0], payload, sizeof(payload)) || recv(sockets[1], payload, 1, 0) != 1)
196 return 21;
197
198 payload[0] = 'x';
199 payload[1] = 'y';
200 struct io_context context = {
201 .descriptor = sockets[0],
202 .payload = payload,
203 .length = 2,
204 .operation = send_operation,
205 .result = -2,
206 };
207 pthread_t worker;
208 if (pthread_create(&worker, 0, run_io, &context) || wait_for_value(&context.entered) ||
209 context.tid <= 0 || interrupt_worker(&context) || pthread_join(worker, 0))
210 return 22;
211
212 if (context.result != 1 || context.error || signal_calls < 1)
213 return 23;
214 return close(sockets[0]) || close(sockets[1]) ? 24 : 0;
215}
216
217static int interrupted_send_child(void) {
218 alarm(10);
219 if (install_signal_handler())
220 return 25;
221
222 int sockets[2];
223 char payload[1024];
224 memset(payload, 'i', sizeof(payload));
225 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets) ||
226 fill_send_queue(sockets[0], payload, sizeof(payload)))
227 return 26;
228
229 struct io_context context = {
230 .descriptor = sockets[0],
231 .payload = payload,
232 .length = 1,
233 .operation = send_operation,
234 .result = -2,
235 };
236 pthread_t worker;
237 if (pthread_create(&worker, 0, run_io, &context) || wait_for_value(&context.entered) ||
238 context.tid <= 0 || interrupt_worker(&context) || pthread_join(worker, 0))
239 return 27;
240
241 if (context.result != -1 || context.error != EINTR || signal_calls < 1)
242 return 28;
243 return close(sockets[0]) || close(sockets[1]) ? 29 : 0;
244}
245
246static int signal_handler_close_child(enum io_operation operation) {
247 alarm(10);
248 if (install_signal_handler())
249 return 50;
250
251 int sockets[2];
252 char payload[1024];
253 memset(payload, 'h', sizeof(payload));
254 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets))
255 return 51;
256 if (operation == send_operation && fill_send_queue(sockets[0], payload, sizeof(payload)))
257 return 52;
258
259 int rights_pipe[2] = {-1, -1};
260 if (operation == send_operation && pipe(rights_pipe))
261 return 53;
262
263 const int operation_side = operation == receive_operation ? 1 : 0;
264 const int close_side = operation == receive_operation ? operation_side : 1 - operation_side;
265 struct io_context context = {
266 .descriptor = sockets[operation_side],
267 .payload = payload,
268 .length = 1,
269 .operation = operation,
270 .use_message = 1,
271 .rights_descriptor = rights_pipe[0],
272 .result = -2,
273 };
274 pthread_t worker;
275 if (pthread_create(&worker, 0, run_io, &context) || wait_for_value(&context.entered))
276 return 54;
277 for (size_t pause = 0; pause < 512; ++pause)
278 sched_yield();
279
280 __atomic_store_n(&close_from_signal, sockets[close_side], __ATOMIC_RELEASE);
281 if (__atomic_load_n(&context.returned, __ATOMIC_ACQUIRE) || context.tid <= 0 ||
282 interrupt_worker(&context) || pthread_join(worker, 0))
283 return 55;
284 if (__atomic_load_n(&signal_close_result, __ATOMIC_ACQUIRE) || context.result != -1 ||
285 context.error != EINTR || signal_calls < 1)
286 return 56;
287
288 sockets[close_side] = -1;
289 if (sockets[0] >= 0 && close(sockets[0]))
290 return 57;
291 if (sockets[1] >= 0 && close(sockets[1]))
292 return 58;
293 if (rights_pipe[0] >= 0 && close(rights_pipe[0]))
294 return 59;
295 if (rights_pipe[1] >= 0 && close(rights_pipe[1]))
296 return 60;
297 return 0;
298}
299
300static int receive_signal_handler_close_child(void) {
301 return signal_handler_close_child(receive_operation);
302}
303
304static int send_signal_handler_close_child(void) {
305 return signal_handler_close_child(send_operation);
306}
307
308static int serialized_wait_child(enum io_operation operation) {
309 alarm(10);
310 if (install_signal_handler())
311 return 40;
312
313 int sockets[2];
314 char payload[1024];
315 memset(payload, 's', sizeof(payload));
316 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets))
317 return 41;
318 if (operation == send_operation && fill_send_queue(sockets[0], payload, sizeof(payload)))
319 return 42;
320
321 const int operation_side = operation == receive_operation ? 1 : 0;
322 struct io_context holder = {
323 .descriptor = sockets[operation_side],
324 .payload = payload,
325 .length = 1,
326 .operation = operation,
327 .result = -2,
328 };
329 struct io_context waiter = {
330 .descriptor = sockets[operation_side],
331 .payload = payload,
332 .length = 1,
333 .operation = operation,
334 .result = -2,
335 };
336 pthread_t holder_thread;
337 pthread_t waiter_thread;
338 if (pthread_create(&holder_thread, 0, run_io, &holder) || wait_for_value(&holder.entered))
339 return 43;
340 for (size_t pause = 0; pause < 512; ++pause)
341 sched_yield();
342 if (__atomic_load_n(&holder.returned, __ATOMIC_ACQUIRE) ||
343 pthread_create(&waiter_thread, 0, run_io, &waiter) || wait_for_value(&waiter.entered) ||
344 waiter.tid <= 0 || interrupt_worker(&waiter) || pthread_join(waiter_thread, 0))
345 return 44;
346
347 if (waiter.result != -1 || waiter.error != EINTR || signal_calls < 1 ||
348 __atomic_load_n(&holder.returned, __ATOMIC_ACQUIRE))
349 return 45;
350 if (close(sockets[0]) || close(sockets[1]) || pthread_join(holder_thread, 0))
351 return 46;
352
353 const int holder_result_ok = operation == receive_operation
354 ? holder.result == 0
355 : holder.result == -1 && holder.error == EPIPE;
356 return holder_result_ok ? 0 : 47;
357}
358
359static int serialized_receive_wait_child(void) {
360 return serialized_wait_child(receive_operation);
361}
362
363static int serialized_send_wait_child(void) {
364 return serialized_wait_child(send_operation);
365}
366
367static int close_wake_child(enum io_operation operation, int close_local) {
368 alarm(10);
369 if (install_signal_handler())
370 return 30;
371
372 int sockets[2];
373 char payload[1024];
374 memset(payload, 'c', sizeof(payload));
375 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets))
376 return 31;
377 if (operation == send_operation && fill_send_queue(sockets[0], payload, sizeof(payload)))
378 return 32;
379
380 const int operation_side = operation == receive_operation ? 1 : 0;
381 const int close_side = close_local ? operation_side : 1 - operation_side;
382 struct io_context context = {
383 .descriptor = sockets[operation_side],
384 .payload = payload,
385 .length = 1,
386 .operation = operation,
387 .result = -2,
388 };
389 pthread_t worker;
390 if (pthread_create(&worker, 0, run_io, &context) || wait_for_value(&context.entered))
391 return 33;
392 for (size_t pause = 0; pause < 256; ++pause)
393 sched_yield();
394
395 if (close(sockets[close_side]))
396 return 34;
397 sockets[close_side] = -1;
398 const int returned_without_rescue = !wait_for_value(&context.returned);
399 if (!returned_without_rescue) {
400 if (sockets[0] >= 0)
401 close(sockets[0]);
402 if (sockets[1] >= 0)
403 close(sockets[1]);
404 (void)syscall(SYS_tkill, context.tid, SIGUSR1);
405 }
406 if (pthread_join(worker, 0))
407 return 35;
408
409 const int result_ok = operation == receive_operation
410 ? context.result == 0
411 : context.result == -1 && context.error == EPIPE;
412 if (!returned_without_rescue || !result_ok)
413 return 36;
414 if (sockets[0] >= 0 && close(sockets[0]))
415 return 37;
416 if (sockets[1] >= 0 && close(sockets[1]))
417 return 38;
418 return 0;
419}
420
421static void run_bounded(int (*child_test)(void)) {
422 const pid_t child = fork();
423 if (child < 0)
424 fail();
425 if (!child)
426 _exit(child_test());
427
428 int status_code = 0;
429 pid_t waited;
430 do {
431 waited = waitpid(child, &status_code, 0);
432 } while (waited < 0 && errno == EINTR);
433 if (waited != child || !WIFEXITED(status_code) || WEXITSTATUS(status_code)) {
434 if (waited == child && WIFEXITED(status_code))
435 fprintf(stderr, "AF_UNIX interruption child failed: code=%d\n", WEXITSTATUS(status_code));
436 fail();
437 }
438}
439
440static int peer_receive_close_child(void) {
441 return close_wake_child(receive_operation, 0);
442}
443
444static int local_receive_close_child(void) {
445 return close_wake_child(receive_operation, 1);
446}
447
448static int peer_send_close_child(void) {
449 return close_wake_child(send_operation, 0);
450}
451
452static int local_send_close_child(void) {
453 return close_wake_child(send_operation, 1);
454}
455
456void test_unix_stream_interruption(void) {
457 puts("Testing AF_UNIX stream interruption and close wakeups... ");
458 fflush(stdout);
459 run_bounded(interrupted_receive_child);
460 run_bounded(interrupted_send_child);
461 run_bounded(receive_signal_handler_close_child);
462 run_bounded(send_signal_handler_close_child);
463 run_bounded(serialized_receive_wait_child);
464 run_bounded(serialized_send_wait_child);
465 run_bounded(partial_send_child);
466 run_bounded(peer_receive_close_child);
467 run_bounded(local_receive_close_child);
468 run_bounded(peer_send_close_child);
469 run_bounded(local_send_close_child);
470 puts("OK\n");
471 fflush(stdout);
472}