diff --git a/CMakeLists.txt b/CMakeLists.txt index 6ce19d67..a86ee2f5 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -377,12 +377,14 @@ if (UA2F_BUILD_TESTS) target_link_libraries(ua2f_handler_test uci) endif () - # Proxy lifecycle tests exercise listener setup, worker selection, and - # startup failure paths without requiring transparent-routing privileges. + # Proxy tests cover lifecycle and real socket/epoll forwarding without + # transparent-routing privileges. The C bridge includes proxy.c to exercise + # private connection state without exposing test hooks in production headers. add_executable( ua2f_proxy_test test/proxy_test.cc - src/proxy.c + test/proxy_transport_test.cc + test/proxy_transport.c src/handler.c src/util.c src/cache.c diff --git a/src/proxy.c b/src/proxy.c index 83f9aea3..87b70701 100644 --- a/src/proxy.c +++ b/src/proxy.c @@ -136,6 +136,7 @@ struct proxy_connection { uint32_t target_armed; }; +// Not registered in epoll, including a paused side with no useful events. #define PROXY_EVENTS_UNSET UINT32_MAX struct proxy_context { @@ -373,7 +374,10 @@ static void proxy_buffer_compact(struct proxy_buffer *buf) { static int epoll_set(int epoll_fd, int op, int fd, struct epoll_ref *ref, uint32_t events) { struct epoll_event event; memset(&event, 0, sizeof(event)); - event.events = events | EPOLLERR | EPOLLHUP | EPOLLRDHUP; + event.events = events | EPOLLERR | EPOLLHUP; + if (events & EPOLLIN) { + event.events |= EPOLLRDHUP; + } event.data.ptr = ref; return epoll_ctl(epoll_fd, op, fd, &event); } @@ -381,10 +385,21 @@ static int epoll_set(int epoll_fd, int op, int fd, struct epoll_ref *ref, uint32 // Re-arm an fd only when the desired interest mask differs from what is armed, // avoiding an EPOLL_CTL_MOD syscall on every event in steady state. static int epoll_rearm(int epoll_fd, int fd, struct epoll_ref *ref, uint32_t *armed, uint32_t desired) { + if (desired == 0) { + // HUP is reported even with an empty interest mask. Remove a blocked + // or fully drained side until the other side makes progress, rather + // than spinning on HUP while there is nowhere to forward its data. + if (*armed != PROXY_EVENTS_UNSET && epoll_ctl(epoll_fd, EPOLL_CTL_DEL, fd, NULL) != 0) { + return -1; + } + *armed = PROXY_EVENTS_UNSET; + return 0; + } if (*armed == desired) { return 0; } - if (epoll_set(epoll_fd, EPOLL_CTL_MOD, fd, ref, desired) != 0) { + const int op = *armed == PROXY_EVENTS_UNSET ? EPOLL_CTL_ADD : EPOLL_CTL_MOD; + if (epoll_set(epoll_fd, op, fd, ref, desired) != 0) { return -1; } *armed = desired; @@ -802,7 +817,7 @@ static void handle_connection_event(struct epoll_ref *ref, uint32_t events) { return; } - if (events & EPOLLOUT) { + if (events & (EPOLLOUT | EPOLLHUP)) { if (side == PROXY_SIDE_CLIENT) { if (flush_buffer(conn->client_fd, &conn->target_to_client) != 0 || pump_splice_to_client(conn) != 0) { connection_schedule_close(conn); @@ -814,7 +829,10 @@ static void handle_connection_event(struct epoll_ref *ref, uint32_t events) { } } - if (events & (EPOLLIN | EPOLLRDHUP)) { + // HUP/RDHUP may arrive with unread bytes. Only recv/splice returning zero + // establishes EOF; a full buffer or pending pipe must resume after flushing. + const bool read_eof = side == PROXY_SIDE_CLIENT ? conn->client_eof : conn->target_eof; + if (!read_eof && (events & (EPOLLIN | EPOLLRDHUP | EPOLLHUP))) { const int read_result = side == PROXY_SIDE_TARGET ? transfer_target_to_client(conn) : read_into_buffer(conn, side); if (read_result != 0) { @@ -835,11 +853,13 @@ static void handle_connection_event(struct epoll_ref *ref, uint32_t events) { } } - if ((events & (EPOLLERR | EPOLLHUP)) != 0) { - if (side == PROXY_SIDE_CLIENT) { - conn->client_eof = true; - } else { - conn->target_eof = true; + if (events & EPOLLERR) { + const int fd = side == PROXY_SIDE_CLIENT ? conn->client_fd : conn->target_fd; + int error = 0; + socklen_t error_len = sizeof(error); + if (getsockopt(fd, SOL_SOCKET, SO_ERROR, &error, &error_len) != 0 || error != 0) { + connection_schedule_close(conn); + return; } } diff --git a/test/proxy_transport.c b/test/proxy_transport.c new file mode 100644 index 00000000..8196a45b --- /dev/null +++ b/test/proxy_transport.c @@ -0,0 +1,86 @@ +// Compile the production proxy here for transport tests. This keeps its private +// event-loop API out of the installed headers, while testing the exact same C +// implementation (including lifecycle tests that call run_proxy). +#include "../src/proxy.c" +#include "proxy_transport.h" + +struct proxy_test_connection { + struct proxy_context ctx; + struct proxy_connection *conn; +}; + +struct proxy_test_connection *proxy_test_create(int client_fd, int target_fd, bool splice_enabled, bool target_in_progress) { + struct proxy_test_connection *test = calloc(1, sizeof(*test)); + if (test == NULL) { + close(client_fd); + close(target_fd); + return NULL; + } + test->ctx.epoll_fd = epoll_create1(EPOLL_CLOEXEC); + if (test->ctx.epoll_fd < 0 || !proxy_try_acquire_connection()) { + if (test->ctx.epoll_fd >= 0) { + close(test->ctx.epoll_fd); + } + close(client_fd); + close(target_fd); + free(test); + return NULL; + } + test->conn = create_connection(&test->ctx, client_fd, target_fd, AF_INET, target_in_progress); + if (test->conn == NULL) { + close_all_connections(&test->ctx); + close(test->ctx.epoll_fd); + free(test); + return NULL; + } + test->conn->rewrite_disabled = true; + if (!splice_enabled) { + test->conn->splice_enabled = false; + } else if (!test->conn->splice_enabled || fcntl(test->conn->splice_pipe[0], F_SETPIPE_SZ, 4096) < 0) { + proxy_test_destroy(test); + return NULL; + } + return test; +} + +void proxy_test_destroy(struct proxy_test_connection *test) { + if (test != NULL) { + close_all_connections(&test->ctx); + close(test->ctx.epoll_fd); + free(test); + } +} + +int proxy_test_step(struct proxy_test_connection *test, int timeout_ms) { + struct epoll_event events[16]; + const int ready = epoll_wait(test->ctx.epoll_fd, events, 16, timeout_ms); + if (ready < 0) { + return -1; + } + for (int i = 0; i < ready; i++) { + handle_connection_event(events[i].data.ptr, events[i].events); + } + // Retain closing connections until test destruction so assertions can inspect + // EOF and pending-byte state after the descriptors have been closed. + return ready; +} + +struct proxy_test_state proxy_test_snapshot(const struct proxy_test_connection *test) { + const struct proxy_connection *conn = test->conn; + return (struct proxy_test_state){ + .client_eof = conn->client_eof, + .target_eof = conn->target_eof, + .client_write_shutdown = conn->client_write_shutdown, + .target_write_shutdown = conn->target_write_shutdown, + .closing = conn->closing, + .splice_enabled = conn->splice_enabled, + .target_connected = conn->target_connected, + .client_registered = !conn->closing && conn->client_armed != PROXY_EVENTS_UNSET, + .target_registered = !conn->closing && conn->target_armed != PROXY_EVENTS_UNSET, + .client_events = conn->client_armed, + .target_events = conn->target_armed, + .request_pending = conn->client_to_target.len - conn->client_to_target.off, + .response_pending = conn->target_to_client.len - conn->target_to_client.off, + .splice_pending = conn->splice_pending, + }; +} diff --git a/test/proxy_transport.h b/test/proxy_transport.h new file mode 100644 index 00000000..bc787304 --- /dev/null +++ b/test/proxy_transport.h @@ -0,0 +1,24 @@ +#ifndef UA2F_PROXY_TRANSPORT_TEST_H +#define UA2F_PROXY_TRANSPORT_TEST_H + +#include +#include +#include + +struct proxy_test_connection; +struct proxy_test_state { + bool client_eof, target_eof; + bool client_write_shutdown, target_write_shutdown; + bool closing, splice_enabled, target_connected; + bool client_registered, target_registered; + uint32_t client_events, target_events; + size_t request_pending, response_pending, splice_pending; +}; + +// Takes ownership of both proxy-side descriptors, including on failure. +struct proxy_test_connection *proxy_test_create(int client_fd, int target_fd, bool splice_enabled, bool target_in_progress); +void proxy_test_destroy(struct proxy_test_connection *test); +int proxy_test_step(struct proxy_test_connection *test, int timeout_ms); +struct proxy_test_state proxy_test_snapshot(const struct proxy_test_connection *test); + +#endif diff --git a/test/proxy_transport_test.cc b/test/proxy_transport_test.cc new file mode 100644 index 00000000..f1aca7a4 --- /dev/null +++ b/test/proxy_transport_test.cc @@ -0,0 +1,322 @@ +#include + +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +extern "C" { +#include "proxy_transport.h" +} + +namespace { + +struct Transport { + bool tcp; + bool splice; +}; + +class Socket { +public: + ~Socket() { if (fd >= 0) close(fd); } + int Release() { const int result = fd; fd = -1; return result; } + int fd = -1; +}; + +using Clock = std::chrono::steady_clock; +using Bytes = std::vector; + +Bytes Payload(size_t size, uint32_t seed) { + Bytes result(size); + for (auto &byte : result) { + seed = seed * 1664525U + 1013904223U; + byte = static_cast(seed >> 24); + } + return result; +} + +uint64_t Checksum(const Bytes &bytes) { + uint64_t result = UINT64_C(14695981039346656037); + for (uint8_t byte : bytes) result = (result ^ byte) * UINT64_C(1099511628211); + return result; +} + +class ProxyTransportTest : public testing::TestWithParam { +protected: + void Pair(Socket *peer, Socket *proxy) { + if (!GetParam().tcp) { + int pair[2]; + ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, pair), 0); + peer->fd = pair[0]; proxy->fd = pair[1]; + } else { + Socket listener; + listener.fd = socket(AF_INET, SOCK_STREAM, 0); + ASSERT_GE(listener.fd, 0); + sockaddr_in address{}; + address.sin_family = AF_INET; + address.sin_addr.s_addr = htonl(INADDR_LOOPBACK); + ASSERT_EQ(bind(listener.fd, reinterpret_cast(&address), sizeof(address)), 0); + ASSERT_EQ(listen(listener.fd, 1), 0); + socklen_t length = sizeof(address); + ASSERT_EQ(getsockname(listener.fd, reinterpret_cast(&address), &length), 0); + peer->fd = socket(AF_INET, SOCK_STREAM, 0); + ASSERT_GE(peer->fd, 0); + ASSERT_EQ(connect(peer->fd, reinterpret_cast(&address), length), 0); + proxy->fd = accept(listener.fd, nullptr, nullptr); + ASSERT_GE(proxy->fd, 0); + const int one = 1; + ASSERT_EQ(setsockopt(peer->fd, IPPROTO_TCP, TCP_NODELAY, &one, sizeof(one)), 0); + ASSERT_EQ(setsockopt(proxy->fd, IPPROTO_TCP, TCP_NODELAY, &one, sizeof(one)), 0); + } + for (int fd : {peer->fd, proxy->fd}) { + const int flags = fcntl(fd, F_GETFL); + ASSERT_GE(flags, 0); + ASSERT_EQ(fcntl(fd, F_SETFL, flags | O_NONBLOCK), 0); + } + } + + void Start(bool small_outputs = false, bool target_in_progress = false) { + Socket client_proxy, target_proxy; + ASSERT_NO_FATAL_FAILURE(Pair(&client_, &client_proxy)); + ASSERT_NO_FATAL_FAILURE(Pair(&target_, &target_proxy)); + if (small_outputs) { + const int small = 4096; + for (int fd : {client_proxy.fd, target_proxy.fd}) { + ASSERT_EQ(setsockopt(fd, SOL_SOCKET, SO_SNDBUF, &small, sizeof(small)), 0); + } + } + test_.reset(proxy_test_create(client_proxy.Release(), target_proxy.Release(), + GetParam().splice, target_in_progress)); + ASSERT_NE(test_, nullptr); + ASSERT_EQ(State().splice_enabled, GetParam().splice); + } + + proxy_test_state State() const { return proxy_test_snapshot(test_.get()); } + int Pump(int timeout = 1) { + const int events = proxy_test_step(test_.get(), timeout); + EXPECT_GE(events, 0); + return events; + } + void Write(int fd, const Bytes &bytes, size_t *offset) { + if (*offset == bytes.size()) return; + const ssize_t n = send(fd, bytes.data() + *offset, bytes.size() - *offset, MSG_NOSIGNAL); + if (n < 0) ASSERT_TRUE(errno == EAGAIN || errno == EWOULDBLOCK) << strerror(errno); + else { ASSERT_GT(n, 0); *offset += static_cast(n); } + } + void Read(int fd, Bytes *bytes, bool *eof) { + std::array buffer{}; + for (;;) { + const ssize_t n = recv(fd, buffer.data(), buffer.size(), 0); + if (n < 0) { + ASSERT_TRUE(errno == EAGAIN || errno == EWOULDBLOCK) << strerror(errno); + return; + } + if (n == 0) { *eof = true; return; } + bytes->insert(bytes->end(), buffer.begin(), buffer.begin() + n); + } + } + void EqualPayload(const Bytes &expected, const Bytes &received) { + ASSERT_EQ(received.size(), expected.size()); + EXPECT_EQ(Checksum(received), Checksum(expected)); + EXPECT_EQ(memcmp(received.data(), expected.data(), expected.size()), 0); + } + void Drain(int peer, const Bytes &expected) { + Bytes received; + bool eof = false; + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!eof) { + ASSERT_LT(Clock::now(), deadline) << "timeout before EOF"; + Pump(); + ASSERT_NO_FATAL_FAILURE(Read(peer, &received, &eof)); + } + ASSERT_NO_FATAL_FAILURE(EqualPayload(expected, received)); + } + + Socket client_, target_; + std::unique_ptr test_{nullptr, proxy_test_destroy}; +}; + +TEST_P(ProxyTransportTest, DrainsClientBytesAfterTargetHalfClose) { + ASSERT_NO_FATAL_FAILURE(Start()); + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!State().target_eof) { ASSERT_LT(Clock::now(), deadline); Pump(); } + ASSERT_TRUE(State().client_write_shutdown); + ASSERT_FALSE(State().client_eof); + EXPECT_EQ(Pump(0), 0) << "a drained read side must not spin while the other side remains open"; + // Regression: master delivered only 16,384 of these 65,536 accepted bytes. + const Bytes request = Payload(65536, 17); + ASSERT_EQ(send(client_.fd, request.data(), request.size(), MSG_NOSIGNAL), static_cast(request.size())); + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); + ASSERT_NO_FATAL_FAILURE(Drain(target_.fd, request)); + EXPECT_TRUE(State().closing); + EXPECT_TRUE(State().client_eof); +} + +TEST_P(ProxyTransportTest, DrainsTargetBytesAfterClientHalfClose) { + ASSERT_NO_FATAL_FAILURE(Start()); + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!State().client_eof) { ASSERT_LT(Clock::now(), deadline); Pump(); } + ASSERT_TRUE(State().target_write_shutdown); + EXPECT_EQ(Pump(0), 0) << "a drained read side must not spin while the other side remains open"; + const Bytes response = Payload(65536, 71); + ASSERT_EQ(send(target_.fd, response.data(), response.size(), MSG_NOSIGNAL), static_cast(response.size())); + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); + ASSERT_NO_FATAL_FAILURE(Drain(client_.fd, response)); + EXPECT_TRUE(State().closing); + EXPECT_TRUE(State().target_eof); +} + +TEST_P(ProxyTransportTest, RegistersClientAfterTargetConnectThenDrainsHalfClosedInput) { + // An already-connected TCP/socketpair endpoint makes connect completion + // deterministic while exercising the production deferred-registration path. + ASSERT_NO_FATAL_FAILURE(Start(false, true)); + ASSERT_FALSE(State().target_connected); + ASSERT_FALSE(State().client_registered); + ASSERT_TRUE(State().target_registered); + ASSERT_EQ(State().target_events, static_cast(EPOLLOUT)); + const Bytes request = Payload(65536, 101); + ASSERT_EQ(send(client_.fd, request.data(), request.size(), MSG_NOSIGNAL), static_cast(request.size())); + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); + ASSERT_NO_FATAL_FAILURE(Drain(target_.fd, request)); + EXPECT_TRUE(State().target_connected); + EXPECT_TRUE(State().closing); +} + +TEST_P(ProxyTransportTest, DrainsPayloadsWhenBothPeersHalfCloseBeforeAnyDispatch) { + ASSERT_NO_FATAL_FAILURE(Start(true)); + const Bytes request = Payload(65536, 113), response = Payload(65536, 127); + ASSERT_EQ(send(client_.fd, request.data(), request.size(), MSG_NOSIGNAL), static_cast(request.size())); + ASSERT_EQ(send(target_.fd, response.data(), response.size(), MSG_NOSIGNAL), static_cast(response.size())); + // Both FINs are queued before the proxy processes either direction. + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); + Bytes received_request, received_response; + bool request_eof = false, response_eof = false; + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!request_eof || !response_eof) { + ASSERT_LT(Clock::now(), deadline); + Pump(); + ASSERT_NO_FATAL_FAILURE(Read(target_.fd, &received_request, &request_eof)); + ASSERT_NO_FATAL_FAILURE(Read(client_.fd, &received_response, &response_eof)); + } + ASSERT_NO_FATAL_FAILURE(EqualPayload(request, received_request)); + ASSERT_NO_FATAL_FAILURE(EqualPayload(response, received_response)); + EXPECT_TRUE(State().closing); +} + +TEST_P(ProxyTransportTest, DrainsBothDirectionsAfterSimultaneousHalfCloseWithBackpressure) { + ASSERT_NO_FATAL_FAILURE(Start(true)); + const Bytes request = Payload(1048613, 131), response = Payload(1048627, 197); + Bytes received_request, received_response; + size_t request_sent = 0, response_sent = 0; + bool request_shutdown = false, response_shutdown = false; + bool request_eof = false, response_eof = false; + bool request_blocked = false, response_blocked = false; + const auto deadline = Clock::now() + std::chrono::seconds(15); + while (!request_eof || !response_eof) { + ASSERT_LT(Clock::now(), deadline) << "duplex transfer stalled"; + ASSERT_NO_FATAL_FAILURE(Write(client_.fd, request, &request_sent)); + ASSERT_NO_FATAL_FAILURE(Write(target_.fd, response, &response_sent)); + if (request_sent == request.size() && !request_shutdown) { + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); request_shutdown = true; + } + if (response_sent == response.size() && !response_shutdown) { + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); response_shutdown = true; + } + Pump(); + const auto state = State(); + request_blocked |= state.request_pending > 0; + response_blocked |= state.response_pending > 0 || state.splice_pending > 0; + if (state.request_pending > 0) EXPECT_FALSE(state.target_write_shutdown); + if (state.response_pending > 0 || state.splice_pending > 0) EXPECT_FALSE(state.client_write_shutdown); + // Pause both applications until both forwarding directions are blocked. + if (request_blocked && response_blocked) { + ASSERT_NO_FATAL_FAILURE(Read(target_.fd, &received_request, &request_eof)); + ASSERT_NO_FATAL_FAILURE(Read(client_.fd, &received_response, &response_eof)); + } + } + EXPECT_TRUE(request_blocked); + EXPECT_TRUE(response_blocked); + ASSERT_NO_FATAL_FAILURE(EqualPayload(request, received_request)); + ASSERT_NO_FATAL_FAILURE(EqualPayload(response, received_response)); + EXPECT_TRUE(State().closing); + EXPECT_EQ(State().request_pending, 0U); + EXPECT_EQ(State().response_pending, 0U); + EXPECT_EQ(State().splice_pending, 0U); +} + +// UNIX sockets make send-buffer backpressure deterministic, without relying on +// TCP delayed ACK timing. The TCP variants above exercise the real HUP flags. +TEST_P(ProxyTransportTest, PausesHungUpBlockedInputAndReaddsItAfterOutputProgress) { + if (GetParam().tcp) GTEST_SKIP() << "deterministic socket-buffer pressure uses UNIX sockets"; + ASSERT_NO_FATAL_FAILURE(Start(true)); + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!State().target_eof) { ASSERT_LT(Clock::now(), deadline); Pump(); } + const Bytes request = Payload(65536, 263); + ASSERT_EQ(send(client_.fd, request.data(), request.size(), MSG_NOSIGNAL), static_cast(request.size())); + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); + for (int i = 0; i < 100 && State().client_registered; i++) Pump(); + ASSERT_GT(State().request_pending, 0U); + ASSERT_FALSE(State().client_eof); + ASSERT_FALSE(State().client_registered); + ASSERT_FALSE(State().target_write_shutdown); + EXPECT_EQ(Pump(0), 0) << "HUP must not repeatedly wake a blocked source"; + ASSERT_NO_FATAL_FAILURE(Drain(target_.fd, request)); + EXPECT_TRUE(State().closing); +} + +TEST_P(ProxyTransportTest, DoesNotPollReadHalfCloseWhileResponseOutputIsBlocked) { + if (GetParam().tcp) GTEST_SKIP() << "deterministic socket-buffer pressure uses UNIX sockets"; + ASSERT_NO_FATAL_FAILURE(Start(true)); + ASSERT_EQ(shutdown(client_.fd, SHUT_WR), 0); + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!State().client_eof) { ASSERT_LT(Clock::now(), deadline); Pump(); } + const Bytes response = Payload(65536, 331); + ASSERT_EQ(send(target_.fd, response.data(), response.size(), MSG_NOSIGNAL), static_cast(response.size())); + ASSERT_EQ(shutdown(target_.fd, SHUT_WR), 0); + for (int i = 0; i < 100; i++) { if (Pump(0) == 0) break; } + const auto state = State(); + ASSERT_GT(state.response_pending + state.splice_pending, 0U); + ASSERT_FALSE(state.client_write_shutdown); + ASSERT_TRUE(state.client_registered); + EXPECT_EQ(state.client_events, static_cast(EPOLLOUT)); + EXPECT_EQ(Pump(0), 0) << "RDHUP must not wake a read-complete socket waiting to write"; + ASSERT_NO_FATAL_FAILURE(Drain(client_.fd, response)); + EXPECT_TRUE(State().closing); +} + +TEST_P(ProxyTransportTest, ClosesCleanlyOnTargetReset) { + if (!GetParam().tcp) GTEST_SKIP() << "TCP reset semantics require TCP sockets"; + ASSERT_NO_FATAL_FAILURE(Start(true)); + const linger reset{1, 0}; + ASSERT_EQ(setsockopt(target_.fd, SOL_SOCKET, SO_LINGER, &reset, sizeof(reset)), 0); + ASSERT_EQ(close(target_.Release()), 0); + const auto deadline = Clock::now() + std::chrono::seconds(10); + while (!State().closing) { ASSERT_LT(Clock::now(), deadline); Pump(); } + EXPECT_EQ(Pump(0), 0); +} + +INSTANTIATE_TEST_SUITE_P(BufferedAndSplice, ProxyTransportTest, + testing::Values(Transport{true, false}, Transport{true, true}, + Transport{false, false}, Transport{false, true}), + [](const testing::TestParamInfo &info) { + return std::string(info.param.tcp ? "Tcp" : "Unix") + + (info.param.splice ? "Splice" : "Buffered"); + }); + +} // namespace