The Pedigree Project 0.1
pipe-transfer-contract-test/concurrency.c
1#define _GNU_SOURCE
2#include <string.h>
3#include <unistd.h>
4
5#include "contract.h"
6
7static int atomic_capacity(void) {
8 int failed = 0, p[2] = {-1, -1};
9 CHECK(!pipe(p) && !fcntl(p[1], F_SETFL, O_NONBLOCK));
10 CHECK(!pt_fill(p[1], PT_CAPACITY - 1, 17));
11 CHECK(write(p[1], "xx", 2) == -1 && errno == EAGAIN);
12 CHECK(!pt_read(p[0], PT_CAPACITY - 1, 0, 17));
13 CHECK(!pt_fill(p[1], PT_CAPACITY, 43));
14 CHECK(write(p[1], "x", 1) == -1 && errno == EAGAIN);
15 CHECK(!pt_read(p[0], PT_CAPACITY, 0, 43));
16out:
17 if (p[0] >= 0)
18 close(p[0]);
19 if (p[1] >= 0)
20 close(p[1]);
21 return failed;
22}
23static int count_bytes(int fd, size_t maximum, size_t totals[2]) {
24 unsigned char bytes[PT_CAPACITY];
25 ssize_t count = read(fd, bytes, maximum);
26 if (count <= 0)
27 return -1;
28 for (ssize_t n = 0; n < count; ++n) {
29 if (bytes[n] != 'a' && bytes[n] != 'b') {
30 errno = EILSEQ;
31 return -1;
32 }
33 ++totals[bytes[n] - 'a'];
34 }
35 return count;
36}
37static int opposite_pair(int kind) {
38 int failed = 0, p[2] = {-1, -1}, q[2] = {-1, -1}, started = 0;
39 struct pt_call calls[2] = {{.kind = kind, .count = PT_CAPACITY, .result = -2},
40 {.kind = kind, .count = PT_CAPACITY, .result = -2}};
41 pthread_t workers[2];
42 char bytes[PT_CAPACITY];
43 size_t totals[2] = {0};
44 CHECK(!pipe(p) && !pipe(q));
45 memset(bytes, 'a', sizeof(bytes));
46 CHECK(write(p[1], bytes, sizeof(bytes)) == sizeof(bytes));
47 memset(bytes, 'b', sizeof(bytes));
48 CHECK(write(q[1], bytes, sizeof(bytes)) == sizeof(bytes));
49 calls[0].input = p[0];
50 calls[0].output = q[1];
51 calls[1].input = q[0];
52 calls[1].output = p[1];
53 for (int n = 0; n < 2; ++n) {
54 CHECK(!pthread_create(&workers[n], NULL, pt_worker, &calls[n]));
55 ++started;
56 }
57 CHECK(!pt_wait(&calls[0].ready, 5000) && !pt_wait(&calls[1].ready, 5000));
58 __atomic_store_n(&calls[0].gate, 1, __ATOMIC_RELEASE);
59 __atomic_store_n(&calls[1].gate, 1, __ATOMIC_RELEASE);
60 /* Leave source data available while creating room on both initially full destinations. */
61 CHECK(count_bytes(p[0], 128, totals) == 128 && count_bytes(q[0], 128, totals) == 128);
62 CHECK(!pt_wait(&calls[0].done, 5000) && !pt_wait(&calls[1].done, 5000));
63 CHECK(calls[0].result > 0 && calls[0].result <= PT_CAPACITY && calls[1].result > 0 &&
64 calls[1].result <= PT_CAPACITY);
65 CHECK(!fcntl(p[0], F_SETFL, O_NONBLOCK) && !fcntl(q[0], F_SETFL, O_NONBLOCK));
66 const int readers[] = {p[0], q[0]};
67 for (int n = 0; n < 2; ++n) {
68 while (count_bytes(readers[n], PT_CAPACITY, totals) > 0) {
69 }
70 CHECK(errno == EAGAIN);
71 }
72 CHECK(totals[0] == PT_CAPACITY + (kind == PT_TEE ? (size_t)calls[0].result : 0));
73 CHECK(totals[1] == PT_CAPACITY + (kind == PT_TEE ? (size_t)calls[1].result : 0));
74out:
75 if (failed) {
76 pt_diagnostic(&calls[0]);
77 pt_diagnostic(&calls[1]);
78 }
79 for (int n = 0; n < started; ++n)
80 __atomic_store_n(&calls[n].gate, 1, __ATOMIC_RELEASE);
81 for (int n = 0; n < started; ++n)
82 pt_join(workers[n], &calls[n]);
83 for (int n = 0; n < 2; ++n) {
84 if (p[n] >= 0)
85 close(p[n]);
86 if (q[n] >= 0)
87 close(q[n]);
88 }
89 return failed;
90}
91int pipe_transfer_concurrency(void) {
92 if (atomic_capacity())
93 return 1;
94 for (int n = 0; n < 4; ++n)
95 if (opposite_pair(PT_SPLICE) || opposite_pair(PT_TEE))
96 return 1;
97 return 0;
98}