Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,15 @@ 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. `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
Expand Down
11 changes: 9 additions & 2 deletions docs/production-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -1111,7 +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), 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
Expand Down
29 changes: 26 additions & 3 deletions examples/bazel-consumer/access_log_acceptance_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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<LogLine> lines() const {
const std::lock_guard<std::mutex> lock(mutex_);
Expand Down Expand Up @@ -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))},
Expand All @@ -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<opal::http::HttpResponse> SendForwarded(const std::string& body) {
opal::Outcome<opal::http::HttpResponse> 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);
}
Expand Down Expand Up @@ -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
Expand Down
13 changes: 12 additions & 1 deletion runtime/include/opal/http/beast_transport.h
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
#include <string>
#include <string_view>
#include <thread>
#include <utility>
#include <vector>

#include "opal/http/http1.h"
Expand Down Expand Up @@ -42,7 +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).
// 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
Expand All @@ -51,11 +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<std::pair<std::string, std::string>>;

struct RejectedRequest {
int status = 0;
std::string peer_address = {};
std::string method = {};
std::string target = {};
Labels labels = {};
};

// A connection the transport terminated without delivering a response
Expand Down Expand Up @@ -207,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<Labels(const Headers&)> label_rejection{};
};

BeastServerTransport() : BeastServerTransport(Options{}) {}
Expand Down
19 changes: 18 additions & 1 deletion runtime/include/opal/server/middleware.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include <optional>
#include <string>
#include <string_view>
#include <utility>
#include <vector>

#include "opal/http/forwarded.h"
Expand Down Expand Up @@ -110,6 +111,13 @@ Middleware HealthEndpoint(std::string path = "/health", std::vector<ReadinessChe
// spell the pivot identically.
inline constexpr std::string_view kUnmatchedRoute = "unmatched";

// 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. The same shape as MetricLabels.
using RequestLabels = std::vector<std::pair<std::string, std::string>>;
using RequestLabeler = std::function<RequestLabels(const http::Headers&)>;

// 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 {
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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<void(const RequestObservation&)> on_complete,
std::function<void(const RequestStart&)> on_start = nullptr,
std::function<std::chrono::steady_clock::time_point()> now = nullptr,
std::optional<http::TrustedProxies> trusted = std::nullopt);
std::optional<http::TrustedProxies> trusted = std::nullopt,
RequestLabeler labeler = nullptr);

// 401 unless the request carries "authorization: Bearer <token>" (scheme
// matched case-insensitively per RFC 6750) and validator(token) returns
Expand Down
16 changes: 11 additions & 5 deletions runtime/src/http/beast_transport.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1243,11 +1243,17 @@ struct BeastServerTransport::State : std::enable_shared_from_this<State> {
if (!opts.on_rejected) {
return;
}
const BeastServerTransport::RejectedRequest rejected{
.status = static_cast<int>(status),
.peer_address = PeerAddressOf(stream),
.method = std::string(partial.method_string()),
.target = std::string(partial.target())};
BeastServerTransport::RejectedRequest rejected{.status = static_cast<int>(status),
.peer_address = PeerAddressOf(stream),
.method = std::string(partial.method_string()),
.target = std::string(partial.target())};
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);
} catch (const std::exception& e) {
Expand Down
16 changes: 12 additions & 4 deletions runtime/src/server/middleware.cc
Original file line number Diff line number Diff line change
Expand Up @@ -162,21 +162,29 @@ Middleware HealthEndpoint(std::string path, std::vector<ReadinessCheck> checks)
Middleware Observe(std::function<void(const RequestObservation&)> on_complete,
std::function<void(const RequestStart&)> on_start,
std::function<std::chrono::steady_clock::time_point()> now,
std::optional<http::TrustedProxies> trusted) {
std::optional<http::TrustedProxies> trusted, RequestLabeler labeler) {
if (on_complete == nullptr) {
opal::internal::Fatal("opal::server::Observe: on_complete may not be null");
}
if (now == nullptr) {
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("");
Expand Down
38 changes: 27 additions & 11 deletions runtime/tests/http/beast_transport_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -488,19 +488,25 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) {
std::vector<BeastServerTransport::RejectedRequest> 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<std::mutex> 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<std::mutex> 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('/'))}};
}});
ASSERT_TRUE(server.Start([](const HttpRequest&) { return HttpResponse{}; }).ok());
SocketHttpClient client("127.0.0.1", server.port());
HttpRequest request;
request.method = "POST";
request.target = "/";
request.headers.Set("user-agent", "games_hub/1.0");
request.headers.Set("authorization", "Bearer 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,
Expand All @@ -517,6 +523,9 @@ TEST(BeastTransportTest, OversizedDeclaredBodyReadsA413) {
EXPECT_EQ(rejected[0].status, 413);
EXPECT_EQ(rejected[0].method, "POST");
EXPECT_EQ(rejected[0].target, "/");
// 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;
}
{
Expand Down Expand Up @@ -857,9 +866,15 @@ TEST(BeastTransportTest, OversizedHeadersReadA431) {
std::mutex mutex;
std::vector<BeastServerTransport::RejectedRequest> rejected;
BeastServerTransport server(BeastServerTransport::Options{
.max_header_bytes = 1024, .on_rejected = [&](const BeastServerTransport::RejectedRequest& r) {
const std::lock_guard<std::mutex> lock(mutex);
rejected.push_back(r);
.max_header_bytes = 1024,
.on_rejected =
[&](const BeastServerTransport::RejectedRequest& r) {
const std::lock_guard<std::mutex> 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());
Expand All @@ -878,6 +893,7 @@ TEST(BeastTransportTest, OversizedHeadersReadA431) {
const std::lock_guard<std::mutex> 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();
Expand Down
Loading
Loading