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));
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);
28 for (ssize_t n = 0; n < count; ++n) {
29 if (bytes[n] !=
'a' && bytes[n] !=
'b') {
33 ++totals[bytes[n] -
'a'];
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}};
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]));
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);
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) {
70 CHECK(errno == EAGAIN);
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));
76 pt_diagnostic(&calls[0]);
77 pt_diagnostic(&calls[1]);
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) {
91int pipe_transfer_concurrency(
void) {
92 if (atomic_capacity())
94 for (
int n = 0; n < 4; ++n)
95 if (opposite_pair(PT_SPLICE) || opposite_pair(PT_TEE))