11static volatile int pipe_signals;
12static void caught(
int number) {
13 if (number == SIGPIPE)
14 __atomic_add_fetch(&pipe_signals, 1, __ATOMIC_RELEASE);
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};
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));
27 call.output = pair[0];
28 CHECK(!pthread_create(&
worker, NULL, pt_worker, &call));
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);
36 CHECK(!pt_fill(p[1], 5, 59) && pt_ready(p[0], POLLIN, 0) == 0);
37 CHECK(!fcntl(p[0], F_SETFL, O_NONBLOCK));
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));
45 CHECK(!pt_wait(&call.done, 1000) && call.result == -1 && call.error == EINTR);
46 CHECK(!pt_join(
worker, &call));
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);
60 shutdown(pair[1], SHUT_RD);
65 for (
int n = 0; n < 2; ++n) {
73static int blocking_interrupt(
int kind) {
74 int failed = 0, p[2] = {-1, -1}, q[2] = {-1, -1}, running = 0;
76 struct iovec vector = {&byte, 1};
78 .kind = kind, .count = 1, .vectors = &vector, .vector_count = 1, .result = -2};
80 CHECK(!pipe(p) && !pipe(q));
83 if (kind == PT_VMSPLICE) {
84 CHECK(!pt_fill(p[1], PT_CAPACITY, 43));
85 CHECK(!fcntl(p[1], F_SETFL, O_NONBLOCK));
88 CHECK(!pthread_create(&
worker, NULL, pt_worker, &call));
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));
97 CHECK(!pt_wait(&call.done, 1000) && call.result == -1 && call.error == EINTR);
98 CHECK(!pt_join(
worker, &call));
100 if (kind == PT_VMSPLICE)
101 CHECK(!pt_read(p[0], PT_CAPACITY, 0, 43));
104 pt_diagnostic(&call);
110 pthread_kill(
worker, SIGUSR1);
113 for (
int n = 0; n < 2; ++n) {
121static int broken_pipe(
int kind) {
122 int failed = 0, p[2] = {-1, -1}, q[2] = {-1, -1};
124 struct iovec vector = {&byte, 1};
125 CHECK(!pipe(p) && !pipe(q) && !pt_fill(p[1], 8, 23));
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));
136 for (
int n = 0; n < 2; ++n) {
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};
149 snprintf(path,
sizeof(path),
"/tmp/pipe-transfer-reopen-%ld", (
long)getpid());
150 CHECK(!mkfifo(path, 0600));
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]);
158 call.output = pair[0];
159 CHECK(!pthread_create(&
worker, NULL, pt_worker, &call));
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));
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));
174 CHECK(!pt_wait(&call.done, 1000) && call.result == -1 && call.error == EINTR);
175 CHECK(!pt_join(
worker, &call));
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));
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));
190 pt_diagnostic(&call);
192 shutdown(pair[1], SHUT_RD);
199 for (
int n = 0; n < 2; ++n)
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))
211 if (sigaction(SIGPIPE, &action, &old_pipe)) {
212 sigaction(SIGUSR1, &old_usr, NULL);
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);