20#include <sys/syscall.h>
25enum { wait_attempts = 100000 };
36 enum io_operation operation;
38 int rights_descriptor;
40 volatile int returned;
47 struct cmsghdr alignment;
48 unsigned char bytes[CMSG_SPACE(
sizeof(
int))];
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;
55static void signal_handler(
int signal_number) {
56 if (signal_number == SIGUSR1) {
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);
65static void* run_io(
void* parameter) {
67 context->tid = syscall(SYS_gettid);
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,
79 .msg_control = control.bytes,
80 .msg_controllen =
sizeof(control.bytes),
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);
90 context->result = recvmsg(context->descriptor, &
message, 0);
92 }
else if (context->operation == receive_operation) {
93 context->result = recv(context->descriptor, context->payload, context->length, 0);
95 context->result = send(context->descriptor, context->payload, context->length, MSG_NOSIGNAL);
97 context->error = errno;
98 __atomic_store_n(&context->returned, 1, __ATOMIC_RELEASE);
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))
111static int install_signal_handler(
void) {
112 struct sigaction action = {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;
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))
125 if (syscall(SYS_tkill, context->tid, SIGUSR1))
127 for (
size_t pause = 0; pause < 32; ++pause) {
128 if (__atomic_load_n(&context->returned, __ATOMIC_ACQUIRE))
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))
144 const ssize_t result = send(descriptor, payload, length, MSG_NOSIGNAL);
146 written += (size_t)result;
149 if (result != -1 || (errno != EAGAIN && errno != EWOULDBLOCK))
154 return written && !fcntl(descriptor, F_SETFL, flags) ? 0 : -1;
157static int interrupted_receive_child(
void) {
159 if (install_signal_handler())
164 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets)) {
165 fprintf(stderr,
"interrupted receive socketpair failed: errno=%d\n", errno);
170 .descriptor = sockets[1],
173 .operation = receive_operation,
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))
181 if (context.result != -1 || context.error != EINTR || signal_calls < 1)
183 return close(sockets[0]) || close(sockets[1]) ? 14 : 0;
186static int partial_send_child(
void) {
188 if (install_signal_handler())
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)
201 .descriptor = sockets[0],
204 .operation = send_operation,
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))
212 if (context.result != 1 || context.error || signal_calls < 1)
214 return close(sockets[0]) || close(sockets[1]) ? 24 : 0;
217static int interrupted_send_child(
void) {
219 if (install_signal_handler())
224 memset(payload,
'i',
sizeof(payload));
225 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets) ||
226 fill_send_queue(sockets[0], payload,
sizeof(payload)))
230 .descriptor = sockets[0],
233 .operation = send_operation,
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))
241 if (context.result != -1 || context.error != EINTR || signal_calls < 1)
243 return close(sockets[0]) || close(sockets[1]) ? 29 : 0;
246static int signal_handler_close_child(
enum io_operation operation) {
248 if (install_signal_handler())
253 memset(payload,
'h',
sizeof(payload));
254 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets))
256 if (operation == send_operation && fill_send_queue(sockets[0], payload,
sizeof(payload)))
259 int rights_pipe[2] = {-1, -1};
260 if (operation == send_operation && pipe(rights_pipe))
263 const int operation_side = operation == receive_operation ? 1 : 0;
264 const int close_side = operation == receive_operation ? operation_side : 1 - operation_side;
266 .descriptor = sockets[operation_side],
269 .operation = operation,
271 .rights_descriptor = rights_pipe[0],
275 if (pthread_create(&
worker, 0, run_io, &context) || wait_for_value(&context.entered))
277 for (
size_t pause = 0; pause < 512; ++pause)
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))
284 if (__atomic_load_n(&signal_close_result, __ATOMIC_ACQUIRE) || context.result != -1 ||
285 context.error != EINTR || signal_calls < 1)
288 sockets[close_side] = -1;
289 if (sockets[0] >= 0 && close(sockets[0]))
291 if (sockets[1] >= 0 && close(sockets[1]))
293 if (rights_pipe[0] >= 0 && close(rights_pipe[0]))
295 if (rights_pipe[1] >= 0 && close(rights_pipe[1]))
300static int receive_signal_handler_close_child(
void) {
301 return signal_handler_close_child(receive_operation);
304static int send_signal_handler_close_child(
void) {
305 return signal_handler_close_child(send_operation);
308static int serialized_wait_child(
enum io_operation operation) {
310 if (install_signal_handler())
315 memset(payload,
's',
sizeof(payload));
316 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets))
318 if (operation == send_operation && fill_send_queue(sockets[0], payload,
sizeof(payload)))
321 const int operation_side = operation == receive_operation ? 1 : 0;
323 .descriptor = sockets[operation_side],
326 .operation = operation,
330 .descriptor = sockets[operation_side],
333 .operation = operation,
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))
340 for (
size_t pause = 0; pause < 512; ++pause)
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))
347 if (
waiter.result != -1 ||
waiter.error != EINTR || signal_calls < 1 ||
348 __atomic_load_n(&holder.returned, __ATOMIC_ACQUIRE))
350 if (close(sockets[0]) || close(sockets[1]) || pthread_join(holder_thread, 0))
353 const int holder_result_ok = operation == receive_operation
355 : holder.result == -1 && holder.error == EPIPE;
356 return holder_result_ok ? 0 : 47;
359static int serialized_receive_wait_child(
void) {
360 return serialized_wait_child(receive_operation);
363static int serialized_send_wait_child(
void) {
364 return serialized_wait_child(send_operation);
367static int close_wake_child(
enum io_operation operation,
int close_local) {
369 if (install_signal_handler())
374 memset(payload,
'c',
sizeof(payload));
375 if (socketpair(AF_UNIX, SOCK_STREAM, 0, sockets))
377 if (operation == send_operation && fill_send_queue(sockets[0], payload,
sizeof(payload)))
380 const int operation_side = operation == receive_operation ? 1 : 0;
381 const int close_side = close_local ? operation_side : 1 - operation_side;
383 .descriptor = sockets[operation_side],
386 .operation = operation,
390 if (pthread_create(&
worker, 0, run_io, &context) || wait_for_value(&context.entered))
392 for (
size_t pause = 0; pause < 256; ++pause)
395 if (close(sockets[close_side]))
397 sockets[close_side] = -1;
398 const int returned_without_rescue = !wait_for_value(&context.returned);
399 if (!returned_without_rescue) {
404 (void)syscall(SYS_tkill, context.tid, SIGUSR1);
406 if (pthread_join(
worker, 0))
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)
414 if (sockets[0] >= 0 && close(sockets[0]))
416 if (sockets[1] >= 0 && close(sockets[1]))
421static void run_bounded(
int (*child_test)(
void)) {
422 const pid_t child = fork();
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));
440static int peer_receive_close_child(
void) {
441 return close_wake_child(receive_operation, 0);
444static int local_receive_close_child(
void) {
445 return close_wake_child(receive_operation, 1);
448static int peer_send_close_child(
void) {
449 return close_wake_child(send_operation, 0);
452static int local_send_close_child(
void) {
453 return close_wake_child(send_operation, 1);
456void test_unix_stream_interruption(
void) {
457 puts(
"Testing AF_UNIX stream interruption and close wakeups... ");
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);