The Pedigree Project 0.1
transfer-contract-test/lifetime.c
1#define _GNU_SOURCE
2#include <poll.h>
3#include <unistd.h>
4
5#include "contract.h"
6#include <sys/socket.h>
7
8static int collect(int socket, struct tf_transfer* transfer) {
9 const int64_t until = tf_now() + 5000000000;
10 size_t total = 0;
11 unsigned char bytes[4096];
12 while (tf_now() < until) {
13 if (__atomic_load_n(&transfer->done, __ATOMIC_ACQUIRE)) {
14 if (transfer->result <= 0 || total > (size_t)transfer->result)
15 return -1;
16 if (total == (size_t)transfer->result)
17 return 0;
18 }
19 struct pollfd entry = {.fd = socket, .events = POLLIN};
20 int ready = poll(&entry, 1, 100);
21 if (ready < 0 && errno == EINTR)
22 continue;
23 if (ready < 0)
24 return -1;
25 if (!ready)
26 continue;
27 ssize_t count = read(socket, bytes, sizeof(bytes));
28 if (count <= 0)
29 return -1;
30 for (ssize_t n = 0; n < count; ++n)
31 if (bytes[n] != tf_pattern(total + n, 47))
32 return -1;
33 total += count;
34 if (total > transfer->count)
35 return -1;
36 }
37 return -1;
38}
39int transfer_lifetime(void) {
40 int failed = 0, pair[2] = {-1, -1}, replacement[2] = {-1, -1};
41 int input_alias = -1, output_alias = -1, running = 0;
42 struct tf_file input = {.fd = -1}, other = {.fd = -1};
43 struct tf_transfer transfer = {.kind = TF_SENDFILE, .result = -2};
44 pthread_t worker;
45 CHECK(!socketpair(AF_UNIX, SOCK_STREAM, 0, pair));
46 ssize_t capacity = tf_saturate(pair[0]);
47 CHECK(capacity > 0 && !tf_socket_read(pair[1], capacity, 0, TF_FILLER));
48 CHECK(!tf_create(&input, "/tmp", capacity + 2 * TF_CHUNK, 47));
49 CHECK(!tf_create(&other, "/tmp", 64, 89));
50 CHECK(!socketpair(AF_UNIX, SOCK_STREAM, 0, replacement));
51 CHECK((input_alias = dup(input.fd)) >= 0 && (output_alias = dup(pair[0])) >= 0);
52 transfer.input = input.fd;
53 transfer.output = pair[0];
54 transfer.count = capacity + 2 * TF_CHUNK;
55 CHECK(!pthread_create(&worker, NULL, tf_worker, &transfer));
56 running = 1;
57 CHECK(!tf_wait(&transfer.ready, 5000));
58 __atomic_store_n(&transfer.gate, 1, __ATOMIC_RELEASE);
59 struct pollfd readable = {.fd = pair[1], .events = POLLIN};
60 /* Only this syscall can make the initially empty receive queue readable. */
61 CHECK(poll(&readable, 1, 5000) == 1 && (readable.revents & POLLIN));
62 CHECK(!__atomic_load_n(&transfer.done, __ATOMIC_ACQUIRE));
63 CHECK(dup2(other.fd, input.fd) == input.fd);
64 CHECK(dup2(replacement[0], pair[0]) == pair[0]);
65 CHECK(!collect(pair[1], &transfer) && !tf_join(worker, &transfer));
66 running = 0;
67 CHECK(transfer.result > 0 && transfer.result <= (ssize_t)transfer.count && transfer.error == 0);
68 CHECK(lseek(input_alias, 0, SEEK_CUR) == transfer.result);
69 CHECK(lseek(input.fd, 0, SEEK_CUR) == 0 && !tf_verify(other.fd, 0, 64, 0, 89));
70 CHECK(!fcntl(replacement[1], F_SETFL, O_NONBLOCK));
71 char byte;
72 CHECK(read(replacement[1], &byte, 1) == -1 && errno == EAGAIN);
73out:
74 if (failed)
75 tf_diagnostic(&transfer);
76 if (running) {
77 shutdown(pair[1], SHUT_RD);
78 tf_join(worker, &transfer);
79 }
80 if (output_alias >= 0)
81 close(output_alias);
82 if (input_alias >= 0)
83 close(input_alias);
84 for (int n = 0; n < 2; ++n) {
85 if (pair[n] >= 0)
86 close(pair[n]);
87 if (replacement[n] >= 0)
88 close(replacement[n]);
89 }
90 tf_destroy(&other);
91 tf_destroy(&input);
92 return failed;
93}