15 return clock_gettime(CLOCK_MONOTONIC, &now) ? -1 : (int64_t)now.tv_sec * 1000000000 + now.tv_nsec;
17void pt_pause(
int milliseconds) {
18 struct timespec pause = {milliseconds / 1000, (milliseconds % 1000) * 1000000L};
19 while (nanosleep(&pause, &pause) && errno == EINTR) {
22int pt_wait(
volatile int* flag,
int milliseconds) {
23 int64_t until = pt_now() + (int64_t)milliseconds * 1000000;
24 while (!__atomic_load_n(flag, __ATOMIC_ACQUIRE)) {
25 if (pt_now() >= until)
31int pt_reap(pid_t child,
int milliseconds) {
32 int64_t until = pt_now() + (int64_t)milliseconds * 1000000;
33 while (pt_now() < until) {
35 pid_t result = waitpid(child, &status, WNOHANG);
36 if (result == child) {
37 int code = WIFEXITED(status) ? WEXITSTATUS(status) : 128 + WTERMSIG(status);
39 fprintf(stderr,
"PIPE-TRANSFER-CONTRACT: child=%ld status=%d\n", (
long)child, code);
42 if (result < 0 && errno != EINTR)
46 fprintf(stderr,
"PIPE-TRANSFER-CONTRACT: child=%ld timeout\n", (
long)child);
48 while (waitpid(child, NULL, 0) < 0 && errno == EINTR) {
52int pt_ready(
int fd,
short events,
int milliseconds) {
53 struct pollfd entry = {.fd = fd, .events = events};
56 result = poll(&entry, 1, milliseconds);
57 while (result < 0 && errno == EINTR);
58 return result < 0 ? -1 : !!(entry.revents & events);
60int pt_wait_readiness(
int fd,
short events,
int expected) {
61 int64_t until = pt_now() + 3000000000;
63 int ready = pt_ready(fd, events, 0);
64 if (ready < 0 || ready == expected)
65 return ready < 0 ? -1 : 0;
67 }
while (pt_now() < until);
70unsigned char pt_pattern(
size_t offset,
int seed) {
71 return seed == PT_FILLER ? 0x7a : ((offset * 17) ^ (offset >> 9) ^ seed) & 255;
73int pt_fill(
int fd,
size_t count,
int seed) {
74 unsigned char bytes[PT_CAPACITY];
75 if (count >
sizeof(bytes))
77 for (
size_t n = 0; n < count; ++n)
78 bytes[n] = pt_pattern(n, seed);
79 return write(fd, bytes, count) == (ssize_t)count ? 0 : -1;
81int pt_read(
int fd,
size_t count,
size_t source_offset,
int seed) {
82 unsigned char bytes[PT_CAPACITY];
84 while (done < count) {
85 if (pt_ready(fd, POLLIN, 3000) != 1)
87 size_t amount = count - done <
sizeof(bytes) ? count - done : sizeof(bytes);
88 ssize_t received = read(fd, bytes, amount);
91 for (ssize_t n = 0; n < received; ++n)
92 if (bytes[n] != pt_pattern(source_offset + done + n, seed))
98int pt_file(
size_t count,
int seed) {
99 int fd = memfd_create(
"pipe-transfer", MFD_ALLOW_SEALING);
102 unsigned char bytes[PT_CAPACITY];
103 for (
size_t offset = 0; offset < count;) {
104 size_t amount = count - offset <
sizeof(bytes) ? count - offset : sizeof(bytes);
105 for (
size_t n = 0; n < amount; ++n)
106 bytes[n] = pt_pattern(offset + n, seed);
107 if (pwrite(fd, bytes, amount, offset) != (ssize_t)amount) {
115int pt_verify(
int fd, off_t offset,
size_t count,
size_t source_offset,
int seed) {
116 unsigned char bytes[PT_CAPACITY];
117 if (count >
sizeof(bytes) || pread(fd, bytes, count, offset) != (ssize_t)count)
119 for (
size_t n = 0; n < count; ++n)
120 if (bytes[n] != pt_pattern(source_offset + n, seed))
124ssize_t pt_socket_full(
int fd) {
125 int flags = fcntl(fd, F_GETFL);
126 if (flags < 0 || fcntl(fd, F_SETFL, flags | O_NONBLOCK))
129 memset(bytes, 0x7a,
sizeof(bytes));
131 while (total < 1024 * 1024) {
132 ssize_t written = write(fd, bytes,
sizeof(bytes));
133 if (written < 0 && errno == EAGAIN)
134 return fcntl(fd, F_SETFL, flags) ? -1 : total;
139 fcntl(fd, F_SETFL, flags);
142int pt_send_fd(
int socket,
int fd) {
144 struct iovec vector = {&byte, 1};
146 struct cmsghdr align;
147 char bytes[CMSG_SPACE(
sizeof(
int))];
149 struct msghdr
message = {.msg_iov = &vector,
151 .msg_control = control.bytes,
152 .msg_controllen =
sizeof(control.bytes)};
153 struct cmsghdr* item = CMSG_FIRSTHDR(&
message);
154 item->cmsg_level = SOL_SOCKET;
155 item->cmsg_type = SCM_RIGHTS;
156 item->cmsg_len = CMSG_LEN(
sizeof(
int));
157 memcpy(CMSG_DATA(item), &fd,
sizeof(fd));
158 return sendmsg(socket, &
message, 0) == 1 ? 0 : -1;
160int pt_receive_fd(
int socket) {
162 struct iovec vector = {&byte, 1};
164 struct cmsghdr align;
165 char bytes[CMSG_SPACE(
sizeof(
int))];
167 struct msghdr
message = {.msg_iov = &vector,
169 .msg_control = control.bytes,
170 .msg_controllen =
sizeof(control.bytes)};
171 if (pt_ready(socket, POLLIN, 3000) != 1 || recvmsg(socket, &
message, MSG_CMSG_CLOEXEC) != 1 ||
172 (
message.msg_flags & MSG_CTRUNC))
174 struct cmsghdr* item = CMSG_FIRSTHDR(&
message);
175 if (!item || item->cmsg_level != SOL_SOCKET || item->cmsg_type != SCM_RIGHTS ||
176 item->cmsg_len != CMSG_LEN(
sizeof(
int)))
179 memcpy(&fd, CMSG_DATA(item),
sizeof(fd));
182void* pt_worker(
void* argument) {
183 struct pt_call* call = argument;
184 __atomic_store_n(&call->ready, 1, __ATOMIC_RELEASE);
185 if (!pt_wait(&call->gate, 5000)) {
187 __atomic_store_n(&call->entered, 1, __ATOMIC_RELEASE);
188 if (call->kind == PT_TEE)
189 call->result = tee(call->input, call->output, call->count, call->flags);
190 else if (call->kind == PT_VMSPLICE)
191 call->result = vmsplice(call->input, call->vectors, call->vector_count, call->flags);
193 call->result = splice(call->input, call->input_offset, call->output, call->output_offset,
194 call->count, call->flags);
197 call->error = ETIMEDOUT;
198 __atomic_store_n(&call->done, 1, __ATOMIC_RELEASE);
201void pt_diagnostic(
const struct pt_call* call) {
202 int done = __atomic_load_n(&call->done, __ATOMIC_ACQUIRE);
203 fprintf(stderr,
"PIPE-TRANSFER-CONTRACT: kind=%d count=%zu done=%d result=%ld/%d\n", call->kind,
204 call->count, done, done ? (
long)call->result : -2, done ? call->error : -2);
207 __atomic_store_n(&call->gate, 1, __ATOMIC_RELEASE);
208 if (pt_wait(&call->done, 5000)) {
212 return pthread_join(
worker, NULL);
214static int run(
const char* name,
int (*test)(
void)) {
215 printf(
"PIPE-TRANSFER-CONTRACT: BEGIN %s\n", name);
217 pid_t child = fork();
225 _exit(result ? 1 : 0);
227 int status = pt_reap(child, 45000);
228 printf(
"PIPE-TRANSFER-CONTRACT: %s %s status=%d\n", status ?
"FAIL" :
"PASS", name, status);
232int main(
int argc,
char** argv) {
233 if (signal(SIGPIPE, SIG_IGN) == SIG_ERR)
238 } suites[] = {{
"splice", pipe_transfer_splice}, {
"tee", pipe_transfer_tee},
239 {
"vmsplice", pipe_transfer_vmsplice}, {
"lifetime", pipe_transfer_lifetime},
240 {
"readiness", pipe_transfer_readiness}, {
"concurrency", pipe_transfer_concurrency}};
242 for (
unsigned n = 0; n <
sizeof(suites) /
sizeof(suites[0]); ++n) {
243 if (argc > 1 && strcmp(argv[1], suites[n].name))
246 if (run(suites[n].name, suites[n].test))
251 puts(
"PIPE-TRANSFER-CONTRACT: END PASS");