From a8fce96050a54b7ce043be270cadab8583f4ed50 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 20:13:56 +0000 Subject: [PATCH 1/6] Observe: label requests from their headers; rejections carry headers --- CHANGELOG.md | 8 +++ docs/production-guide.md | 10 +++- runtime/include/opal/http/beast_transport.h | 4 +- runtime/include/opal/server/middleware.h | 19 ++++++- runtime/src/http/beast_transport.cc | 12 +++-- runtime/src/server/middleware.cc | 16 ++++-- runtime/tests/http/beast_transport_test.cc | 2 + runtime/tests/server/middleware_test.cc | 57 +++++++++++++++++++++ 8 files changed, 115 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 66f9d5aa..028f264e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -90,6 +90,14 @@ policy in [docs/versioning.md](docs/versioning.md). ### Added +- **`Observe` labels requests from their headers** (issue #235). An optional + fifth argument, `RequestLabeler`, maps a request's headers to + `RequestLabels` once, before dispatch; the result rides on both + `RequestStart::labels` and `RequestObservation::labels`, so a sink can key + its completion series on something only the headers carry. A throwing + labeler is logged and labels nothing. `RejectedRequest::headers` carries + the fields the parser read before a 413/431, so transport rejections can + be labeled the same way. - **`SessionRegistry` delivery classes** (issue #227). `SendTo` and every `Broadcast` overload take an optional `DeliveryClass`: `Reliable()` (the default), `Droppable()` (evicted first when a queue is full, dropped diff --git a/docs/production-guide.md b/docs/production-guide.md index ce60c550..0644d005 100644 --- a/docs/production-guide.md +++ b/docs/production-guide.md @@ -378,7 +378,12 @@ transport.Start(opal::server::Chain( // gauge +1 (labeled by s.method/s.target; the operation is not // known until the router runs). }, - nullptr, trusted), + nullptr, trusted, + // Optional: bounded labels read off the headers once, on both + // s.labels and o.labels. + [](const opal::http::Headers& h) -> opal::server::RequestLabels { + return {{"caller", h.Get("x-forwarded-for") ? "edge" : "direct"}}; + }), // Liveness: GET or HEAD /livez -> 200 {"status":"healthy"}. A HEAD // gets that body's Content-Length and none of its octets, framed by // the transport; everything else passes through to the router. @@ -1111,7 +1116,8 @@ the budget without reading may still see a reset, which is inherent to the recipe. These rejections are written by the transport itself, before a handler chain exists, so `Observe` middleware never sees them — set `Options::on_rejected` to observe them (status, peer address, and whatever -the parser got to), wired to the same sink as your `Observe` callbacks. +the parser got to, headers included), wired to the same sink as your +`Observe` callbacks. The connections that die without any response are observable the same way (ADR-0013): set `Options::on_connection_event` for TLS handshake failures diff --git a/runtime/include/opal/http/beast_transport.h b/runtime/include/opal/http/beast_transport.h index cca28d42..2458457d 100644 --- a/runtime/include/opal/http/beast_transport.h +++ b/runtime/include/opal/http/beast_transport.h @@ -42,7 +42,8 @@ class BeastServerTransport : public HttpServerTransport { // A request the transport rejected itself — the over-limit 413/431 answers // written before a handler chain exists, which Observe middleware therefore // never sees (issue #46). method/target may be empty when the request never - // parsed that far (a 431 can fire mid-headers). + // parsed that far (a 431 can fire mid-headers), and headers holds the + // fields that did. // // The `= {}` on the strings is not redundant with their default constructor // (issue #193). Clang's -Wmissing-designated-field-initializers, on under @@ -56,6 +57,7 @@ class BeastServerTransport : public HttpServerTransport { std::string peer_address = {}; std::string method = {}; std::string target = {}; + Headers headers = {}; }; // A connection the transport terminated without delivering a response diff --git a/runtime/include/opal/server/middleware.h b/runtime/include/opal/server/middleware.h index 7fb2cd5e..4d740446 100644 --- a/runtime/include/opal/server/middleware.h +++ b/runtime/include/opal/server/middleware.h @@ -4,6 +4,7 @@ #include #include #include +#include #include #include #include @@ -112,6 +113,13 @@ inline constexpr std::string_view kUnmatchedRoute = "unmatched"; // One served request, as seen from outside the router. FormatAccessLog // (opal/server/access_log.h) renders one as a JSON access-log line. +// Names a sink attaches to a request, read off its headers by Observe's +// labeler. Each value becomes a metrics series or a log field, so a labeler +// maps what it reads onto a bounded vocabulary rather than passing a header +// through. +using RequestLabels = std::map; +using RequestLabeler = std::function; + struct RequestObservation { std::string method; std::string target; @@ -154,6 +162,8 @@ struct RequestObservation { // which carries the same trace id (ADR-0011). Always false under // -fno-exceptions, where the path is compiled out. bool handler_threw = false; + // What Observe's labeler returned for this request; empty without one. + RequestLabels labels = {}; }; // What on_start sees, before the router runs. The Smithy operation is not @@ -162,6 +172,8 @@ struct RequestObservation { struct RequestStart { std::string method; std::string target; + // The same labels the request's completion carries. + RequestLabels labels = {}; }; // Middleware reporting every request to callbacks — the structured-logging @@ -187,10 +199,15 @@ struct RequestStart { // address parse, and string building on every request, which is pure waste // in a chain whose sinks never read client — RecordMetrics deliberately // does not, so the metrics-only composition pays nothing here. +// +// `labeler` runs once per request, before on_start, and its result rides on +// both RequestStart and RequestObservation: the one place a sink reads the +// request's headers. A throwing labeler is logged and labels nothing. Middleware Observe(std::function on_complete, std::function on_start = nullptr, std::function now = nullptr, - std::optional trusted = std::nullopt); + std::optional trusted = std::nullopt, + RequestLabeler labeler = nullptr); // 401 unless the request carries "authorization: Bearer " (scheme // matched case-insensitively per RFC 6750) and validator(token) returns diff --git a/runtime/src/http/beast_transport.cc b/runtime/src/http/beast_transport.cc index 5f165c41..d2191f04 100644 --- a/runtime/src/http/beast_transport.cc +++ b/runtime/src/http/beast_transport.cc @@ -1243,11 +1243,13 @@ struct BeastServerTransport::State : std::enable_shared_from_this { if (!opts.on_rejected) { return; } - const BeastServerTransport::RejectedRequest rejected{ - .status = static_cast(status), - .peer_address = PeerAddressOf(stream), - .method = std::string(partial.method_string()), - .target = std::string(partial.target())}; + BeastServerTransport::RejectedRequest rejected{.status = static_cast(status), + .peer_address = PeerAddressOf(stream), + .method = std::string(partial.method_string()), + .target = std::string(partial.target())}; + for (const auto& field : partial) { + rejected.headers.Add(std::string(field.name_string()), std::string(field.value())); + } try { opts.on_rejected(rejected); } catch (const std::exception& e) { diff --git a/runtime/src/server/middleware.cc b/runtime/src/server/middleware.cc index 5b7e3cca..d1f2b43f 100644 --- a/runtime/src/server/middleware.cc +++ b/runtime/src/server/middleware.cc @@ -162,7 +162,7 @@ Middleware HealthEndpoint(std::string path, std::vector checks) Middleware Observe(std::function on_complete, std::function on_start, std::function now, - std::optional trusted) { + std::optional trusted, RequestLabeler labeler) { if (on_complete == nullptr) { opal::internal::Fatal("opal::server::Observe: on_complete may not be null"); } @@ -170,13 +170,21 @@ Middleware Observe(std::function on_complete, now = [] { return std::chrono::steady_clock::now(); }; } return [on_complete = std::move(on_complete), on_start = std::move(on_start), - now = std::move(now), trusted = std::move(trusted)](http::RequestHandler next) { - return [on_complete, on_start, now, trusted, + now = std::move(now), trusted = std::move(trusted), + labeler = std::move(labeler)](http::RequestHandler next) { + return [on_complete, on_start, now, trusted, labeler, next = std::move(next)](const http::HttpRequest& request) { + RequestLabels labels; + if (labeler != nullptr) { + CallContained([&](const http::Headers& headers) { labels = labeler(headers); }, + request.headers, "Observe labeler"); + } if (on_start != nullptr) { - CallContained(on_start, RequestStart{request.method, request.target}, "Observe on_start"); + CallContained(on_start, RequestStart{request.method, request.target, labels}, + "Observe on_start"); } RequestObservation observation; + observation.labels = std::move(labels); observation.method = request.method; observation.target = request.target; observation.trace_parent = request.headers.Get("traceparent").value_or(""); diff --git a/runtime/tests/http/beast_transport_test.cc b/runtime/tests/http/beast_transport_test.cc index 95171298..e38931ab 100644 --- a/runtime/tests/http/beast_transport_test.cc +++ b/runtime/tests/http/beast_transport_test.cc @@ -501,6 +501,7 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { HttpRequest request; request.method = "POST"; request.target = "/"; + request.headers.Set("user-agent", "games_hub/1.0"); request.body = std::string(64 * 1024, 'x'); const auto response = client.Send(request); // Issue #94: a declared Content-Length over the limit is the deterministic, @@ -517,6 +518,7 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { EXPECT_EQ(rejected[0].status, 413); EXPECT_EQ(rejected[0].method, "POST"); EXPECT_EQ(rejected[0].target, "/"); + EXPECT_EQ(rejected[0].headers.Get("user-agent"), "games_hub/1.0"); EXPECT_EQ(rejected[0].peer_address.rfind("127.0.0.1:", 0), 0u) << rejected[0].peer_address; } { diff --git a/runtime/tests/server/middleware_test.cc b/runtime/tests/server/middleware_test.cc index 21acb543..5bd2be51 100644 --- a/runtime/tests/server/middleware_test.cc +++ b/runtime/tests/server/middleware_test.cc @@ -281,6 +281,63 @@ TEST(ObserveTest, OnStartFiresBeforeDispatch) { EXPECT_EQ(log, (std::vector{"start:POST /tasks", "handler", "complete"})); } +// A caller's name from its User-Agent, the shape of labeler a metrics sink +// wants: one bounded value per request, read off the headers once. +RequestLabels CallerLabel(const http::Headers& headers) { + const std::string agent = headers.Get("user-agent").value_or(""); + return {{"caller", agent.substr(0, agent.find('/'))}}; +} + +TEST(ObserveTest, LabelsFromTheHeadersRideOnStartAndCompletion) { + std::vector starts; + std::vector completions; + auto handler = Chain({Observe([&](const RequestObservation& o) { completions.push_back(o); }, + [&](const RequestStart& s) { starts.push_back(s); }, nullptr, + std::nullopt, CallerLabel)}, + [](const http::HttpRequest&) { return Ok("served"); }); + + http::HttpRequest request; + request.headers.Set("user-agent", "games_hub/1.0"); + (void)handler(request); + + const RequestLabels expected = {{"caller", "games_hub"}}; + ASSERT_EQ(starts.size(), 1u); + EXPECT_EQ(starts[0].labels, expected); + ASSERT_EQ(completions.size(), 1u); + EXPECT_EQ(completions[0].labels, expected); +} + +TEST(ObserveTest, LabelsRideOnTheThrownCompletionToo) { + std::vector completions; + auto handler = Chain({Observe([&](const RequestObservation& o) { completions.push_back(o); }, + nullptr, nullptr, std::nullopt, CallerLabel)}, + [](const http::HttpRequest&) -> http::HttpResponse { + throw std::runtime_error("handler exploded"); + }); + + http::HttpRequest request; + request.headers.Set("user-agent", "mcpserver"); + EXPECT_THROW((void)handler(request), std::runtime_error); + ASSERT_EQ(completions.size(), 1u); + EXPECT_EQ(completions[0].labels, (RequestLabels{{"caller", "mcpserver"}})); +} + +TEST(ObserveTest, AThrowingLabelerIsContainedAndLabelsNothing) { + std::vector completions; + auto handler = Chain({Observe([&](const RequestObservation& o) { completions.push_back(o); }, + nullptr, nullptr, std::nullopt, + [](const http::Headers&) -> RequestLabels { + throw std::runtime_error("labeler down"); + })}, + [](const http::HttpRequest&) { return Ok("served"); }); + + http::HttpResponse response; + EXPECT_NO_THROW(response = handler({})); + EXPECT_EQ(response.body, "served"); + ASSERT_EQ(completions.size(), 1u); + EXPECT_TRUE(completions[0].labels.empty()); +} + TEST(ObserveTest, PairsCompleteWithStartWhenDispatchThrows) { int started = 0; std::vector completions; From 633e8c529a135ef510de07162d81debe8ea259ce Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 20:38:51 +0000 Subject: [PATCH 2/6] Observe labels share MetricLabels' shape; consumer example --- .../access_log_acceptance_test.cc | 29 +++++++++++++++++-- runtime/include/opal/server/middleware.h | 10 +++---- 2 files changed, 31 insertions(+), 8 deletions(-) diff --git a/examples/bazel-consumer/access_log_acceptance_test.cc b/examples/bazel-consumer/access_log_acceptance_test.cc index bdf517fd..a49eb4e8 100644 --- a/examples/bazel-consumer/access_log_acceptance_test.cc +++ b/examples/bazel-consumer/access_log_acceptance_test.cc @@ -60,8 +60,16 @@ struct LogLine { bool handler_threw = false; std::string client; opal::http::DerivedClient::Source client_source = opal::http::DerivedClient::Source::kUnknown; + opal::server::RequestLabels labels; }; +// The caller a request names in its User-Agent's first product token, the +// labeler a sink hands Observe to key its series on who sent a request. +opal::server::RequestLabels CallerOf(const opal::http::Headers& headers) { + const std::string agent = headers.Get("user-agent").value_or(""); + return {{"caller", agent.substr(0, agent.find_first_of("/ "))}}; +} + class AccessLog { public: void Write(const opal::server::RequestObservation& o) { @@ -75,7 +83,8 @@ class AccessLog { .response_bytes = o.response_bytes, .handler_threw = o.handler_threw, .client = o.client.address, - .client_source = o.client.source}); + .client_source = o.client.source, + .labels = o.labels}); } std::vector lines() const { const std::lock_guard lock(mutex_); @@ -155,7 +164,7 @@ class AccessLogAcceptanceTest : public ::testing::Test { // client it was rejected for. opal::server::Observe( [log](const opal::server::RequestObservation& o) { log->Write(o); }, nullptr, - nullptr, observe_trust), + nullptr, observe_trust, CallerOf), opal::server::PerClientRateLimit( [limiter](const std::string& client) { return limiter->Allow(client); }, *trusted, std::chrono::seconds(1))}, @@ -167,13 +176,17 @@ class AccessLogAcceptanceTest : public ::testing::Test { // A request carrying an x-forwarded-for, the way one arrives through a // proxy. The client is the header's entry; the peer is the proxy. - opal::Outcome SendForwarded(const std::string& body) { + opal::Outcome SendForwarded(const std::string& body, + const std::string& user_agent = "") { opal::http::BeastHttpClient raw({.host = "127.0.0.1", .port = transport_->port()}); opal::http::HttpRequest request; request.method = "POST"; request.target = "/tasks"; request.headers.Set("content-type", "application/json"); request.headers.Set("x-forwarded-for", kForwardedClient); + if (!user_agent.empty()) { + request.headers.Set("user-agent", user_agent); + } request.body = body; return raw.Send(request); } @@ -210,6 +223,16 @@ TEST_F(AccessLogAcceptanceTest, TheLoggedClientIsTheBucketTheLimiterKeyedOn) { EXPECT_EQ(lines[0].client_source, opal::http::DerivedClient::Source::kForwarded); } +TEST_F(AccessLogAcceptanceTest, TheLineCarriesTheLabelsItsHeadersEarned) { + const auto served = SendForwarded(R"({"title":"ship it"})", "games_hub/1.0"); + ASSERT_TRUE(served.ok()) << served.error().message(); + ASSERT_EQ(served->status, 200); + + const auto lines = log_->lines(); + ASSERT_EQ(lines.size(), 1u); + EXPECT_EQ(lines[0].labels, (opal::server::RequestLabels{{"caller", "games_hub"}})); +} + TEST_F(AccessLogAcceptanceTest, ARejectionIsLoggedWithTheClientItWasRejectedFor) { // "Whose bucket did that 429 come from" — unanswerable before #202, because // the observation carried only the raw header, which is not what the diff --git a/runtime/include/opal/server/middleware.h b/runtime/include/opal/server/middleware.h index 4d740446..af7c6412 100644 --- a/runtime/include/opal/server/middleware.h +++ b/runtime/include/opal/server/middleware.h @@ -4,10 +4,10 @@ #include #include #include -#include #include #include #include +#include #include #include "opal/http/forwarded.h" @@ -111,15 +111,15 @@ Middleware HealthEndpoint(std::string path = "/health", std::vector; +// through. The same shape as MetricLabels. +using RequestLabels = std::vector>; using RequestLabeler = std::function; +// One served request, as seen from outside the router. FormatAccessLog +// (opal/server/access_log.h) renders one as a JSON access-log line. struct RequestObservation { std::string method; std::string target; From 128033fc4e31e175c8e1b6476e34acb8e12fdf16 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 20:53:46 +0000 Subject: [PATCH 3/6] Rejections never hand credential headers to observers --- CHANGELOG.md | 5 +++-- runtime/include/opal/http/beast_transport.h | 2 +- runtime/src/http/beast_transport.cc | 15 ++++++++++++++- runtime/tests/http/beast_transport_test.cc | 7 +++++++ 4 files changed, 25 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 028f264e..2e11bba5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -96,8 +96,9 @@ policy in [docs/versioning.md](docs/versioning.md). `RequestStart::labels` and `RequestObservation::labels`, so a sink can key its completion series on something only the headers carry. A throwing labeler is logged and labels nothing. `RejectedRequest::headers` carries - the fields the parser read before a 413/431, so transport rejections can - be labeled the same way. + the fields the parser read before a 413/431, less `authorization`, + `proxy-authorization` and `cookie`, so transport rejections can be labeled + the same way. - **`SessionRegistry` delivery classes** (issue #227). `SendTo` and every `Broadcast` overload take an optional `DeliveryClass`: `Reliable()` (the default), `Droppable()` (evicted first when a queue is full, dropped diff --git a/runtime/include/opal/http/beast_transport.h b/runtime/include/opal/http/beast_transport.h index 2458457d..18053056 100644 --- a/runtime/include/opal/http/beast_transport.h +++ b/runtime/include/opal/http/beast_transport.h @@ -43,7 +43,7 @@ class BeastServerTransport : public HttpServerTransport { // written before a handler chain exists, which Observe middleware therefore // never sees (issue #46). method/target may be empty when the request never // parsed that far (a 431 can fire mid-headers), and headers holds the - // fields that did. + // fields that did, less authorization, proxy-authorization and cookie. // // The `= {}` on the strings is not redundant with their default constructor // (issue #193). Clang's -Wmissing-designated-field-initializers, on under diff --git a/runtime/src/http/beast_transport.cc b/runtime/src/http/beast_transport.cc index d2191f04..2aa1d278 100644 --- a/runtime/src/http/beast_transport.cc +++ b/runtime/src/http/beast_transport.cc @@ -102,6 +102,17 @@ HttpRequest ToSmithyRequest(bhttp::request wire) { return request; } +// Headers that carry credentials. A rejection observer is a logging and +// metrics hook, so it never sees them. +bool IsCredentialHeader(std::string_view name) { + for (const std::string_view credential : {"authorization", "proxy-authorization", "cookie"}) { + if (HeaderNameEquals(name, credential)) { + return true; + } + } + return false; +} + // The transport is authoritative for framing: keep_alive()/prepare_payload() // below own these fields, and a handler-set copy would ride along beside // them — a duplicate or conflicting content-length / transfer-encoding is @@ -1248,7 +1259,9 @@ struct BeastServerTransport::State : std::enable_shared_from_this { .method = std::string(partial.method_string()), .target = std::string(partial.target())}; for (const auto& field : partial) { - rejected.headers.Add(std::string(field.name_string()), std::string(field.value())); + if (!IsCredentialHeader(field.name_string())) { + rejected.headers.Add(std::string(field.name_string()), std::string(field.value())); + } } try { opts.on_rejected(rejected); diff --git a/runtime/tests/http/beast_transport_test.cc b/runtime/tests/http/beast_transport_test.cc index e38931ab..01f3f516 100644 --- a/runtime/tests/http/beast_transport_test.cc +++ b/runtime/tests/http/beast_transport_test.cc @@ -502,6 +502,9 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { request.method = "POST"; request.target = "/"; request.headers.Set("user-agent", "games_hub/1.0"); + request.headers.Set("authorization", "Bearer s3cret"); + request.headers.Set("Proxy-Authorization", "Basic cHJveHk="); + request.headers.Set("cookie", "session=s3cret"); request.body = std::string(64 * 1024, 'x'); const auto response = client.Send(request); // Issue #94: a declared Content-Length over the limit is the deterministic, @@ -519,6 +522,10 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { EXPECT_EQ(rejected[0].method, "POST"); EXPECT_EQ(rejected[0].target, "/"); EXPECT_EQ(rejected[0].headers.Get("user-agent"), "games_hub/1.0"); + // Credentials never reach an observer, which may well log what it gets. + EXPECT_FALSE(rejected[0].headers.Get("authorization").has_value()); + EXPECT_FALSE(rejected[0].headers.Get("proxy-authorization").has_value()); + EXPECT_FALSE(rejected[0].headers.Get("cookie").has_value()); EXPECT_EQ(rejected[0].peer_address.rfind("127.0.0.1:", 0), 0u) << rejected[0].peer_address; } { From 8697cdac376b1a247814000cfb3f89b38690b7cd Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 21:09:53 +0000 Subject: [PATCH 4/6] Credential header check via ranges::any_of for clang-tidy --- runtime/src/http/beast_transport.cc | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/runtime/src/http/beast_transport.cc b/runtime/src/http/beast_transport.cc index 2aa1d278..7809b989 100644 --- a/runtime/src/http/beast_transport.cc +++ b/runtime/src/http/beast_transport.cc @@ -105,12 +105,11 @@ HttpRequest ToSmithyRequest(bhttp::request wire) { // Headers that carry credentials. A rejection observer is a logging and // metrics hook, so it never sees them. bool IsCredentialHeader(std::string_view name) { - for (const std::string_view credential : {"authorization", "proxy-authorization", "cookie"}) { - if (HeaderNameEquals(name, credential)) { - return true; - } - } - return false; + static constexpr std::array kCredentials = {"authorization", + "proxy-authorization", "cookie"}; + return std::ranges::any_of(kCredentials, [name](std::string_view credential) { + return HeaderNameEquals(name, credential); + }); } // The transport is authoritative for framing: keep_alive()/prepare_payload() From 4f1b66c487acf54bbb9265c8f0687f9e4fdeef99 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 21:44:27 +0000 Subject: [PATCH 5/6] Rejections carry labeler output, never raw headers --- CHANGELOG.md | 8 ++-- docs/production-guide.md | 5 ++- runtime/include/opal/http/beast_transport.h | 15 +++++-- runtime/src/http/beast_transport.cc | 18 +++------ runtime/tests/http/beast_transport_test.cc | 43 ++++++++++++--------- 5 files changed, 49 insertions(+), 40 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2e11bba5..b87440e4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -95,10 +95,10 @@ policy in [docs/versioning.md](docs/versioning.md). `RequestLabels` once, before dispatch; the result rides on both `RequestStart::labels` and `RequestObservation::labels`, so a sink can key its completion series on something only the headers carry. A throwing - labeler is logged and labels nothing. `RejectedRequest::headers` carries - the fields the parser read before a 413/431, less `authorization`, - `proxy-authorization` and `cookie`, so transport rejections can be labeled - the same way. + labeler is logged and labels nothing. `BeastServerTransport::Options:: + label_rejection` takes the same function for the transport's own 413/431 + rejections: `RejectedRequest::labels` carries what it returns, and the + headers themselves never reach `on_rejected`. - **`SessionRegistry` delivery classes** (issue #227). `SendTo` and every `Broadcast` overload take an optional `DeliveryClass`: `Reliable()` (the default), `Droppable()` (evicted first when a queue is full, dropped diff --git a/docs/production-guide.md b/docs/production-guide.md index 0644d005..63295c8a 100644 --- a/docs/production-guide.md +++ b/docs/production-guide.md @@ -1116,8 +1116,9 @@ the budget without reading may still see a reset, which is inherent to the recipe. These rejections are written by the transport itself, before a handler chain exists, so `Observe` middleware never sees them — set `Options::on_rejected` to observe them (status, peer address, and whatever -the parser got to, headers included), wired to the same sink as your -`Observe` callbacks. +the parser got to), wired to the same sink as your `Observe` callbacks. Set +`Options::label_rejection` to the labeler you hand `Observe` and the +rejections carry the same labels; the raw headers never reach the hook. The connections that die without any response are observable the same way (ADR-0013): set `Options::on_connection_event` for TLS handshake failures diff --git a/runtime/include/opal/http/beast_transport.h b/runtime/include/opal/http/beast_transport.h index 18053056..eea8b8be 100644 --- a/runtime/include/opal/http/beast_transport.h +++ b/runtime/include/opal/http/beast_transport.h @@ -8,6 +8,7 @@ #include #include #include +#include #include #include "opal/http/http1.h" @@ -42,8 +43,9 @@ class BeastServerTransport : public HttpServerTransport { // A request the transport rejected itself — the over-limit 413/431 answers // written before a handler chain exists, which Observe middleware therefore // never sees (issue #46). method/target may be empty when the request never - // parsed that far (a 431 can fire mid-headers), and headers holds the - // fields that did, less authorization, proxy-authorization and cookie. + // parsed that far (a 431 can fire mid-headers). labels is what + // Options::label_rejection made of the headers; the headers themselves + // never leave the transport. // // The `= {}` on the strings is not redundant with their default constructor // (issue #193). Clang's -Wmissing-designated-field-initializers, on under @@ -52,12 +54,15 @@ class BeastServerTransport : public HttpServerTransport { // case above, `{.status = 431}`, would not compile under this repo's // --config=werror, and would have to spell out the empties to say what the // omission already said. + // The shape of opal::server::RequestLabels, so one labeler serves both. + using Labels = std::vector>; + struct RejectedRequest { int status = 0; std::string peer_address = {}; std::string method = {}; std::string target = {}; - Headers headers = {}; + Labels labels = {}; }; // A connection the transport terminated without delivering a response @@ -139,6 +144,10 @@ class BeastServerTransport : public HttpServerTransport { // the same sink as opal::server::Observe so over-limit abuse is // visible in the same metrics. std::function on_rejected{}; + // Projects a rejected request's headers onto RejectedRequest::labels, + // the same function a service hands Observe. Only its result reaches + // on_rejected; a throwing labeler is logged and labels nothing. + std::function label_rejection{}; // Observation hook for connections the transport terminated without a // response (one call per ConnectionEvent; ADR-0013). Same contract as // on_rejected: io thread, concurrent across connections, cheap and diff --git a/runtime/src/http/beast_transport.cc b/runtime/src/http/beast_transport.cc index 7809b989..e8359cf2 100644 --- a/runtime/src/http/beast_transport.cc +++ b/runtime/src/http/beast_transport.cc @@ -102,16 +102,6 @@ HttpRequest ToSmithyRequest(bhttp::request wire) { return request; } -// Headers that carry credentials. A rejection observer is a logging and -// metrics hook, so it never sees them. -bool IsCredentialHeader(std::string_view name) { - static constexpr std::array kCredentials = {"authorization", - "proxy-authorization", "cookie"}; - return std::ranges::any_of(kCredentials, [name](std::string_view credential) { - return HeaderNameEquals(name, credential); - }); -} - // The transport is authoritative for framing: keep_alive()/prepare_payload() // below own these fields, and a handler-set copy would ride along beside // them — a duplicate or conflicting content-length / transfer-encoding is @@ -1257,10 +1247,12 @@ struct BeastServerTransport::State : std::enable_shared_from_this { .peer_address = PeerAddressOf(stream), .method = std::string(partial.method_string()), .target = std::string(partial.target())}; - for (const auto& field : partial) { - if (!IsCredentialHeader(field.name_string())) { - rejected.headers.Add(std::string(field.name_string()), std::string(field.value())); + if (opts.label_rejection) { + Headers headers; + for (const auto& field : partial) { + headers.Add(std::string(field.name_string()), std::string(field.value())); } + InvokeCompletion("label_rejection", [&] { rejected.labels = opts.label_rejection(headers); }); } try { opts.on_rejected(rejected); diff --git a/runtime/tests/http/beast_transport_test.cc b/runtime/tests/http/beast_transport_test.cc index 01f3f516..04f0bcff 100644 --- a/runtime/tests/http/beast_transport_test.cc +++ b/runtime/tests/http/beast_transport_test.cc @@ -488,14 +488,18 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { std::vector rejected; ConnectionEventRecorder events; // must stay empty: on_rejected already // observed this connection (ADR-0013) - BeastServerTransport server( - BeastServerTransport::Options{.max_body_bytes = 1024, - .on_rejected = - [&](const BeastServerTransport::RejectedRequest& r) { - const std::lock_guard lock(mutex); - rejected.push_back(r); - }, - .on_connection_event = events.Hook()}); + BeastServerTransport server(BeastServerTransport::Options{ + .max_body_bytes = 1024, + .on_rejected = + [&](const BeastServerTransport::RejectedRequest& r) { + const std::lock_guard lock(mutex); + rejected.push_back(r); + }, + .label_rejection = [](const Headers& headers) -> BeastServerTransport::Labels { + const std::string agent = headers.Get("user-agent").value_or(""); + return {{"caller", agent.substr(0, agent.find('/'))}}; + }, + .on_connection_event = events.Hook()}); ASSERT_TRUE(server.Start([](const HttpRequest&) { return HttpResponse{}; }).ok()); SocketHttpClient client("127.0.0.1", server.port()); HttpRequest request; @@ -503,8 +507,6 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { request.target = "/"; request.headers.Set("user-agent", "games_hub/1.0"); request.headers.Set("authorization", "Bearer s3cret"); - request.headers.Set("Proxy-Authorization", "Basic cHJveHk="); - request.headers.Set("cookie", "session=s3cret"); request.body = std::string(64 * 1024, 'x'); const auto response = client.Send(request); // Issue #94: a declared Content-Length over the limit is the deterministic, @@ -521,11 +523,9 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { EXPECT_EQ(rejected[0].status, 413); EXPECT_EQ(rejected[0].method, "POST"); EXPECT_EQ(rejected[0].target, "/"); - EXPECT_EQ(rejected[0].headers.Get("user-agent"), "games_hub/1.0"); - // Credentials never reach an observer, which may well log what it gets. - EXPECT_FALSE(rejected[0].headers.Get("authorization").has_value()); - EXPECT_FALSE(rejected[0].headers.Get("proxy-authorization").has_value()); - EXPECT_FALSE(rejected[0].headers.Get("cookie").has_value()); + // Only what the labeler projects out of the headers reaches the + // observer; the headers themselves, credentials included, never do. + EXPECT_EQ(rejected[0].labels, (BeastServerTransport::Labels{{"caller", "games_hub"}})); EXPECT_EQ(rejected[0].peer_address.rfind("127.0.0.1:", 0), 0u) << rejected[0].peer_address; } { @@ -866,9 +866,15 @@ TEST(BeastTransportTest, OversizedHeadersReadA431) { std::mutex mutex; std::vector rejected; BeastServerTransport server(BeastServerTransport::Options{ - .max_header_bytes = 1024, .on_rejected = [&](const BeastServerTransport::RejectedRequest& r) { - const std::lock_guard lock(mutex); - rejected.push_back(r); + .max_header_bytes = 1024, + .on_rejected = + [&](const BeastServerTransport::RejectedRequest& r) { + const std::lock_guard lock(mutex); + rejected.push_back(r); + }, + // A throwing labeler costs the labels, never the observation. + .label_rejection = [](const Headers&) -> BeastServerTransport::Labels { + throw std::runtime_error("labeler down"); }}); ASSERT_TRUE(server.Start([](const HttpRequest&) { return HttpResponse{200, {}, ""}; }).ok()); SocketHttpClient client("127.0.0.1", server.port()); @@ -887,6 +893,7 @@ TEST(BeastTransportTest, OversizedHeadersReadA431) { const std::lock_guard lock(mutex); ASSERT_EQ(rejected.size(), 1u); EXPECT_EQ(rejected[0].status, 431); + EXPECT_TRUE(rejected[0].labels.empty()); EXPECT_EQ(rejected[0].peer_address.rfind("127.0.0.1:", 0), 0u) << rejected[0].peer_address; } server.Stop(); From 4d4f206ddeddedc47fd9628f197f03aba184c780 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 21:52:18 +0000 Subject: [PATCH 6/6] label_rejection goes last in Options --- runtime/include/opal/http/beast_transport.h | 8 ++++---- runtime/tests/http/beast_transport_test.cc | 4 ++-- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/runtime/include/opal/http/beast_transport.h b/runtime/include/opal/http/beast_transport.h index eea8b8be..6fbc0f5e 100644 --- a/runtime/include/opal/http/beast_transport.h +++ b/runtime/include/opal/http/beast_transport.h @@ -144,10 +144,6 @@ class BeastServerTransport : public HttpServerTransport { // the same sink as opal::server::Observe so over-limit abuse is // visible in the same metrics. std::function on_rejected{}; - // Projects a rejected request's headers onto RejectedRequest::labels, - // the same function a service hands Observe. Only its result reaches - // on_rejected; a throwing labeler is logged and labels nothing. - std::function label_rejection{}; // Observation hook for connections the transport terminated without a // response (one call per ConnectionEvent; ADR-0013). Same contract as // on_rejected: io thread, concurrent across connections, cheap and @@ -218,6 +214,10 @@ class BeastServerTransport : public HttpServerTransport { // h2-only) is refused at the handshake. mTLS is tracked with #90. std::string tls_certificate_chain_pem{}; std::string tls_private_key_pem{}; + // Projects a rejected request's headers onto RejectedRequest::labels, + // the same function a service hands Observe. Only its result reaches + // on_rejected; a throwing labeler is logged and labels nothing. + std::function label_rejection{}; }; BeastServerTransport() : BeastServerTransport(Options{}) {} diff --git a/runtime/tests/http/beast_transport_test.cc b/runtime/tests/http/beast_transport_test.cc index 04f0bcff..194c2b71 100644 --- a/runtime/tests/http/beast_transport_test.cc +++ b/runtime/tests/http/beast_transport_test.cc @@ -495,11 +495,11 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) { const std::lock_guard lock(mutex); rejected.push_back(r); }, + .on_connection_event = events.Hook(), .label_rejection = [](const Headers& headers) -> BeastServerTransport::Labels { const std::string agent = headers.Get("user-agent").value_or(""); return {{"caller", agent.substr(0, agent.find('/'))}}; - }, - .on_connection_event = events.Hook()}); + }}); ASSERT_TRUE(server.Start([](const HttpRequest&) { return HttpResponse{}; }).ok()); SocketHttpClient client("127.0.0.1", server.port()); HttpRequest request;