The Pedigree Project 0.1
readiness.c
1#define _GNU_SOURCE
2#include <poll.h>
3#include <signal.h>
4#include <unistd.h>
5
6#include "contract.h"
7#include <sys/epoll.h>
8#include <sys/socket.h>
9#include <sys/stat.h>
10
11static volatile int pipe_signals;
12static void caught(int number) {
13 if (number == SIGPIPE)
14 __atomic_add_fetch(&pipe_signals, 1, __ATOMIC_RELEASE);
15}
16static int cancel_reservation(void) {
17 int failed = 0, p[2] = {-1, -1}, pair[2] = {-1, -1}, epoll = -1, running = 0;
18 struct pt_call call = {.kind = PT_SPLICE, .count = 32, .result = -2};
19 pthread_t worker;
20 struct epoll_event event = {.events = EPOLLIN | EPOLLET, .data.u64 = 17};
21 CHECK(!pipe(p) && !socketpair(AF_UNIX, SOCK_STREAM, 0, pair) && !pt_fill(p[1], 32, 37));
22 ssize_t capacity = pt_socket_full(pair[0]);
23 CHECK(capacity > 0 && (epoll = epoll_create1(0)) >= 0);
24 CHECK(!epoll_ctl(epoll, EPOLL_CTL_ADD, p[0], &event));
25 CHECK(epoll_wait(epoll, &event, 1, 1000) == 1 && (event.events & EPOLLIN));
26 call.input = p[0];
27 call.output = pair[0];
28 CHECK(!pthread_create(&worker, NULL, pt_worker, &call));
29 running = 1;
30 CHECK(!pt_wait(&call.ready, 5000));
31 __atomic_store_n(&call.gate, 1, __ATOMIC_RELEASE);
32 CHECK(!pt_wait_readiness(p[0], POLLIN, 0));
33 CHECK(!__atomic_load_n(&call.done, __ATOMIC_ACQUIRE));
34 CHECK(epoll_wait(epoll, &event, 1, 0) == 0);
35 /* Appending behind a reserved head must not make that head readable. */
36 CHECK(!pt_fill(p[1], 5, 59) && pt_ready(p[0], POLLIN, 0) == 0);
37 CHECK(!fcntl(p[0], F_SETFL, O_NONBLOCK));
38 char byte;
39 CHECK(read(p[0], &byte, 1) == -1 && errno == EAGAIN);
40 const int64_t until = pt_now() + 3000000000;
41 while (!__atomic_load_n(&call.done, __ATOMIC_ACQUIRE) && pt_now() < until) {
42 CHECK(!pthread_kill(worker, SIGUSR1));
43 pt_pause(2);
44 }
45 CHECK(!pt_wait(&call.done, 1000) && call.result == -1 && call.error == EINTR);
46 CHECK(!pt_join(worker, &call));
47 running = 0;
48 CHECK(pt_ready(p[0], POLLIN, 1000) == 1);
49 CHECK(epoll_wait(epoll, &event, 1, 1000) == 1 && (event.events & EPOLLIN) &&
50 event.data.u64 == 17);
51 CHECK(epoll_wait(epoll, &event, 1, 0) == 0);
52 CHECK(!pt_read(p[0], 32, 0, 37) && !pt_read(p[0], 5, 0, 59));
53 CHECK(!pt_read(pair[1], capacity, 0, PT_FILLER));
54 CHECK(!fcntl(pair[1], F_SETFL, O_NONBLOCK));
55 CHECK(read(pair[1], &byte, 1) == -1 && errno == EAGAIN);
56out:
57 if (failed)
58 pt_diagnostic(&call);
59 if (running) {
60 shutdown(pair[1], SHUT_RD);
61 pt_join(worker, &call);
62 }
63 if (epoll >= 0)
64 close(epoll);
65 for (int n = 0; n < 2; ++n) {
66 if (p[n] >= 0)
67 close(p[n]);
68 if (pair[n] >= 0)
69 close(pair[n]);
70 }
71 return failed;
72}
73static int blocking_interrupt(int kind) {
74 int failed = 0, p[2] = {-1, -1}, q[2] = {-1, -1}, running = 0;
75 char byte = 'x';
76 struct iovec vector = {&byte, 1};
77 struct pt_call call = {
78 .kind = kind, .count = 1, .vectors = &vector, .vector_count = 1, .result = -2};
79 pthread_t worker;
80 CHECK(!pipe(p) && !pipe(q));
81 call.input = p[0];
82 call.output = q[1];
83 if (kind == PT_VMSPLICE) {
84 CHECK(!pt_fill(p[1], PT_CAPACITY, 43));
85 CHECK(!fcntl(p[1], F_SETFL, O_NONBLOCK));
86 call.input = p[1];
87 }
88 CHECK(!pthread_create(&worker, NULL, pt_worker, &call));
89 running = 1;
90 CHECK(!pt_wait(&call.ready, 5000));
91 __atomic_store_n(&call.gate, 1, __ATOMIC_RELEASE);
92 const int64_t until = pt_now() + 3000000000;
93 while (!__atomic_load_n(&call.done, __ATOMIC_ACQUIRE) && pt_now() < until) {
94 CHECK(!pthread_kill(worker, SIGUSR1));
95 pt_pause(2);
96 }
97 CHECK(!pt_wait(&call.done, 1000) && call.result == -1 && call.error == EINTR);
98 CHECK(!pt_join(worker, &call));
99 running = 0;
100 if (kind == PT_VMSPLICE)
101 CHECK(!pt_read(p[0], PT_CAPACITY, 0, 43));
102out:
103 if (failed)
104 pt_diagnostic(&call);
105 if (running) {
106 close(q[0]);
107 q[0] = -1;
108 close(p[0]);
109 p[0] = -1;
110 pthread_kill(worker, SIGUSR1);
111 pt_join(worker, &call);
112 }
113 for (int n = 0; n < 2; ++n) {
114 if (p[n] >= 0)
115 close(p[n]);
116 if (q[n] >= 0)
117 close(q[n]);
118 }
119 return failed;
120}
121static int broken_pipe(int kind) {
122 int failed = 0, p[2] = {-1, -1}, q[2] = {-1, -1};
123 char byte = 'x';
124 struct iovec vector = {&byte, 1};
125 CHECK(!pipe(p) && !pipe(q) && !pt_fill(p[1], 8, 23));
126 CHECK(!close(q[0]));
127 q[0] = -1;
128 __atomic_store_n(&pipe_signals, 0, __ATOMIC_RELEASE);
129 ssize_t result = kind == PT_TEE ? tee(p[0], q[1], 1, 0)
130 : kind == PT_VMSPLICE ? vmsplice(q[1], &vector, 1, 0)
131 : splice(p[0], NULL, q[1], NULL, 1, 0);
132 CHECK(result == -1 && errno == EPIPE);
133 CHECK(!pt_wait(&pipe_signals, 1000) && __atomic_load_n(&pipe_signals, __ATOMIC_ACQUIRE) == 1);
134 CHECK(!pt_read(p[0], 8, 0, 23));
135out:
136 for (int n = 0; n < 2; ++n) {
137 if (p[n] >= 0)
138 close(p[n]);
139 if (q[n] >= 0)
140 close(q[n]);
141 }
142 return failed;
143}
144static int fifo_reopen(void) {
145 int failed = 0, reader = -1, writer = -1, pair[2] = {-1, -1}, running = 0, created = 0;
146 char path[128], byte;
147 struct pt_call call = {.kind = PT_SPLICE, .count = 32, .result = -2};
148 pthread_t worker;
149 snprintf(path, sizeof(path), "/tmp/pipe-transfer-reopen-%ld", (long)getpid());
150 CHECK(!mkfifo(path, 0600));
151 created = 1;
152 CHECK((reader = open(path, O_RDONLY | O_NONBLOCK)) >= 0);
153 CHECK((writer = open(path, O_WRONLY | O_NONBLOCK)) >= 0);
154 CHECK(!pt_fill(writer, 32, 61) && !socketpair(AF_UNIX, SOCK_STREAM, 0, pair));
155 ssize_t capacity = pt_socket_full(pair[0]);
156 CHECK(capacity > 0);
157 call.input = reader;
158 call.output = pair[0];
159 CHECK(!pthread_create(&worker, NULL, pt_worker, &call));
160 running = 1;
161 CHECK(!pt_wait(&call.ready, 5000));
162 __atomic_store_n(&call.gate, 1, __ATOMIC_RELEASE);
163 CHECK(!pt_wait_readiness(reader, POLLIN, 0));
164 CHECK(!__atomic_load_n(&call.done, __ATOMIC_ACQUIRE));
165 CHECK(!close(writer));
166 writer = -1;
167 CHECK((writer = open(path, O_WRONLY | O_NONBLOCK)) >= 0);
168 CHECK(!pt_fill(writer, 8, 79) && pt_ready(reader, POLLIN, 0) == 0);
169 const int64_t until = pt_now() + 3000000000;
170 while (!__atomic_load_n(&call.done, __ATOMIC_ACQUIRE) && pt_now() < until) {
171 CHECK(!pthread_kill(worker, SIGUSR1));
172 pt_pause(2);
173 }
174 CHECK(!pt_wait(&call.done, 1000) && call.result == -1 && call.error == EINTR);
175 CHECK(!pt_join(worker, &call));
176 running = 0;
177 CHECK(!pt_read(reader, 32, 0, 61) && !pt_read(reader, 8, 0, 79));
178 CHECK(!pt_read(pair[1], capacity, 0, PT_FILLER));
179 CHECK(!pt_fill(writer, 16, 43));
180 CHECK(!close(writer));
181 writer = -1;
182 CHECK(!close(reader));
183 reader = -1;
184 CHECK((reader = open(path, O_RDONLY | O_NONBLOCK)) >= 0);
185 CHECK((writer = open(path, O_WRONLY | O_NONBLOCK)) >= 0);
186 CHECK(read(reader, &byte, 1) == -1 && errno == EAGAIN);
187 CHECK(!pt_fill(writer, 3, 19) && !pt_read(reader, 3, 0, 19));
188out:
189 if (failed)
190 pt_diagnostic(&call);
191 if (running) {
192 shutdown(pair[1], SHUT_RD);
193 pt_join(worker, &call);
194 }
195 if (writer >= 0)
196 close(writer);
197 if (reader >= 0)
198 close(reader);
199 for (int n = 0; n < 2; ++n)
200 if (pair[n] >= 0)
201 close(pair[n]);
202 if (created)
203 unlink(path);
204 return failed;
205}
206int pipe_transfer_readiness(void) {
207 struct sigaction action = {.sa_handler = caught}, old_usr, old_pipe;
208 sigemptyset(&action.sa_mask);
209 if (sigaction(SIGUSR1, &action, &old_usr))
210 return 1;
211 if (sigaction(SIGPIPE, &action, &old_pipe)) {
212 sigaction(SIGUSR1, &old_usr, NULL);
213 return 1;
214 }
215 int failed = cancel_reservation() || fifo_reopen();
216 for (int kind = 0; !failed && kind < 3; ++kind)
217 failed = blocking_interrupt(kind) || broken_pipe(kind);
218 sigaction(SIGPIPE, &old_pipe, NULL);
219 sigaction(SIGUSR1, &old_usr, NULL);
220 return failed;
221}
Definition waits.c:9