8static int collect(
int socket,
struct tf_transfer* transfer) {
9 const int64_t until = tf_now() + 5000000000;
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)
16 if (total == (
size_t)transfer->result)
19 struct pollfd entry = {.fd = socket, .events = POLLIN};
20 int ready = poll(&entry, 1, 100);
21 if (ready < 0 && errno == EINTR)
27 ssize_t count = read(socket, bytes,
sizeof(bytes));
30 for (ssize_t n = 0; n < count; ++n)
31 if (bytes[n] != tf_pattern(total + n, 47))
34 if (total > transfer->count)
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};
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));
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};
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));
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));
72 CHECK(read(replacement[1], &
byte, 1) == -1 && errno == EAGAIN);
75 tf_diagnostic(&transfer);
77 shutdown(pair[1], SHUT_RD);
78 tf_join(
worker, &transfer);
80 if (output_alias >= 0)
84 for (
int n = 0; n < 2; ++n) {
87 if (replacement[n] >= 0)
88 close(replacement[n]);