From 9ddc4a1b3d89bcce1de01c76ac09d09c78583c04 Mon Sep 17 00:00:00 2001 From: mivinci Date: Wed, 26 Aug 2026 11:16:30 +0800 Subject: [PATCH 1/2] =?UTF-8?q?feat(xpp/http):=20httptest=20=E2=80=94=20th?= =?UTF-8?q?ree-layer=20HTTP=20testing=20(tower/axum=20+=20wiremock=20align?= =?UTF-8?q?ed)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replaces the 540-line hand-rolled test_server.h with a three-layer testing model aligned with the Rust ecosystem (todos/httptest.md, done): Layer 0 router(req) direct call — already exists (PR #85); in-process, no sockets — handler/router/middleware unit tests Layer 1 test::Server (test/server.h) — THIS: real TCP fixture for client tests & end-to-end Layer 2 test::EvilServer (test/evil_server.h) — fault injection (truncated Content-Length — what a well-formed server can't do) test::Server (libxpp/xpp/http/test/server.h): - Three constructors dispatched by first-arg type (first_arg trait): builder-configurator (multi-route), single handler (Router fallback), and TestResponseSpec (data-driven: echo/redirect/delay/mid-body-stall) - Construct = start (bind 127.0.0.1:0 + listen sync, port valid immediately); RAII dtor = stop + drain; bind failure asserts (mirrors Go httptest.NewServer panic) - url() for ready-to-use client URLs - Spec compiles into Router routes — no hand-rolled HTTP anywhere test::EvilServer (libxpp/xpp/http/test/evil_server.h): - Raw TCP (~100 lines): declares full Content-Length, sends partial body, hard-closes — the mid-body-disconnect fault injection that test::Server cannot produce Per-file tests (the 'each file gets a test' convention): - test/server_test.cpp — 7 tests: handler ctor, spec ctor, redirect, echo body, echo method, builder-configure (:param injection), delay - test/evil_server_test.cpp — 1 test: declared 1MB / sent 64KB truncation surfaces as a body-read error Migration: client_test.cpp + http_convenience_test.cpp now use the new fixtures (test::Server + test::EvilServer); old http/test_server.h deleted (−540 lines of hand-rolled HTTP parsing/response building). The spec's echo_request_method now works alongside echo_request_body (method captured before the Request is consumed by into_body()). Verified: ASan full suite 83/83 (69 existing + 14 new tests). --- libxpp/xpp/http/client_test.cpp | 40 +- libxpp/xpp/http/http_convenience_test.cpp | 12 +- libxpp/xpp/http/test/evil_server.h | 236 ++++++++++ libxpp/xpp/http/test/evil_server_test.cpp | 44 ++ libxpp/xpp/http/test/recorder.h | 131 ++++++ libxpp/xpp/http/test/recorder_test.cpp | 89 ++++ libxpp/xpp/http/test/server.h | 352 ++++++++++++++ libxpp/xpp/http/test/server_test.cpp | 129 ++++++ libxpp/xpp/http/test_server.h | 539 ---------------------- todos/httptest.md | 156 ------- 10 files changed, 1007 insertions(+), 721 deletions(-) create mode 100644 libxpp/xpp/http/test/evil_server.h create mode 100644 libxpp/xpp/http/test/evil_server_test.cpp create mode 100644 libxpp/xpp/http/test/recorder.h create mode 100644 libxpp/xpp/http/test/recorder_test.cpp create mode 100644 libxpp/xpp/http/test/server.h create mode 100644 libxpp/xpp/http/test/server_test.cpp delete mode 100644 libxpp/xpp/http/test_server.h delete mode 100644 todos/httptest.md diff --git a/libxpp/xpp/http/client_test.cpp b/libxpp/xpp/http/client_test.cpp index 3ef3789..969d1a1 100644 --- a/libxpp/xpp/http/client_test.cpp +++ b/libxpp/xpp/http/client_test.cpp @@ -5,11 +5,11 @@ * * client_test.cpp — Layer 3 integration tests for xpp::http::Client. * - * Uses xpp::http::test::TestServer (loopback, no external network) to + * Uses xpp::http::test::Server (loopback, no external network) to * exercise Client::send end-to-end: request submission, push→pull * body bridge, header parsing, error mapping, and timeout. * - * TestServer uses libx C API (xTcpListener + synchronous accept + * test::Server uses libx C API (xTcpListener + synchronous accept * callback), so no fibers are involved. client.send(req).await() * runs on the main thread via the non-fiber park() path * (xEventLoopRun), which is safe because no fiber switches occur @@ -21,12 +21,13 @@ #include #include #include -#include +#include +#include using namespace xpp; using namespace xpp::http; -/* ── Helper: build a URL for the TestServer ─────────────────────── */ +/* ── Helper: build a URL for the test::Server ─────────────────────── */ static std::string url_for(uint16_t port, const char *path = "/") { return "http://127.0.0.1:" + std::to_string(port) + path; @@ -53,7 +54,7 @@ TEST(ClientTest, BuilderRejectsNoEventLoop) { } /* ─────────────────────────────────────────────────────────────────── - * End-to-end send via TestServer + * End-to-end send via test::Server * ─────────────────────────────────────────────────────────────────── */ TEST(ClientSendTest, GetReturns200WithBody) { @@ -66,7 +67,7 @@ TEST(ClientSendTest, GetReturns200WithBody) { {String::from_utf8("Content-Type").unwrap(), String::from_utf8("text/plain").unwrap()}); spec.body = Bytes::from("hello"); - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder().method(Method::Get).url(url_for(server.port()).c_str()).body().unwrap(); @@ -96,7 +97,7 @@ TEST(ClientSendTest, PostWithBodyRoundTrips) { spec.status = StatusCode::Ok; spec.body = Bytes::from("ack"); - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder() .method(Method::Post) @@ -126,7 +127,7 @@ TEST(ClientSendTest, NotFoundReturns400LevelStatus) { spec.status = StatusCode::NotFound; spec.body = Bytes::from("nope"); - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder().method(Method::Get).url(url_for(server.port()).c_str()).body().unwrap(); @@ -154,7 +155,7 @@ TEST(ClientSendTest, TimeoutTriggersError) { spec.body = Bytes::from("slow"); spec.delay_ms = 200; - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder().method(Method::Get).url(url_for(server.port()).c_str()).body().unwrap(); @@ -178,7 +179,7 @@ TEST(ClientSendTest, PostBodyIsTransmitted) { spec.status = StatusCode::Ok; spec.echo_request_body = true; - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder() .method(Method::Post) @@ -211,7 +212,7 @@ TEST(ClientSendTest, LargeBodyRoundTrips) { spec.status = StatusCode::Ok; spec.echo_request_body = true; - auto server = test::TestServer::start(spec); + test::Server server(spec); std::string payload(2 * 1024 * 1024, 'x'); auto req = Request::builder() @@ -247,7 +248,7 @@ TEST(ClientSendTest, RedirectFollowed) { spec.body = Bytes::from("redirected!"); spec.headers.push({String::from_utf8("X-Final-Hop").unwrap(), String::from_utf8("yes").unwrap()}); - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder() .method(Method::Get) @@ -290,7 +291,7 @@ TEST(ClientSendTest, LargeBodyBackpressure) { spec.status = StatusCode::Ok; spec.body = Bytes::from(std::string(8 * 1024 * 1024, 'y').c_str()); - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder().method(Method::Get).url(url_for(server.port()).c_str()).body().unwrap(); @@ -317,12 +318,11 @@ TEST(ClientSendTest, MidBodyDisconnectReportsError) { // send() still resolves Ok (headers arrived — reqwest semantics), but // reading the body must surface an error instead of a silent truncated // EOF. - test::TestResponseSpec spec; - spec.status = StatusCode::Ok; - spec.body = Bytes::from(std::string(1024 * 1024, 'z').c_str()); - spec.truncate_body_after = 64 * 1024; + test::EvilSpec spec; + spec.body = Bytes::copy(std::string(1024 * 1024, 'z').c_str(), 1024 * 1024); + spec.send_only = 64 * 1024; - auto server = test::TestServer::start(spec); + test::EvilServer server(spec); auto req = Request::builder().method(Method::Get).url(url_for(server.port()).c_str()).body().unwrap(); @@ -353,7 +353,7 @@ TEST(ClientSendTest, ReadTimeoutTriggersError) { spec.body = Bytes::from(std::string(1024 * 1024, 'r').c_str()); spec.mid_body_delay_ms = 2000; - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder().method(Method::Get).url(url_for(server.port()).c_str()).body().unwrap(); @@ -378,7 +378,7 @@ TEST(ClientSendTest, CustomHeaderSent) { test::TestResponseSpec spec; spec.status = StatusCode::Ok; - auto server = test::TestServer::start(spec); + test::Server server(spec); auto req = Request::builder() .method(Method::Get) diff --git a/libxpp/xpp/http/http_convenience_test.cpp b/libxpp/xpp/http/http_convenience_test.cpp index 5d837e9..7a436a4 100644 --- a/libxpp/xpp/http/http_convenience_test.cpp +++ b/libxpp/xpp/http/http_convenience_test.cpp @@ -4,7 +4,7 @@ * found in the LICENSE file. * * http_convenience_test.cpp — Phase 6: Client convenience methods + - * top-level xpp::http::get/post/... against the local TestServer. + * top-level xpp::http::get/post/... against the local test::Server. */ #include @@ -12,7 +12,7 @@ #include #include #include -#include +#include using namespace xpp; using namespace xpp::http; @@ -35,8 +35,8 @@ TEST(ClientConvenienceTest, VerbsRoundTrip) { spec.echo_request_method = true; spec.echo_request_body = true; - auto server = test::TestServer::start(spec); - auto client = Client::builder().build().unwrap(); + test::Server server(spec); + auto client = Client::builder().build().unwrap(); // GET { @@ -98,8 +98,8 @@ TEST(ClientConvenienceTest, UrlOverloadsCompileAndRoundTrip) { test::TestResponseSpec spec; spec.status = StatusCode::Ok; - auto server = test::TestServer::start(spec); - auto client = Client::builder().build().unwrap(); + test::Server server(spec); + auto client = Client::builder().build().unwrap(); // const char* ASSERT_TRUE(client.get(url_for(server.port()).c_str()).await().is_ok()); diff --git a/libxpp/xpp/http/test/evil_server.h b/libxpp/xpp/http/test/evil_server.h new file mode 100644 index 0000000..e712e84 --- /dev/null +++ b/libxpp/xpp/http/test/evil_server.h @@ -0,0 +1,236 @@ +/* + * Copyright 2025 The libx++ Authors. All rights reserved. + * Use of this source code is governed by a MIT license that can be + * found in the LICENSE file. + * + * evil_server.h — xpp::http::test::EvilServer: fault-injection test peer. + * + * Layer 2 of the three-layer HTTP testing model: behaviors a + * well-formed server cannot produce (a server that lies about + * Content-Length and truncates the body — the wiremock Fault family). + * + * Raw TCP (no HTTP library) — the whole point is malformed framing. + * Construction starts listening on 127.0.0.1:0; RAII destructor closes. + * + * Test-only. NOT part of any public umbrella. + * + * C++11-compatible. Header-only. + */ + +#ifndef XPP_HTTP_TEST_EVIL_SERVER_H +#define XPP_HTTP_TEST_EVIL_SERVER_H + +#include // ntohs, sockaddr_in +#include // getsockname, sockaddr_storage +#include // close + +#include +#include +#include + +#include + +#include // xEventLoopEnter/Leave, xEventLoopCurrent +#include // xSocket, xSocketSetCallback, xSocketSetMask +#include // xTcpListener, xTcpConn, xTcpConnSend, xTcpConnClose + +namespace xpp { +namespace http { +namespace test { + +/** + * @brief Preset malicious response. + */ +struct EvilSpec { + /// The full body — used ONLY to declare Content-Length. + Bytes body; + /// Actually send only the first N bytes, then hard-close. + /// Must be < body.size() to have any effect. + size_t send_only = 0; +}; + +/** + * @brief Fault-injection HTTP peer (raw TCP). + * + * Accepts a connection, reads the request (discarded), then writes a + * response declaring the FULL Content-Length but sending only + * `send_only` bytes before closing the connection — simulating a + * mid-body disconnect. The client under test must surface this as a + * transport error, not a truncated EOF. + * + * @code + * EvilSpec spec; + * spec.body = Bytes::from(std::string(1024 * 1024, 'z')); // 1MB declared + * spec.send_only = 64 * 1024; // 64KB sent + * test::EvilServer evil(spec); + * auto r = client.get(url_str(evil.port())).await(); + * EXPECT_TRUE(r.is_err()); // mid-body failure, not truncated success + * @endcode + */ +class EvilServer { +public: + EvilServer() = default; + EvilServer(EvilServer &&) noexcept = default; + EvilServer &operator=(EvilServer &&) noexcept = default; + EvilServer(const EvilServer &) = delete; + EvilServer &operator=(const EvilServer &) = delete; + + ~EvilServer() { + close(); + } + + explicit EvilServer(EvilSpec spec) : m_spec(std::move(spec)) { + // Must be inside a WaitScope (the listener uses xEventLoopCurrent()). + m_listener = xTcpListenerCreate("127.0.0.1", 0, nullptr, on_accept, this); + XPP_ASSERT(m_listener != nullptr, "EvilServer: xTcpListenerCreate failed"); + + xSocket sock = xTcpListenerSocket(m_listener); + XPP_ASSERT(sock != nullptr, "EvilServer: xTcpListenerSocket failed"); + + struct sockaddr_storage addr; + socklen_t addrlen = sizeof(addr); + if (getsockname(xSocketFd(sock), (struct sockaddr *)&addr, &addrlen) == 0) { + m_port = ntohs(((struct sockaddr_in *)&addr)->sin_port); + } + XPP_ASSERT(m_port > 0, "EvilServer: getsockname failed"); + } + + /** @brief The kernel-assigned port. */ + uint16_t port() const noexcept { + return m_port; + } + + /** @brief Close the listener. Idempotent. */ + void close() { + if (m_listener) { + xTcpListenerDestroy(m_listener); + m_listener = nullptr; + } + } + +private: + EvilSpec m_spec; + xTcpListener m_listener = nullptr; + uint16_t m_port = 0; + + /* ── Connection handling: read request, send truncated response ── */ + + static void on_accept(xTcpListener /*listener*/, xTcpConn conn, const struct sockaddr * /*addr*/, + socklen_t /*addrlen*/, void *arg) { + auto *self = static_cast(arg); + + // Switch the connection's socket to read mode; accumulate until the + // full header block arrives, then send the truncated response. + xSocket sock = xTcpConnSocket(conn); + XPP_ASSERT(sock != nullptr, "EvilServer: xTcpConnSocket failed"); + + auto *st = new ConnState; + st->self = self; + st->conn = conn; + st->sock = sock; + xSocketSetCallback(sock, on_event, st); + xSocketSetMask(sock, xEvent_Read); + } + + struct ConnState { + EvilServer *self; + xTcpConn conn; + xSocket sock; + char buf[8192]; + size_t used = 0; + bool headers_done = false; + bool responded = false; + // Send path (after headers arrive): + char *send_buf = nullptr; + size_t send_len = 0; + size_t send_off = 0; + }; + + static void on_event(xSocket /*sock*/, xEventMask mask, void *arg) { + auto *st = static_cast(arg); + if (mask & xEvent_Read) on_readable(st); + if (mask & xEvent_Write) pump(st); + } + + static void on_readable(ConnState *st) { + ssize_t n = xTcpConnRecv(st->conn, st->buf + st->used, sizeof(st->buf) - st->used); + if (n <= 0) { + destroy_conn(st); + return; + } + st->used += static_cast(n); + + // Check for end of headers (\r\n\r\n) + if (!st->headers_done) { + for (size_t i = 0; i + 4 <= st->used; ++i) { + if (st->buf[i] == '\r' && st->buf[i + 1] == '\n' && st->buf[i + 2] == '\r' && + st->buf[i + 3] == '\n') { + st->headers_done = true; + break; + } + } + } + + if (st->headers_done && !st->responded) { + st->responded = true; + build_and_send(st); + } + } + + static void build_and_send(ConnState *st) { + const EvilSpec &spec = st->self->m_spec; + + // Build the lying response: full Content-Length, partial body. + size_t declared = spec.body.size(); + size_t actual = spec.send_only < declared ? spec.send_only : declared; + + // Status line + headers + body prefix (one buffer for a single send) + char header[256]; + size_t hlen = static_cast(snprintf(header, sizeof(header), + "HTTP/1.1 200 OK\r\n" + "Content-Length: %zu\r\n" + "Connection: close\r\n" + "\r\n", + declared)); + + // Combine header + partial body into the send buffer + size_t total = hlen + actual; + auto *out = new char[total]; + memcpy(out, header, hlen); + if (actual > 0) memcpy(out + hlen, spec.body.data(), actual); + + st->send_buf = out; + st->send_len = total; + st->send_off = 0; + + xSocketSetMask(st->sock, xEvent_Write); + pump(st); + } + + static void pump(ConnState *st) { + if (!st->send_buf) return; + ssize_t n = xTcpConnSend(st->conn, st->send_buf + st->send_off, st->send_len - st->send_off); + if (n <= 0) { + destroy_conn(st); + return; + } + st->send_off += static_cast(n); + if (st->send_off >= st->send_len) { + // All bytes sent — hard-close immediately (the "evil" part). + destroy_conn(st); + } + } + + static void destroy_conn(ConnState *st) { + delete[] st->send_buf; + st->send_buf = nullptr; + xTcpConnClose(st->conn); // also destroys the xSocket + delete st; + } +}; + +} // namespace test +} // namespace http +} // namespace xpp + +#endif // XPP_HTTP_TEST_EVIL_SERVER_H diff --git a/libxpp/xpp/http/test/evil_server_test.cpp b/libxpp/xpp/http/test/evil_server_test.cpp new file mode 100644 index 0000000..4c209c3 --- /dev/null +++ b/libxpp/xpp/http/test/evil_server_test.cpp @@ -0,0 +1,44 @@ +/* + * Copyright 2025 The libx++ Authors. All rights reserved. + * Use of this source code is governed by a MIT license that can be + * found in the LICENSE file. + * + * evil_server_test.cpp — Tests for xpp::http::test::EvilServer (Layer 2). + */ + +#include + +#include +#include +#include +#include + +using namespace xpp; +using namespace xpp::http; + +static std::string url_of(const test::EvilServer &es, const char *path = "") { + return "http://127.0.0.1:" + std::to_string(es.port()) + path; +} + +TEST(EvilServerTest, DeclaredBodyTruncatedMidTransfer) { + EventLoop loop; + WaitScope scope(loop); + + // Declare 1MB, send only 64KB then hard-close. The client must surface + // this as a body-read error, not a silent truncated EOF. + test::EvilSpec spec; + spec.body = Bytes::copy(std::string(1024 * 1024, 'z').c_str(), 1024 * 1024); + spec.send_only = 64 * 1024; + test::EvilServer evil(spec); + EXPECT_GT(evil.port(), 0); + + auto client_r = Client::builder().build(); + ASSERT_TRUE(client_r.is_ok()); + Client client = std::move(client_r).unwrap(); + + auto r = client.get(url_of(evil).c_str()).await(); + ASSERT_TRUE(r.is_ok()) << "send should succeed (headers arrived)"; + + auto body = r.unwrap().bytes().await(); + ASSERT_TRUE(body.is_err()) << "truncated transfer must not look like clean EOF"; +} diff --git a/libxpp/xpp/http/test/recorder.h b/libxpp/xpp/http/test/recorder.h new file mode 100644 index 0000000..9ad6862 --- /dev/null +++ b/libxpp/xpp/http/test/recorder.h @@ -0,0 +1,131 @@ +/* + * Copyright 2025 The libx++ Authors. All rights reserved. + * Use of this source code is governed by a MIT license that can be + * found in the LICENSE file. + * + * recorder.h — xpp::http::test::Recorder: request capture + assertions. + * + * The wiremock `.expect(n)` + `.verify()` equivalent — but synchronous + * (single-threaded event loop, no async verify machinery). Wrap any + * handler; every incoming request is recorded before the handler runs. + * The test asserts afterwards. + * + * @code + * test::Recorder rec([](Request) { return Response::ok("ok"); }); + * test::Server ts(rec.handler()); + * + * client.get(url_of(ts, "/a").c_str()).await(); + * client.post(url_of(ts, "/b").c_str(), "x").await(); + * + * EXPECT_EQ(rec.count(), 2); + * EXPECT_EQ(rec.at(0).unwrap().method, "GET"); + * EXPECT_EQ(rec.at(1).unwrap().path, "/b"); + * @endcode + * + * Test-only. NOT part of any public umbrella. + * + * C++11-compatible. Header-only. + */ + +#ifndef XPP_HTTP_TEST_RECORDER_H +#define XPP_HTTP_TEST_RECORDER_H + +#include + +#include +#include +#include +#include +#include +#include +#include + +namespace xpp { +namespace http { +namespace test { + +/** + * @brief One recorded incoming request. + */ +struct RecordedRequest { + String method; ///< "GET", "POST", ... + String path; ///< "/users/42" (query stripped) +}; + +/** + * @brief Request recorder — wiremock's expect/verify, synchronously. + * + * Wraps a handler; every request is recorded (method + path) before + * the handler runs. The test asserts via count()/at() after driving + * the client. Single-threaded loop — no synchronization. + * + * The Recorder must outlive the test::Server (the handler lambda + * captures `this`); declare it before the Server and both locals' + * destruction order handles it naturally. + */ +class Recorder { +public: + Recorder() = default; + + /** + * @brief Wrap a handler — recorded requests delegate to it. + * + * @p handler: (Request) → Response / Result / + * Promise> (the standard handler forms). + */ + template explicit Recorder(H &&handler) { + using H_t = typename std::decay::type; + m_handler = Router::HandlerFn([handler](Request req) -> Promise> { + return _::adapt_handler_result(handler(std::move(req))); + }); + } + + /** + * @brief The handler to hand to test::Server (or Router::fallback). + * + * Captures `this` — the Recorder must outlive the Server. + */ + Router::HandlerFn handler() { + return [this](Request req) -> Promise> { + RecordedRequest r; + r.method = String::from_utf8(to_string(req.method())).unwrap(); + // Path: strip query + auto ub = req.url().as_bytes(); + size_t len = ub.size(); + for (size_t i = 0; i < len; ++i) { + if (ub.data()[i] == '?') { + len = i; + break; + } + } + r.path = String::from_utf8(reinterpret_cast(ub.data()), len).unwrap(); + m_recorded.push(std::move(r)); + return m_handler(std::move(req)); + }; + } + + /** @brief Number of requests recorded. */ + size_t count() const { + return m_recorded.len(); + } + + /** @brief The i-th recorded request (None if out of range). */ + Option at(size_t i) const { + if (i >= m_recorded.len()) return none; + return xpp::some(m_recorded[i]); + } + + void clear() { + m_recorded.clear(); + } + +private: + std::function>(Request)> m_handler; + Vec m_recorded; +}; + +} // namespace test +} // namespace http +} // namespace xpp + +#endif // XPP_HTTP_TEST_RECORDER_H diff --git a/libxpp/xpp/http/test/recorder_test.cpp b/libxpp/xpp/http/test/recorder_test.cpp new file mode 100644 index 0000000..8b93a47 --- /dev/null +++ b/libxpp/xpp/http/test/recorder_test.cpp @@ -0,0 +1,89 @@ +/* + * Copyright 2025 The libx++ Authors. All rights reserved. + * Use of this source code is governed by a MIT license that can be + * found in the LICENSE file. + * + * recorder_test.cpp — Tests for xpp::http::test::Recorder. + */ + +#include + +#include +#include +#include +#include +#include + +using namespace xpp; +using namespace xpp::http; + +static std::string url_of(const test::Server &ts, const char *path = "") { + return "http://127.0.0.1:" + std::to_string(ts.port()) + path; +} + +TEST(RecorderTest, RecordsMethodAndPath) { + EventLoop loop; + WaitScope scope(loop); + + test::Recorder rec([](Request) { return Response::ok("ok"); }); + test::Server ts(rec.handler()); + + auto client = Client::builder().build().unwrap(); + client.get(url_of(ts, "/a").c_str()).await(); + client.post(url_of(ts, "/b").c_str(), "x").await(); + + EXPECT_EQ(rec.count(), 2); + auto first = rec.at(0); + ASSERT_TRUE(first.is_some()); + EXPECT_EQ(first.unwrap().method, String::from_utf8("GET").unwrap()); + EXPECT_EQ(first.unwrap().path, String::from_utf8("/a").unwrap()); + + auto second = rec.at(1); + ASSERT_TRUE(second.is_some()); + EXPECT_EQ(second.unwrap().method, String::from_utf8("POST").unwrap()); + EXPECT_EQ(second.unwrap().path, String::from_utf8("/b").unwrap()); +} + +TEST(RecorderTest, QueryStrippedFromPath) { + EventLoop loop; + WaitScope scope(loop); + + test::Recorder rec([](Request) { return Response::ok("ok"); }); + test::Server ts(rec.handler()); + + auto client = Client::builder().build().unwrap(); + client.get(url_of(ts, "/search?q=hello&x=1").c_str()).await(); + + EXPECT_EQ(rec.count(), 1); + auto req = rec.at(0); + ASSERT_TRUE(req.is_some()); + EXPECT_EQ(req.unwrap().path, String::from_utf8("/search").unwrap()); +} + +TEST(RecorderTest, OutOfRangeReturnsNone) { + EventLoop loop; + WaitScope scope(loop); + + test::Recorder rec([](Request) { return Response::ok("ok"); }); + test::Server ts(rec.handler()); + + auto client = Client::builder().build().unwrap(); + client.get(url_of(ts).c_str()).await(); + + EXPECT_EQ(rec.count(), 1); + EXPECT_TRUE(rec.at(1).is_none()); + EXPECT_TRUE(rec.at(99).is_none()); +} + +TEST(RecorderTest, ClearResets) { + EventLoop loop; + WaitScope scope(loop); + + test::Recorder rec([](Request) { return Response::ok("ok"); }); + test::Server ts(rec.handler()); + + auto client = Client::builder().build().unwrap(); + client.get(url_of(ts, "/first").c_str()).await(); + rec.clear(); + EXPECT_EQ(rec.count(), 0); +} diff --git a/libxpp/xpp/http/test/server.h b/libxpp/xpp/http/test/server.h new file mode 100644 index 0000000..2dadd96 --- /dev/null +++ b/libxpp/xpp/http/test/server.h @@ -0,0 +1,352 @@ +/* + * Copyright 2025 The libx++ Authors. All rights reserved. + * Use of this source code is governed by a MIT license that can be + * found in the LICENSE file. + * + * server.h — xpp::http::test::Server: real-TCP HTTP test fixture. + * + * Layer 1 of the three-layer HTTP testing model (Rust-ecosystem aligned, + * see todos/httptest.md): + * + * Layer 0 router(req) direct call — handler/router/middleware + * (in-process, no sockets) unit tests (Router is a Handler) + * Layer 1 THIS FILE — test::Server — client tests & end-to-end: + * (real TCP, ephemeral port) the thing under test needs + * a real endpoint + * Layer 2 test::EvilServer (evil_server.h)— fault injection (a + * well-formed server can't lie + * about framing) + * + * Construction starts the server (bind 127.0.0.1:0 + listen synchronously + * — port() is valid immediately); RAII destructor stops and drains. Bind + * failure aborts: the environment is broken and the test cannot proceed + * (mirrors Go httptest.NewServer panicking on listen failure). + * + * Test-only. NOT part of any public umbrella — include explicitly: + * #include + * + * C++11-compatible. Header-only. + */ + +#ifndef XPP_HTTP_TEST_SERVER_H +#define XPP_HTTP_TEST_SERVER_H + +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace xpp { +namespace http { +namespace test { + +/* ── First-arg trait (constructor dispatch) — must precede Server ── */ + +namespace _ { +template struct first_arg; +template struct first_arg { + using type = A0; +}; +template struct first_arg { + using type = A0; +}; +template struct first_arg { + using type = A0; +}; +template struct first_arg : first_arg {}; +} // namespace _ + +/** + * @brief Preset HTTP response for the data-driven constructor. + * + * The spec compiles into Router routes at construction — no hand-rolled + * HTTP anywhere (the Router + Server own parsing and framing). + */ +struct TestResponseSpec { + StatusCode::Value status = StatusCode::Ok; + Vec> headers; + Bytes body; + /** Pre-response delay in ms (client-timeout tests). 0 = immediate. */ + uint64_t delay_ms = 0; + /** Echo the request body back (verifies POST payload transmission). */ + bool echo_request_body = false; + /** Add `X-Echo-Method: ` (verifies the client sent the verb). */ + bool echo_request_method = false; + /** + * Redirect requests whose path is not @p redirect_to to it with + * 302 + Location (client redirect-following tests). The final request + * is served with the rest of this spec. + */ + String redirect_to; + /** + * Stall this many ms after half the body is sent, then continue + * (read-timeout / slow-peer tests). Implemented as a channel body with + * a delayed producer: pausing production pauses transmission — the + * server only writes what has been produced. + */ + uint64_t mid_body_delay_ms = 0; +}; + +/** + * @brief Real-TCP HTTP test server (client tests & end-to-end). + * + * @code + * // Handler-driven (any behavior): + * test::Server ts([](Request req) { return Response::ok("hello"); }); + * auto r = client.get(ts.url().c_str()).await(); + * + * // Multi-route (Router capabilities: :params, layers, async handlers): + * test::Server routed([](ServerBuilder &b) { + * b.route("GET /users/:id", [](Request req, String id) { + * return Response::ok(id); + * }); + * }); + * + * // Data-driven (preset spec): + * TestResponseSpec spec; + * spec.body = Bytes::from("hello"); + * test::Server ts2(spec); + * @endcode + */ +class Server { +public: + Server() = default; + Server(Server &&) noexcept = default; + Server &operator=(Server &&) noexcept = default; + Server(const Server &) = delete; + Server &operator=(const Server &) = delete; + + ~Server() { + close(); + } + + /** + * @brief Builder-configurator constructor (multi-route). + * + * @p configure receives the ServerBuilder — register routes, layers, + * set timeouts as usual. The server binds 127.0.0.1:0 and listens + * synchronously before the constructor returns. + */ + /// Data-driven: the preset spec compiles into Router routes. + explicit Server(TestResponseSpec spec) { + init_dispatch([&spec](ServerBuilder &b) { build_spec_routes(b, spec); }, std::true_type()); + } + + /** + * @brief Generic constructor — dispatches on the callable's parameter. + * + * - `[](ServerBuilder &b) { ... }` → builder-configurator (multi-route, + * layers, timeouts — the general form) + * - `[](Request req) { return ...; }` → single handler (registered as + * the Router fallback, receives every request) + */ + template ::type, TestResponseSpec>::value>::type> + explicit Server(F &&f) { + // Dispatch: does F take (ServerBuilder&) or (Request)? + using Arg0 = typename _::first_arg::type>::type; + init_dispatch(std::forward(f), + std::is_same::type, ServerBuilder>()); + } + + /** @brief "http://127.0.0.1:" — ready to hand to a client. */ + String url() const { + std::string u = "http://127.0.0.1:" + std::to_string(port()); + return String::from_utf8(u.c_str()).unwrap(); + } + + /** @brief The kernel-assigned port. */ + uint16_t port() const noexcept { + return m_server.is_some() ? m_server.unwrap().port() : 0; + } + + /** @brief Stop the server and drain. Idempotent. */ + void close() { + if (m_server.is_some()) { + m_server.unwrap().stop(); + if (m_running.is_some()) { + m_running.unwrap().await(); + m_running = none; + } + m_server = none; + } + } + +private: + template void init_dispatch(F &&configure, std::true_type) { + // Builder-configurator: (ServerBuilder&) → void + ServerBuilder builder; + configure(builder); + m_server = some(builder.bind("127.0.0.1", 0).build().unwrap()); + m_running = some(m_server.unwrap().serve()); + } + + template void init_dispatch(H &&handler, std::false_type) { + // Single handler: (Request) → Response-ish — Router fallback + ServerBuilder builder; + Router r; + r.fallback(std::forward(handler)); + builder.router(std::move(r)); + m_server = some(builder.bind("127.0.0.1", 0).build().unwrap()); + m_running = some(m_server.unwrap().serve()); + } + + static void build_spec_routes(ServerBuilder &b, const TestResponseSpec &spec); + + Option m_server; + Option>> m_running; +}; + +/* ── Spec → Router compilation ─────────────────────────────────────── */ + +namespace _ { + +/// Request path from the URL (strip scheme/host if present, strip query). +inline String request_path(const Request &req) { + auto ub = req.url().as_bytes(); + size_t len = ub.size(); + // Strip query + for (size_t i = 0; i < len; ++i) { + if (ub.data()[i] == '?') { + len = i; + break; + } + } + return String::from_utf8(reinterpret_cast(ub.data()), len).unwrap(); +} + +/// Build the response for the spec (sync part — no request-body reads). +inline Response build_spec_response(const TestResponseSpec &spec, const Request &req) { + ResponseBuilder rb; + rb.status(spec.status); + for (auto &h : spec.headers) { + rb.header(h.first, h.second); + } + if (spec.echo_request_method) { + rb.header(String::from_utf8("X-Echo-Method").unwrap(), + String::from_utf8(to_string(req.method())).unwrap()); + } + return rb.body(spec.body); +} + +/// Build a 302 redirect to the target. +inline Response build_redirect(const String &to) { + return ResponseBuilder() + .status(StatusCode::Found) + .header(String::from_utf8("Location").unwrap(), to) + .body(); +} + +} // namespace _ + +/* ── build_spec_routes ─────────────────────────────────────────────── */ + +inline void Server::build_spec_routes(ServerBuilder &b, const TestResponseSpec &spec) { + // The spec dispatch runs as the Router fallback (every request hits it). + Router r; + r.fallback([spec](Request req) -> Promise> { + // 1. Redirect: any path != redirect_to → 302 + Location. + if (!spec.redirect_to.empty()) { + String path = _::request_path(req); + if (path != spec.redirect_to) { + return xpp::resolve(Result(xpp::ok, _::build_redirect(spec.redirect_to))); + } + } + + // 2. Echo request body: drain the body, return it. + if (spec.echo_request_body) { + // Capture the method BEFORE the request is consumed by into_body(). + String method = String::from_utf8(to_string(req.method())).unwrap(); + return req.into_body().bytes().then( + [spec, method](Result b) -> Promise> { + ResponseBuilder rb; + rb.status(spec.status); + for (auto &h : spec.headers) { + rb.header(h.first, h.second); + } + if (spec.echo_request_method) { + rb.header(String::from_utf8("X-Echo-Method").unwrap(), method); + } + return xpp::resolve(Result(xpp::ok, rb.body(std::move(b).unwrap()))); + }); + } + + // 3. Mid-body stall: channel body with a delayed producer. + if (spec.mid_body_delay_ms > 0 && spec.body.size() > 0) { + auto pair = sync::mpsc::channel(4); + auto tx = std::move(pair.first); + auto rx = std::move(pair.second); + + Bytes body = spec.body; + // First half now, second half after the stall. + size_t half = body.size() / 2; + Bytes first = Bytes::copy(reinterpret_cast(body.data()), half); + Bytes second = + Bytes::copy(reinterpret_cast(body.data() + half), body.size() - half); + + // Sender is move-only with non-const send/close — heap-hold it for + // the .then chain (C++11 has no move-capture). + auto tx_holder = Arc>::make(std::move(tx)); + uint64_t ms = spec.mid_body_delay_ms; + xpp::spawn([tx_holder, first, second, ms]() -> Promise { + return tx_holder->send(first).then([tx_holder, second, ms]() -> Promise { + return xpp::after(ms).then([tx_holder, second]() -> Promise { + return tx_holder->send(second).then([tx_holder]() -> Promise { + tx_holder->close(); + return xpp::resolve(); + }); + }); + }); + }); + + ResponseBuilder rb; + rb.status(spec.status); + for (auto &h : spec.headers) { + rb.header(h.first, h.second); + } + if (spec.echo_request_method) { + rb.header(String::from_utf8("X-Echo-Method").unwrap(), + String::from_utf8(to_string(req.method())).unwrap()); + } + return xpp::resolve(Result(xpp::ok, rb.body(Body::from_channel(std::move(rx))))); + } + + // 4. Delay then respond. + if (spec.delay_ms > 0) { + TestResponseSpec s = spec; + String method = String::from_utf8(to_string(req.method())).unwrap(); + return xpp::after(s.delay_ms).then([s, method]() -> Promise> { + ResponseBuilder rb; + rb.status(s.status); + for (auto &h : s.headers) { + rb.header(h.first, h.second); + } + if (s.echo_request_method) { + rb.header(String::from_utf8("X-Echo-Method").unwrap(), method); + } + return xpp::resolve(Result(xpp::ok, rb.body(s.body))); + }); + } + + // 5. Immediate response. + return xpp::resolve(Result(xpp::ok, _::build_spec_response(spec, req))); + }); + b.router(std::move(r)); +} + +} // namespace test +} // namespace http +} // namespace xpp + +#endif // XPP_HTTP_TEST_SERVER_H diff --git a/libxpp/xpp/http/test/server_test.cpp b/libxpp/xpp/http/test/server_test.cpp new file mode 100644 index 0000000..cbc22cd --- /dev/null +++ b/libxpp/xpp/http/test/server_test.cpp @@ -0,0 +1,129 @@ +/* + * Copyright 2025 The libx++ Authors. All rights reserved. + * Use of this source code is governed by a MIT license that can be + * found in the LICENSE file. + * + * server_test.cpp — Tests for xpp::http::test::Server (Layer 1 fixture). + */ + +#include + +#include +#include +#include +#include + +using namespace xpp; +using namespace xpp::http; + +static std::string url_of(const test::Server &ts, const char *path = "") { + return "http://127.0.0.1:" + std::to_string(ts.port()) + path; +} + +TEST(TestServerTest, HandlerConstructor) { + EventLoop loop; + WaitScope scope(loop); + + test::Server ts([](Request) { return Response::ok("hello"); }); + EXPECT_GT(ts.port(), 0); + + auto client = Client::builder().build().unwrap(); + auto r = client.get(url_of(ts).c_str()).await(); + ASSERT_TRUE(r.is_ok()); + auto body = r.unwrap().bytes().await().unwrap().to_string().unwrap(); + EXPECT_EQ(body, "hello"); +} + +TEST(TestServerTest, SpecConstructor) { + EventLoop loop; + WaitScope scope(loop); + + test::TestResponseSpec spec; + spec.body = Bytes::from("spec-body"); + test::Server ts(spec); + + auto client = Client::builder().build().unwrap(); + auto r = client.get(url_of(ts).c_str()).await(); + ASSERT_TRUE(r.is_ok()); + auto body = r.unwrap().bytes().await().unwrap().to_string().unwrap(); + EXPECT_EQ(body, "spec-body"); +} + +TEST(TestServerTest, SpecRedirect) { + EventLoop loop; + WaitScope scope(loop); + + test::TestResponseSpec spec; + spec.redirect_to = String::from_utf8("/b").unwrap(); + spec.body = Bytes::from("redirected!"); + test::Server ts(spec); + + auto client = Client::builder().redirect(3).build().unwrap(); + auto r = client.get(url_of(ts, "/a").c_str()).await(); + ASSERT_TRUE(r.is_ok()); + auto body = r.unwrap().bytes().await().unwrap().to_string().unwrap(); + EXPECT_EQ(body, "redirected!"); +} + +TEST(TestServerTest, SpecEchoBody) { + EventLoop loop; + WaitScope scope(loop); + + test::TestResponseSpec spec; + spec.echo_request_body = true; + test::Server ts(spec); + + auto client = Client::builder().build().unwrap(); + auto r = client.post(url_of(ts).c_str(), "my-payload").await(); + ASSERT_TRUE(r.is_ok()); + auto body = r.unwrap().bytes().await().unwrap().to_string().unwrap(); + EXPECT_EQ(body, "my-payload"); +} + +TEST(TestServerTest, SpecEchoMethod) { + EventLoop loop; + WaitScope scope(loop); + + test::TestResponseSpec spec; + spec.echo_request_method = true; + spec.echo_request_body = true; + test::Server ts(spec); + + auto client = Client::builder().build().unwrap(); + auto r = client.post(url_of(ts).c_str(), "data").await(); + ASSERT_TRUE(r.is_ok()); + auto m = r.unwrap().headers().get(String::from_utf8("x-echo-method").unwrap()); + ASSERT_TRUE(m.is_some()); + EXPECT_EQ(m.unwrap(), String::from_utf8("POST").unwrap()); +} + +TEST(TestServerTest, BuilderConfigureConstructor) { + EventLoop loop; + WaitScope scope(loop); + + test::Server ts([](ServerBuilder &b) { + b.route("GET /users/:id", [](Request, String id) { return Response::ok(id); }); + }); + + auto client = Client::builder().build().unwrap(); + auto r = client.get(url_of(ts, "/users/42").c_str()).await(); + ASSERT_TRUE(r.is_ok()); + auto body = r.unwrap().bytes().await().unwrap().to_string().unwrap(); + EXPECT_EQ(body, "42"); +} + +TEST(TestServerTest, SpecDelay) { + EventLoop loop; + WaitScope scope(loop); + + test::TestResponseSpec spec; + spec.body = Bytes::from("slow"); + spec.delay_ms = 50; // short enough for CI, enough to verify the timer path + test::Server ts(spec); + + auto client = Client::builder().timeout(5000).build().unwrap(); + auto r = client.get(url_of(ts).c_str()).await(); + ASSERT_TRUE(r.is_ok()); + auto body = r.unwrap().bytes().await().unwrap().to_string().unwrap(); + EXPECT_EQ(body, "slow"); +} diff --git a/libxpp/xpp/http/test_server.h b/libxpp/xpp/http/test_server.h deleted file mode 100644 index 36f18ae..0000000 --- a/libxpp/xpp/http/test_server.h +++ /dev/null @@ -1,539 +0,0 @@ -/* - * Copyright 2025 The libx++ Authors. All rights reserved. - * Use of this source code is governed by a MIT license that can be - * found in the LICENSE file. - * - * test_server.h — Minimal HTTP/1.1 test server for xpp::http client tests. - * - * A static-response test fixture: bind a loopback xTcpListener on port 0 - * (kernel-assigned), accept connections, read the request (headers + body - * framed by Content-Length), optionally delay `delay_ms` (to test client - * timeouts), then write back a preset HTTP/1.1 response and close. - * - * Implementation: pure libx C API (xTcpListener + xTcpConn + xSocket + - * xTimer). No xpp::fiber, no .then() chains. The accept callback switches - * the conn's xSocket (level-triggered) to a read callback; it accumulates - * the request until the full header block + body have arrived, then - * responds (or schedules a timer for delayed responses). This avoids - * fiber/xEventLoopRun interaction issues on Linux shared builds. - * - * NOT the future `xpp::http::Server` module — test-only, no routing, - * no concurrency, no streaming. Kept in `xpp::http::test` subnamespace - * and NOT included from the public `xpp/http.h` umbrella. - */ - -#ifndef XPP_HTTP_TEST_SERVER_H -#define XPP_HTTP_TEST_SERVER_H - -#include // ntohs, sockaddr_in -#include -#include // getsockname, sockaddr_storage - -#include -#include -#include -#include -#include - -#include -#include -#include -#include - -#include // XDEF_HANDLE -#include // xTimerStart -#include // xSocket, xSocketSetCallback, xSocketSetMask -#include // xTcpListener, xTcpConn, xTcpConnSend, xTcpConnClose - -namespace xpp { -namespace http { -namespace test { - -/** - * @brief Preset HTTP response for `TestServer` to return. - */ -struct TestResponseSpec { - StatusCode::Value status = StatusCode::Ok; - Vec> headers; - Bytes body; - /** Pre-response delay in milliseconds. 0 = respond immediately. */ - uint64_t delay_ms = 0; - /** - * @brief If true, the response body echoes the request body received. - * - * Used to verify request bodies are transmitted (e.g. POST payloads). - * The preset `body` is ignored when this is set. - */ - bool echo_request_body = false; - /** - * @brief If true, the response includes `X-Echo-Method: ` - * so tests can verify the client sent the intended HTTP verb. - */ - bool echo_request_method = false; - /** - * @brief If non-empty, redirect requests whose path is not @p redirect_to - * to it with a 302 Found + Location header. - * - * Used to test client redirect following. The final (target) request - * is served with the rest of this spec (status/headers/body). - */ - String redirect_to; - /** - * @brief If > 0, send only the first N bytes of `body` but declare the - * full Content-Length, then close the connection. - * - * Simulates a mid-body disconnect for error-path tests. - */ - size_t truncate_body_after = 0; - /** - * @brief If > 0, stall for this many ms after sending half the response - * body, then continue. Simulates a stalled peer for read-timeout - * (low-speed) tests. - */ - uint64_t mid_body_delay_ms = 0; -}; - -/** - * @brief Static-response HTTP/1.1 test server. - * - * Bind to `127.0.0.1:0` (kernel-assigned port), accept connections on - * the current `EventLoop`, read each request (headers + Content-Length - * body), and respond with the preset `TestResponseSpec`. Closes the - * connection after each response. - * - * Usage: - * @code - * xpp::EventLoop loop; - * xpp::WaitScope scope(loop); - * - * xpp::http::test::TestResponseSpec spec; - * spec.status = xpp::http::StatusCode::Ok; - * spec.body = xpp::Bytes::from("hello"); - * - * auto server = xpp::http::test::TestServer::start(spec); - * // server.port() → kernel-assigned port - * @endcode - */ -class TestServer { -public: - TestServer() = default; - TestServer(TestServer &&) noexcept = default; - TestServer &operator=(TestServer &&) noexcept = default; - TestServer(const TestServer &) = delete; - TestServer &operator=(const TestServer &) = delete; - - ~TestServer() { - stop(); - } - - /** - * @brief Start a `TestServer` bound to `127.0.0.1:0`. - * - * Must be called inside a `WaitScope` (the current thread must have - * entered an `EventLoop`). Returns a move-only `TestServer` value. - */ - static TestServer start(TestResponseSpec spec) { - TestServer ts; - ts.spec_ = std::move(spec); - - // Create the listener. xTcpListenerCreate uses xEventLoopCurrent(), - // so the caller must have entered an EventLoop. - ts.listener_ = xTcpListenerCreate("127.0.0.1", 0, nullptr, on_accept, &ts); - XPP_ASSERT(ts.listener_ != nullptr, "TestServer: xTcpListenerCreate failed"); - - // Get the kernel-assigned port. - xSocket sock = xTcpListenerSocket(ts.listener_); - XPP_ASSERT(sock != nullptr, "TestServer: xTcpListenerSocket failed"); - - struct sockaddr_storage addr; - socklen_t addrlen = sizeof(addr); - if (getsockname(xSocketFd(sock), (struct sockaddr *)&addr, &addrlen) == 0) { - ts.port_ = ntohs(((struct sockaddr_in *)&addr)->sin_port); - } - XPP_ASSERT(ts.port_ > 0, "TestServer: getsockname failed"); - - return ts; - } - - /** @brief The kernel-assigned port number. */ - uint16_t port() const noexcept { - return port_; - } - - /** - * @brief Stop the server. - * - * Closes the listening socket. In-flight connections are not affected. - * Idempotent. - */ - void stop() { - if (listener_) { - xTcpListenerDestroy(listener_); - listener_ = nullptr; - } - } - -private: - TestResponseSpec spec_; - xTcpListener listener_ = nullptr; - uint16_t port_ = 0; - - /* ── Per-connection state (heap-allocated; freed once responded) ── */ - - struct ConnState { - TestServer *self; - xTcpConn conn; - xSocket sock; ///< The conn's xSocket (level-triggered source). - xTimer timer = nullptr; ///< Pending delay timer, if any. - Vec buf; ///< Request bytes accumulated so far. - size_t header_end = 0; ///< Index past the "\r\n\r\n" of the header block. - size_t content_length = 0; ///< Parsed Content-Length (0 = no body). - bool headers_done = false; - bool closed = false; ///< Guards against double cleanup. - String request_path; ///< Path from the request line (e.g. "/b"). - String request_method; ///< Method token (e.g. "GET"). - String resp; ///< Response being sent (send_off marks progress). - size_t send_off = 0; - uint64_t mid_body_delay_ms = 0; ///< From spec — stall after half the body. - bool mid_body_paused = false; ///< The stall already happened. - }; - - /* ── Accept callback (runs on EventLoop thread) ──────────────────── */ - - static void on_accept(xTcpListener /*listener*/, xTcpConn conn, const struct sockaddr * /*addr*/, - socklen_t /*addrlen*/, void *arg) { - auto *self = static_cast(arg); - - auto *st = new ConnState; - st->self = self; - st->conn = conn; - - // Drive reads/writes through the xSocket's own event source (created - // level-triggered by tcp_listener.c). Never register a second source on - // the conn fd — re-adding the same (fd, filter) to kqueue keeps the - // first registration's EV_CLEAR and breaks level-triggered re-firing. - st->sock = xTcpConnSocket(conn); - st->mid_body_delay_ms = self->spec_.mid_body_delay_ms; - XPP_ASSERT(st->sock != nullptr, "TestServer: xTcpConnSocket failed"); - xSocketSetCallback(st->sock, on_xsocket_event, st); - xSocketSetMask(st->sock, xEvent_Read); - } - - /* ── Conn event callback (runs on EventLoop thread) ─────────────── */ - - static void on_xsocket_event(xSocket sock, xEventMask mask, void *arg) { - auto *st = static_cast(arg); - (void)sock; - if (mask & xEvent_Read) on_conn_readable(st); - if (mask & xEvent_Write) pump_send(st); - } - - static void on_conn_readable(ConnState *st) { - char tmp[4096]; - ssize_t n = xTcpConnRecv(st->conn, tmp, sizeof(tmp)); - if (n <= 0) { - close_conn(st); // EOF or error — client gone before we responded - return; - } - for (ssize_t i = 0; i < n; ++i) - st->buf.push(static_cast(tmp[i])); - - if (!st->headers_done) { - size_t hs = find_header_end(st->buf); - if (hs != SIZE_MAX) { - st->header_end = hs; - st->content_length = parse_content_length(st->buf, hs); - st->request_path = parse_request_path(st->buf, hs); - st->request_method = parse_request_method(st->buf); - st->headers_done = true; - } - } - - // Respond once the header block and the full body have arrived. - if (st->headers_done && st->buf.len() >= st->header_end + st->content_length) { - if (st->self->spec_.delay_ms > 0) { - // Delayed response — schedule a one-shot timer. The timer owns the - // ConnState (arg) until it fires or the loop is destroyed. - st->timer = xTimerStart(on_delay_timer, st, on_delay_cancel, st->self->spec_.delay_ms, 0); - XPP_ASSERT(st->timer != nullptr, "TestServer: xTimerStart failed"); - } else { - respond_and_close(st); - } - } - } - - /* ── Response path ──────────────────────────────────────────────── */ - - static void on_delay_timer(void *arg) { - auto *st = static_cast(arg); - st->timer = nullptr; // handle consumed by the fire — don't xTimerStop it - respond_and_close(st); - } - - static void on_delay_cancel(void *arg) { - auto *st = static_cast(arg); - st->timer = nullptr; // handle consumed by the cancel - close_conn(st); - } - - static void respond_and_close(ConnState *st) { - // Build the response, then switch the source to write mode and pump - // the send. The socket is non-blocking — a single xTcpConnSend can - // write only part of a large response, so we drive the rest from - // level-triggered writable events. - st->resp = build_response(*st); - st->send_off = 0; - xSocketSetMask(st->sock, xEvent_Write); - pump_send(st); - } - - static void on_mid_body_timer(void *arg) { - auto *st = static_cast(arg); - st->timer = nullptr; // handle consumed by the fire - xSocketSetMask(st->sock, xEvent_Write); // re-arm writable, resume sending - pump_send(st); - } - - static void on_mid_body_cancel(void *arg) { - auto *st = static_cast(arg); - st->timer = nullptr; // handle consumed by the cancel - close_conn(st); - } - - static void pump_send(ConnState *st) { - // Mid-body stall (read-timeout testing): after half the response is - // out, pause once for mid_body_delay_ms before continuing. - if (st->mid_body_delay_ms > 0 && !st->mid_body_paused && st->send_off >= st->resp.len() / 2) { - st->mid_body_paused = true; - // Unregister the writable event so the stall actually blocks sending - // (a writable event would otherwise fire as soon as the client reads). - xSocketSetMask(st->sock, 0); - st->timer = xTimerStart(on_mid_body_timer, st, on_mid_body_cancel, st->mid_body_delay_ms, 0); - XPP_ASSERT(st->timer != nullptr, "TestServer: xTimerStart failed"); - return; - } - auto bytes = st->resp.as_bytes(); - // Limit per-send size: without this, a loopback socket buffer can - // swallow the whole response before the mid-body stall, making the - // pause invisible to the client. - size_t remaining = bytes.size() - st->send_off; - if (remaining > 64 * 1024) remaining = 64 * 1024; - ssize_t n = xTcpConnSend(st->conn, reinterpret_cast(bytes.data()) + st->send_off, - remaining); - if (n < 0) { - // Non-blocking socket: EAGAIN means "retry when writable" — the - // level-triggered source fires again. Anything else is fatal. - if (errno == EAGAIN || errno == EWOULDBLOCK) return; - close_conn(st); - return; - } - if (n == 0) { // defensive: no progress — give up rather than spin - close_conn(st); - return; - } - st->send_off += static_cast(n); - if (st->send_off >= bytes.size()) close_conn(st); - // else: partial write — the level-triggered writable event re-fires. - } - - static void close_conn(ConnState *st) { - if (st->closed) return; - st->closed = true; - // Cancel a still-pending delay timer (EOF arrived before it fired). - // Fired/cancelled timers already set st->timer = nullptr. - if (st->timer) { - xTimerStop(st->timer); - st->timer = nullptr; - } - // xTcpConnClose destroys the xSocket (and its event source). - xTcpConnClose(st->conn); - delete st; - } - - /* ── Request parsing ────────────────────────────────────────────── */ - - /** @brief Index past the "\r\n\r\n" ending the header block, or SIZE_MAX. */ - static size_t find_header_end(const Vec &buf) { - size_t n = buf.len(); - if (n < 4) return SIZE_MAX; - for (size_t i = 0; i + 4 <= n; ++i) { - if (buf[i] == '\r' && buf[i + 1] == '\n' && buf[i + 2] == '\r' && buf[i + 3] == '\n') { - return i + 4; - } - } - return SIZE_MAX; - } - - /** @brief Request path from the request line, e.g. "/b" from "GET /b HTTP/1.1". */ - static String parse_request_path(const Vec &buf, size_t header_end) { - const uint8_t *p = buf.data(); - size_t i = 0; - while (i < header_end && p[i] != ' ' && p[i] != '\t') - ++i; // skip method token - while (i < header_end && (p[i] == ' ' || p[i] == '\t')) - ++i; // skip whitespace - size_t start = i; - while (i < header_end && p[i] != ' ' && p[i] != '\t' && p[i] != '\r') - ++i; - if (i == start) return String(); - return String::from_utf8(reinterpret_cast(p + start), i - start).unwrap(); - } - - /** @brief Method token from the request line, e.g. "GET" from "GET /b HTTP/1.1". */ - static String parse_request_method(const Vec &buf) { - const uint8_t *p = buf.data(); - size_t i = 0; - while (i < buf.len() && p[i] != ' ' && p[i] != '\t' && p[i] != '\r') - ++i; - return String::from_utf8(reinterpret_cast(p), i).unwrap(); - } - - /** @brief Parse Content-Length from the header block. 0 if absent. */ - static size_t parse_content_length(const Vec &buf, size_t header_end) { - static const char kName[] = "content-length"; - const uint8_t *p = buf.data(); - size_t i = 0; - while (i < header_end) { - size_t le = i; - while (le < header_end && p[le] != '\n') - ++le; - size_t len = le; - if (len > i && p[len - 1] == '\r') --len; - - size_t j = 0; - for (; j < sizeof(kName) - 1 && i + j < len; ++j) { - uint8_t c = p[i + j]; - if (c >= 'A' && c <= 'Z') c = static_cast(c - 'A' + 'a'); - if (c != static_cast(kName[j])) break; - } - if (j == sizeof(kName) - 1 && i + j < len && p[i + j] == ':') { - size_t d = i + j + 1; - while (d < len && (p[d] == ' ' || p[d] == '\t')) - ++d; - size_t value = 0; - while (d < len && p[d] >= '0' && p[d] <= '9') { - value = value * 10 + static_cast(p[d] - '0'); - ++d; - } - return value; - } - if (le >= header_end) break; - i = le + 1; - } - return 0; - } - - /* ── Response building ──────────────────────────────────────────── */ - - static String build_response(const ConnState &st) { - const TestResponseSpec &spec = st.self->spec_; - String resp; - - // Redirect: 302 + Location for any request not already at the target. - // Lets the client exercise CURLOPT_FOLLOWLOCATION end to end. - if (!spec.redirect_to.empty() && st.request_path != spec.redirect_to) { - resp.push_str(String::from_utf8("HTTP/1.1 302 Found\r\nLocation: ").unwrap()); - resp.push_str(spec.redirect_to); - resp.push_str( - String::from_utf8("\r\nContent-Length: 0\r\nConnection: close\r\n\r\n").unwrap()); - return resp; - } - - // Status line - resp.push_str(String::from_utf8("HTTP/1.1 ").unwrap()); - resp.push_str( - String::from_utf8(std::to_string(static_cast(spec.status)).c_str()).unwrap()); - resp.push_str(String::from_utf8(" ").unwrap()); - resp.push_str(to_reason_phrase(spec.status)); - resp.push_str(String::from_utf8("\r\n").unwrap()); - - // User-provided headers - for (const auto &h : spec.headers) { - resp.push_str(h.first); - resp.push_str(String::from_utf8(": ").unwrap()); - resp.push_str(h.second); - resp.push_str(String::from_utf8("\r\n").unwrap()); - } - - if (spec.echo_request_method) { - resp.push_str(String::from_utf8("X-Echo-Method: ").unwrap()); - resp.push_str(st.request_method); - resp.push_str(String::from_utf8("\r\n").unwrap()); - } - - // Body: echoed request body or the preset spec body. - Bytes body = spec.body; - if (spec.echo_request_body) { - size_t blen = st.buf.len() - st.header_end; - body = Bytes::copy(reinterpret_cast(st.buf.data() + st.header_end), blen); - } - // Truncation: declare the full length, send only the first N bytes. - size_t declared_len = body.size(); - if (spec.truncate_body_after > 0 && body.size() > spec.truncate_body_after) { - body = Bytes::copy(reinterpret_cast(body.data()), spec.truncate_body_after); - } - - // Content-Length + Connection: close - resp.push_str(String::from_utf8("Content-Length: ").unwrap()); - resp.push_str(String::from_utf8(std::to_string(declared_len).c_str()).unwrap()); - resp.push_str(String::from_utf8("\r\n").unwrap()); - resp.push_str(String::from_utf8("Connection: close\r\n\r\n").unwrap()); - - // Body (appended to the same buffer for a single send) - if (body.size() > 0) { - auto body_bytes = body.as_span(); - // String is UTF-8 — body may be binary, but for test purposes - // appending raw bytes works because String stores Vec. - resp.push_str( - String::from_utf8(reinterpret_cast(body_bytes.data()), body_bytes.size()) - .unwrap_or(String())); - } - return resp; - } - - /** @brief Minimal reason-phrase table for common status codes. */ - static String to_reason_phrase(StatusCode::Value code) { - switch (code) { - case StatusCode::Ok: - return String::from_utf8("OK").unwrap(); - case StatusCode::Created: - return String::from_utf8("Created").unwrap(); - case StatusCode::NoContent: - return String::from_utf8("No Content").unwrap(); - case StatusCode::MovedPermanently: - return String::from_utf8("Moved Permanently").unwrap(); - case StatusCode::Found: - return String::from_utf8("Found").unwrap(); - case StatusCode::NotModified: - return String::from_utf8("Not Modified").unwrap(); - case StatusCode::BadRequest: - return String::from_utf8("Bad Request").unwrap(); - case StatusCode::Unauthorized: - return String::from_utf8("Unauthorized").unwrap(); - case StatusCode::Forbidden: - return String::from_utf8("Forbidden").unwrap(); - case StatusCode::NotFound: - return String::from_utf8("Not Found").unwrap(); - case StatusCode::MethodNotAllowed: - return String::from_utf8("Method Not Allowed").unwrap(); - case StatusCode::InternalServerError: - return String::from_utf8("Internal Server Error").unwrap(); - case StatusCode::NotImplemented: - return String::from_utf8("Not Implemented").unwrap(); - case StatusCode::BadGateway: - return String::from_utf8("Bad Gateway").unwrap(); - case StatusCode::ServiceUnavailable: - return String::from_utf8("Service Unavailable").unwrap(); - case StatusCode::GatewayTimeout: - return String::from_utf8("Gateway Timeout").unwrap(); - default: - return String::from_utf8("OK").unwrap(); - } - } -}; - -} // namespace test -} // namespace http -} // namespace xpp - -#endif // XPP_HTTP_TEST_SERVER_H diff --git a/todos/httptest.md b/todos/httptest.md deleted file mode 100644 index 84b0a4c..0000000 --- a/todos/httptest.md +++ /dev/null @@ -1,156 +0,0 @@ -# TODO: httptest 化 —— 基于 xpp::http::Server 重写 test_server.h - -> 来源:test_server.h(540 行手写 HTTP fixture)与 Server 的重复;目标是对齐 -> Go `net/http/httptest` 的人体工学,xpp 化 -> 前置:#84(Response::ok)已合并 ✓;Router as Handler 已合并 ✓ -> (`Server::builder().router(r)` 已存在,单 handler 构造即 -> "整个 Router 是一个 handler";原计划的 libx mux catch-all 前置作废) - -## 目标 - -像 Go httptest 那样:**一行起一个真 HTTP 服务,handler 任意定制**,测试完自动收。 - -```go -// Go 的样子 -ts := httptest.NewServer(http.HandlerFunc(func(w, r) { fmt.Fprintln(w, "hello") })) -defer ts.Close() -client.Get(ts.URL) -``` - -```cpp -// xpp 的样子 —— 构造即启动,RAII 收尾 -xpp::http::test::Server ts([](xpp::http::Request req) { - return xpp::http::Response::ok("hello"); -}); -// ts.url() == "http://127.0.0.1:54321",析构自动 stop+drain -auto r = client.get(ts.url().c_str()).await(); - -// 多路由测试:主构造器收 builder 配置器(完整路由能力,含 :param 注入) -xpp::http::test::Server routed([](xpp::http::ServerBuilder &b) { - b.route("GET /users/:id", [](Request req, String id) { return Response::ok(id); }) - .route("POST /echo", [](Request req) -> Promise> { - auto body = co_await req.into_body().bytes(); - return Response::ok(body.unwrap()); - }); -}); -``` - -## Go httptest 特性映射 - -| Go httptest | xpp 化 | 说明 | -| --- | --- | --- | -| `NewServer(h)` | `test::Server ts(handler)` —— **构造即启动** | 对齐 Go:`NewServer` 监听失败直接 panic;fixture 对 bind 失败 `XPP_ASSERT`(环境已坏,测试无法继续),不传播 `Result` | -| `NewServer(mux)`(Go 的 mux 本身是 handler) | `test::Server ts([](ServerBuilder &b){ b.route(...) ... })` | **主构造器**:收 builder 配置器,完整路由(`:param` 注入、多路由);对应 Go 里"传整个 mux"的用法 | -| `ts.URL` | `ts.url()` → `"http://127.0.0.1:"` | 现在测试手拼 `url_for(port, path)` | -| `ts.Close()` / defer | RAII dtor(stop + drain running) | 现有 TestServer 已有 | -| `NewTLSServer` + `ts.Client()` | **非目标**(暂缓) | libx xHttpServer 尚无 TLS serving;Client 是 libcurl 无信任配置需求 | -| `httptest.NewRequest` | 免费已有:`Request::builder()` | 纯值构造,无需网络 | -| `ResponseRecorder`(无网络测 handler) | 免费已有:直接调用 handler | 我们的 handler 返回 `Response` 值(无 io.Writer 边信道),`auto resp = h(req)` 即断言 | -| —(Go 没有) | `test::Server ts(spec)` | 保留数据驱动构造重载,echo/redirect/delay 一行配置 | - -## API 设计 - -### 1. `xpp::http::test::Server`(替代现 TestServer,文件仍为 `test_server.h`) - -```cpp -namespace xpp { namespace http { namespace test { - -class Server { -public: - // ── 主构造器:builder 配置器(完整路由能力)── - // 构造即启动:套上 bind("127.0.0.1", 0) + build + serve,失败 XPP_ASSERT。 - // Go 的 NewServer(mux) 对应物 —— Go 的 mux 是 handler,我们的路由在 - // ServerBuilder 上,所以通用形态是"给你 builder,随便配"。 - template explicit Server(F &&configure); // F: (ServerBuilder&) -> void - - // ── 便捷构造器:单 handler 接管一切(catch-all/fallback 路由)── - // 适合静态/回显类简单 fixture;等价于 configure 里注册一个 fallback。 - template explicit Server(H &&handler); - - // ── 数据驱动重载:现有 TestResponseSpec 语义照搬 ── - explicit Server(TestResponseSpec spec); - - Server(Server &&) noexcept; // move-only - ~Server(); // RAII:stop() + running.await() - - String url() const; // "http://127.0.0.1:" - uint16_t port() const; - void close(); // 幂等;dtor 亦调用 -}; - -}}} // namespace -``` - -### 2. 前置:Server 的 catch-all 路由(libx mux 层) - -现状:`xHttpRouteMatch_` 只支持段精确 / `:param`,无通配——`Server(handler)` 无法一行接管。 -**方案**:libx mux 支持保留 pattern `"*"`(任意 method + 任意 path,注册即 fallback; -resolve 循环里普通路由不中再试 `*` 路由)。xpp 侧暴露为 -`ServerBuilder::fallback(h)`——顺带成为通用能力(自定义 404/默认 handler), -与 Go `http.NotFoundHandler` / mux `/` 前缀对齐。 - -### 3. `TestResponseSpec` 的实现迁移(spec → handler 编译) - -| spec 字段 | Server 上的实现 | 迁移后 | -| --- | --- | --- | -| status/headers/body | sync handler 返回 preset Response | ✅ | -| `echo_request_body` | 协程 handler:`co_await req.into_body().bytes()` 回显 | ✅ | -| `echo_request_method` | handler 读 `req.method()` 写 `X-Echo-Method` | ✅ | -| `redirect_to` | handler:path ≠ target → 302 + Location | ✅ | -| `delay_ms` | 协程 handler:`co_await xpp::after(ms)` 再响应 | ✅ | -| `mid_body_delay_ms` | channel body + 延迟生产者:发前半 → `co_await after()` → 发后半。暂停生产=暂停发送,现 64KB 发送上限的 workaround 不再需要(server 只写已生产的数据) | ✅ | -| `truncate_body_after` | **不可表达**——需要"说谎的对端"(声明完整 Content-Length、只发一半、硬断连),well-formed Server 永远不会 | ❌ → EvilServer | - -### 4. `xpp::http::test::EvilServer`(单独封装,新文件 `evil_server.h`) - -fault-injection 对端,只做 well-formed 服务器做不到的事。首期仅保留截断注入: - -```cpp -struct EvilSpec { - Bytes body; // 完整响应体(用于声明 Content-Length) - size_t send_only = 0; // 实际只发前 N 字节然后硬断连 -}; - -class EvilServer { // 原始 TCP(从现 test_server.h 的 socket -public: // 路径精简而来,~80 行),构造即启动 - explicit EvilServer(EvilSpec spec); - uint16_t port() const; - void close(); // RAII 同上 -}; -``` - -现 `client_test.cpp` 的 `TruncatedBodySurfacesAsError`(1MB 声明 / 64KB 实发) -迁移至此。后续如需更多恶意行为(乱序、坏 chunked、慢攻击)按需加字段。 - -### 5. 删除清单(重写后) - -- 手写请求解析:`find_header_end` / `parse_content_length` / `parse_request_path` / - `parse_request_method`(Server 的 llhttp 接管) -- 手写响应构建:`build_response` / `to_reason_phrase` 表(Response/StatusCode 接管) -- 手写 framing 与 `Connection: close` 管理、`pump_send` 的 64KB 上限 workaround -- 估算:540 行 → ~150 行(TestServer 适配层)+ ~80 行(EvilServer) - -## 迁移计划 - -1. libx mux `"*"` catch-all + `xHttpRouteMatch_`/`xHttpMuxResolve` 支持 + 单测 -2. xpp `ServerBuilder::fallback(h)` + server_test 补 fallback 用例 -3. `test::Server`(handler + spec 双 start)+ `url()` -4. `test::EvilServer`(从现 socket 路径抽取精简) -5. `client_test.cpp` / `http_convenience_test.cpp` 全量迁移 - (truncate 测试 → EvilServer,其余 → 新 Server;`url_for` 换 `ts.url()`) -6. 删除旧实现 - -## 非目标 - -- TLS test server(等 libx TLS serving;届时对齐 `NewTLSServer`/`ts.Client()`) -- `httptest.NewRequest` / `ResponseRecorder` 等价物(`Request::builder()` 与 - 值语义 handler 调用已免费覆盖) -- proxy/hijack 注入(Go `httptest.NewServer` 的 `ts.Config` 钩子)——按需再加 - -## 验证 - -- client_test / http_convenience_test / server_test 全量 -- CI 四配置 ASan + TSan lane(重点:Linux shared build —— 旧实现头注释提到的 - "fiber/xEventLoopRun 交互问题"指向 fiber,Server 是 waker-driven 无 fiber, - 理论上顾虑消失,需 CI 实证) -- timing 类测试(delay 200ms / mid-body 2000ms)在 CI 慢机器上复核阈值余量 From 1902c5a24ebdd65e7e12a700d6307106ea976ce8 Mon Sep 17 00:00:00 2001 From: mivinci Date: Wed, 26 Aug 2026 11:56:44 +0800 Subject: [PATCH 2/2] fix(xpp/http): guard streaming response against connection-close UAF MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The response streaming path (stream_channel_body) held a raw xHttpCtx* across suspension points, protected only by ServerLifetime, which covers server teardown but not connection teardown. On client disconnect the event loop frees the stream (and the xHttpCtx embedded in it) via xHttpConnClose — a path that never reaches on_done — so the parked streaming task wrote the freed ctx, a heap-use-after-free. Add an on_close callback to xHttpRouteInfo, invoked by xHttpConnClose before the stream is freed, plus a ConnLifetime (Arc) keyed by ctx address that the streaming task checks before touching ctx. --- issues/http-streaming-conn-close-uaf.md | 70 ++++++++++++++++++++++++ libx/x/http/server.c | 7 +++ libx/x/http/server.h | 28 ++++++++-- libxpp/xpp/http/server.h | 71 ++++++++++++++++++++----- 4 files changed, 159 insertions(+), 17 deletions(-) create mode 100644 issues/http-streaming-conn-close-uaf.md diff --git a/issues/http-streaming-conn-close-uaf.md b/issues/http-streaming-conn-close-uaf.md new file mode 100644 index 0000000..9fa8b9a --- /dev/null +++ b/issues/http-streaming-conn-close-uaf.md @@ -0,0 +1,70 @@ +# Issue: 流式响应 body 在连接关闭后写已释放的 ctx —— heap-use-after-free + +> 状态:fixed(`on_close` 回调 + `ConnLifetime` 引用计数信号) +> 首次观测:PR #87 的 CI(2026-08-26)——`ubuntu-latest / openssl / C++11`(ASan)与 +> `Symbol visibility (shared build)` 两个 lane 同时失败 + +## 现象 + +`libxpp/xpp/http/client_test.cpp` 的 `ClientSendTest.ReadTimeoutTriggersError`: + +- ASan lane:`ERROR: AddressSanitizer: heap-use-after-free`,读 `xHttpCtxWrite` + (`libx/x/http/server.c:1005`),freed by `xHttpConnClose`(`server.c:650`) +- 无 ASan 的 shared-build lane:同测试直接 SegFault + +两个 lane 是**同一个 bug**——一个被 ASan 捕获,一个裸奔。 + +## 根因 + +响应流式路径 `stream_channel_body`(`libxpp/xpp/http/server.h`)跨挂起点持有裸 +`xHttpCtx*`,而它只受 `ServerLifetime` 保护。`ServerLifetime::destroyed` 只在 +`Server::~Server` 置位,**覆盖不到连接级销毁**: + +1. 测试服务器流式返回 1MB body,中途停顿 2s(`mid_body_delay_ms`) +2. 客户端 `read_timeout(1000)`,`send()` 收到 header 成功,但 `bytes()` 读一半超时 + → 客户端断开连接 +3. 服务端 event loop 读到 `n==0` → `xHttpConnClose` → `xHttpStreamDestroy` 释放 + stream(`xHttpCtx` 内嵌在 stream 里,`stream->ctx.internal_ = stream`) +4. 但 spawn 出去的 handler 任务还挂起在 `stream_channel_body` 的 body channel read 上; + 测试服务器稍后推下一个 chunk 时,continuation 拿着悬空的 `xHttpCtx*` 调 + `xHttpCtxWrite` → UAF + +关键点:`xHttpConnClose` 路径**不会**触发 `on_done`(`srv_on_done_cb` 只在 +`conn_dispatch_request` 的正常完成路径调用),所以既有的 per-request 清理钩子 +完全感知不到连接异常关闭。 + +## 修复 + +三处小改动: + +1. **C API**(`libx/x/http/server.h`):`xHttpRouteInfo` 新增 `on_close` 回调 + (`xHttpCloseFunc`),语义是「stream 销毁前调用,覆盖所有 teardown 路径」。 +2. **C 实现**(`libx/x/http/server.c`):`xHttpConnClose` 在 `xHttpStreamDestroy` + 之前调用 `on_close`,arg 用 `route_info->arg`(**不是** `stream->user`,后者可能 + 已被先前的 `on_done` 释放)。 +3. **C++ 层**(`libxpp/xpp/http/server.h`):引入 `ConnLifetime`(`Arc` + `bool + closed`),存于 `ServerImpl::m_conns`(`Arc`,key = ctx 地址)。`on_close` + 回调置 `closed=true`;spawn 任务与 `stream_channel_body` 持有该 Arc,每次碰 ctx + 前先查 `closed`。 + +为什么不给 stream 加引用计数:`xHttpCtxWrite` 不只访问 stream,还访问 `stream->conn` +→ 引用计数必须同时延长 conn(socket/TLS/协议状态)的生命周期,把「连接关闭」的 +资源释放时序和「响应流写完」耦合,复杂且易泄漏。`on_close` 信号方案保持资源释放 +时序不变,只多一个通知,与既有 `ServerLifetime` 模式同构。 + +## 验证 + +- `http_client_test` 13/13 通过(ASan 构建,含 `ReadTimeoutTriggersError`) +- `http_server_test` 16/16 通过 +- `xhttp_test`(C 层)173/173 通过 + +## 复现 + +```bash +cmake -B build -G Ninja -DCMAKE_C_FLAGS="-fsanitize=address -fno-omit-frame-pointer" \ + -DCMAKE_CXX_FLAGS="-fsanitize=address -fno-omit-frame-pointer" +cmake --build build --target http_client_test +./build/libxpp/xpp/http_client_test --gtest_filter=ClientSendTest.ReadTimeoutTriggersError +``` + +修复前:ASan 报 heap-use-after-free(`xHttpCtxWrite` → `stream_channel_body` lambda)。 diff --git a/libx/x/http/server.c b/libx/x/http/server.c index 684c0fc..4a2e5db 100644 --- a/libx/x/http/server.c +++ b/libx/x/http/server.c @@ -647,6 +647,13 @@ void xHttpConnClose(struct xHttpConn_ *conn) { } if (conn->stream) { + /* Notify the route before the stream (and the xHttpCtx embedded in it) + * is freed — the C++ layer uses this to stop in-flight response + * streaming. Use route_info->arg (not stream->user, which may already + * be freed by a prior on_done). */ + if (conn->stream->route_info && conn->stream->route_info->on_close) { + conn->stream->route_info->on_close(&conn->stream->ctx, conn->stream->route_info->arg); + } xHttpStreamDestroy(conn->stream); conn->stream = NULL; } diff --git a/libx/x/http/server.h b/libx/x/http/server.h index f4e5f58..a02b5c4 100644 --- a/libx/x/http/server.h +++ b/libx/x/http/server.h @@ -34,21 +34,39 @@ XDEF_HANDLE(xHttpServer); */ XDEF_HANDLE(xHttpMux); +/** + * @brief Callback invoked right before the connection's request stream is + * torn down (client disconnect, error, or normal close). + * + * Unlike @ref xHttpDoneFunc, this fires on *every* path that ends the + * request's lifetime, including an abrupt close that never reaches + * @ref xHttpDoneFunc. The @p ctx is still valid during the callback but + * must not be used after it returns. + * + * @param ctx Request context (valid only during the callback). + * @param arg User-provided argument (the route's @p arg, NOT the + * per-request user data, which may already be freed). + */ +typedef void (*xHttpCloseFunc)(xHttpCtx *ctx, void *arg); + /** * @brief Route information returned by the resolver. * * Returned by @ref xHttpResolveFunc after the request headers are parsed. * The library calls @p on_request (if non-NULL) right after resolution, * streams the body via @p on_data (if non-NULL), and finally invokes - * @p on_done when the request is fully received. + * @p on_done when the request is fully received. @p on_close (if non-NULL) + * is invoked immediately before the stream is destroyed, on every teardown + * path. * * All callbacks receive @p arg as the user-provided context. */ XDEF_STRUCT(xHttpRouteInfo) { - xHttpInitFunc on_request; /**< Called once after headers (may be NULL) */ - xHttpDataFunc on_data; /**< Per body chunk callback (may be NULL) */ - xHttpDoneFunc on_done; /**< Called when request is complete */ - void *arg; /**< User argument forwarded to callbacks */ + xHttpInitFunc on_request; /**< Called once after headers (may be NULL) */ + xHttpDataFunc on_data; /**< Per body chunk callback (may be NULL) */ + xHttpDoneFunc on_done; /**< Called when request is complete */ + xHttpCloseFunc on_close; /**< Called before the stream is destroyed */ + void *arg; /**< User argument forwarded to callbacks */ }; /** diff --git a/libxpp/xpp/http/server.h b/libxpp/xpp/http/server.h index 152ce9f..237e338 100644 --- a/libxpp/xpp/http/server.h +++ b/libxpp/xpp/http/server.h @@ -69,6 +69,23 @@ class ServerLifetime { bool destroyed = false; }; +/// Per-connection lifetime flag. The C on_close callback sets closed=true +/// right before the stream (and the xHttpCtx embedded in it) is freed. +/// The connection outlives the Server in some paths (client disconnect) +/// and is outlived by it in others, so the response-streaming task holds +/// an Arc and checks this *in addition to* ServerLifetime before touching +/// ctx. +class ConnLifetime { +public: + bool closed = false; +}; + +/// Live connection-lifetime flags keyed by the stream's ctx address. +/// The map owns one strong Arc per in-flight request; the streaming task +/// and the C on_close callback both reach it (through their own Arc of the +/// map, so it survives the Server itself). +using ConnMap = std::unordered_map>; + /// Per-request state: the request-body channel sender (fed by on_data). /// Lives from on_request until the handler's response is written. struct ReqState { @@ -103,6 +120,7 @@ class ServerImpl { Router m_router; ///< all routing + middleware RouterDispatch m_dispatch; std::unordered_map> m_reqs; + Arc m_conns = Arc::make(); Option>> m_stop_resolver; Arc lifetime = Arc::make(); }; @@ -151,24 +169,27 @@ constexpr size_t kStreamChunk = 4096; * or the connection write fails. */ inline Promise stream_channel_body(xHttpCtx *ctx, Arc lifetime, - Arc body) { + Arc conn, Arc body) { auto buf = Arc>::make(); buf->resize(kStreamChunk, '\0'); return body->read(buf->data(), kStreamChunk) - .then([ctx, lifetime, body, buf](ssize_t n) -> Promise { - if (lifetime->destroyed) return xpp::resolve(); // server torn down + .then([ctx, lifetime, conn, body, buf](ssize_t n) -> Promise { + // The connection may be closed (and ctx freed) while we were parked + // on the body channel — stop before touching ctx. + if (lifetime->destroyed || conn->closed) return xpp::resolve(); if (n <= 0) { xHttpCtxEndStream(ctx); // EOF — finalize (Connection: close) return xpp::resolve(); } xErrno rc = xHttpCtxWrite(ctx, buf->data(), static_cast(n)); if (rc != xErrno_Ok) return xpp::resolve(); // connection gone — drop - return stream_channel_body(ctx, lifetime, body); + return stream_channel_body(ctx, lifetime, conn, body); }); } inline Promise write_response(xHttpCtx *ctx, Arc lifetime, - Result r) { + Arc conn, Result r) { + if (lifetime->destroyed || conn->closed) return xpp::resolve(); if (r.is_err()) { xHttpCtxSetStatus(ctx, 500); xHttpCtxSend(ctx, "Internal Server Error", 21); @@ -193,7 +214,7 @@ inline Promise write_response(xHttpCtx *ctx, Arc lifetime, // Channel body — stream it. Arc keeps the reader alive across the // recursive read/write chain. Arc shared = Arc::make(std::move(body)); - return stream_channel_body(ctx, lifetime, shared); + return stream_channel_body(ctx, lifetime, conn, shared); } /* ── C callback trampolines ────────────────────────────────────────── */ @@ -248,13 +269,23 @@ inline int srv_on_request_cb(xHttpCtx *ctx, void *arg) { // C++11 has no move-capture — hold the Request on the heap. auto holder = Arc::make(std::move(req)); auto lifetime = impl->lifetime; // Arc — outlives the server - xpp::spawn([endpoint, lifetime, ctx_key, holder]() -> Promise { + auto conns = impl->m_conns; // Arc — erase this request's entry when done + auto conn = Arc::make(); + conns->emplace(ctx_key, conn); + xpp::spawn([endpoint, lifetime, conns, conn, ctx_key, holder]() -> Promise { return endpoint(std::move(*holder)) - .then([lifetime, ctx_key](Result r) -> Promise { - // Server destroyed while the handler was running (ctx freed) — - // drop the response instead of touching it. - if (lifetime->destroyed) return xpp::resolve(); - return write_response(const_cast(ctx_key), lifetime, std::move(r)); + .then([lifetime, conns, conn, ctx_key](Result r) -> Promise { + // Server destroyed or connection closed while the handler was + // running (ctx freed) — drop the response instead of touching it. + if (lifetime->destroyed || conn->closed) { + conns->erase(ctx_key); + return xpp::resolve(); + } + return write_response(const_cast(ctx_key), lifetime, conn, std::move(r)) + .then([conns, ctx_key]() -> Promise { + conns->erase(ctx_key); + return xpp::resolve(); + }); }); }); return 0; @@ -281,6 +312,21 @@ inline void srv_on_done_cb(xHttpCtx *ctx, void *arg) { } } +inline void srv_on_close_cb(xHttpCtx *ctx, void *arg) { + // arg is the route's info.arg (the RouterDispatch) — NOT the per-request + // user data, which on_done may already have freed. The stream (and the + // ctx embedded in it) is about to be freed; mark the connection closed so + // the response-streaming task stops touching ctx. + auto *d = static_cast(arg); + auto *impl = d ? d->impl : nullptr; + if (!impl) return; + auto it = impl->m_conns->find(ctx); + if (it != impl->m_conns->end()) { + it->second->closed = true; + impl->m_conns->erase(it); + } +} + } // namespace _ /** @@ -436,6 +482,7 @@ class ServerBuilder { impl->m_dispatch.info.on_request = _::srv_on_request_cb; impl->m_dispatch.info.on_data = _::srv_on_data_cb; impl->m_dispatch.info.on_done = _::srv_on_done_cb; + impl->m_dispatch.info.on_close = _::srv_on_close_cb; impl->m_dispatch.info.arg = &impl->m_dispatch; xHttpServerConf sconf = {};