Set Kafka timestamps and expose record headers - #3407
jeremydmiller merged 7 commits into
Conversation
Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
jeremydmiller
left a comment
There was a problem hiding this comment.
Thanks for this — the direction is right and the gap is real. For a raw-JSON listener there genuinely is no Wolverine metadata on the wire, so stamping SentAt from the broker's record timestamp and surfacing the record headers is exactly the missing capability. I also like that you scoped it to JsonOnlyMapper and left the Wolverine-to-Wolverine path alone, where sent-at and the headers already ride in the envelope protocol. The unit tests are well-chosen.
Two things before I can merge, one of them a genuine security issue.
1. Blocking: the header copy must exclude Wolverine's reserved keys
The copy loop is unfiltered:
foreach (var header in incoming.Headers)
{
if (TryReadHeader(incoming, header.Key, out var headerValue))
{
envelope.Headers.TryAdd(header.Key, headerValue);
}
}Kafka header keys are arbitrary strings chosen by the producer, and on a raw-JSON topic that producer is by definition not Wolverine and not necessarily trusted. EnvelopeConstants reserves a set of those names — id, message-type, saga-id, tenant-id, correlation-id, destination, reply-uri, and others.
Right now that looks inert, because JsonOnlyMapper doesn't promote headers into typed properties. But it does not stay inert. In EnvelopeSerializer.writeHeaders (src/Wolverine/Runtime/Serialization/EnvelopeSerializer.cs:287), the typed properties are written first and then every env.Headers entry is appended verbatim, with no reserved-key filter:
writer.WriteProp(ref count, EnvelopeConstants.TenantIdKey, env.TenantId);
...
foreach (var pair in env.Headers) // <-- unfiltered, and it comes last
{
count++;
writer.Write(pair.Key);
writer.Write(pair.Value);
}and the reader switches reserved keys straight back into typed properties (env.TenantId = value, env.SagaId = value, env.Id = ...), processing entries in order — so the appended copy wins.
So on a durable Kafka listener the sequence is:
- An external producer puts
tenant-id: acmeon the topic. - This PR copies it into
envelope.Headers["tenant-id"]. Still inert —envelope.TenantIdis null. - The durable inbox persists the envelope.
writeHeadersappendstenant-id = acmeafter the (null) typed prop. - On read back,
case EnvelopeConstants.TenantIdKey: env.TenantId = value;fires.
A header supplied by an untrusted external producer has become the envelope's TenantId. saga-id gets you someone else's saga state, and id rewrites Envelope.Id, which is the inbox's dedupe identity. That's a privilege boundary, not a style nit.
It's worth noting your own test picks tenant-id as its example custom header — which is the trap in miniature, and honestly why I want the guard rather than a doc note.
Fix: skip reserved keys when copying. Something like a static readonly FrozenSet<string> of the EnvelopeConstants key names, and if (Reserved.Contains(header.Key)) continue; in the loop. Please add a test that an incoming Kafka tenant-id / saga-id / id header does not land on the envelope, so this can't regress.
(If you'd rather have the reserved-key filter live in EnvelopeSerializer.writeHeaders instead, say so — that arguably hardens every transport at once and I'd take that as a separate PR. But this PR is what makes the path reachable, so it needs to be closed here either way.)
2. The integration timestamp test is vacuous
received.SentAt.ShouldBeGreaterThan(DateTimeOffset.UtcNow.Subtract(5.Minutes()));Envelope.SentAt is initialized to DateTimeOffset.UtcNow at construction (src/Wolverine/Envelope.cs:248), so a freshly received envelope already satisfies this assertion with or without your change. It passes on main today. It's false confidence rather than coverage.
Your unit tests already prove the mapping properly (they null out SentAt to MinValue first and assert the exact value), so this one just needs teeth. The easy strengthening: you're already standing up a raw ProducerBuilder in copies_kafka_headers_onto_the_received_envelope — do the same here, produce with an explicit Timestamp well in the past (say CreateTime an hour ago), and assert the received SentAt equals it. That distinguishes "came from the record" from "came from the constructor", which is the whole point of the change.
Minor
if (incoming.Headers is null) return;— the codebase braces its single-statementifs; please match.- Not a blocker, but the PR title says "Set Kafka timestamps" while the change only reads the incoming record timestamp (Confluent assigns the outgoing one at produce time). "Read" would describe it better.
Happy to take this once (1) is closed and (2) has teeth — and thank you for the tests, the unit coverage in particular is better than a lot of what comes through here.
|
@jakub-petrylak-onerail Can you please take a look at Jeremy's comments and address the same? |
Assert SentAt from an explicit past CreateTime and verify reserved broker headers are not copied onto the envelope. Co-authored-by: Cursor <cursoragent@cursor.com>
|
Thanks for the review — I've addressed Jeremy's feedback:
While working on (2), I think I also uncovered a pre-existing bug in // before
return UseInterop((e, _) => new JsonOnlyMapper(e, options ?? new()));
// after
return UseInterop((_, e) => new JsonOnlyMapper(e, options ?? new()));
With With I confirmed this with the integration test: with the old parameter order, Could you double-check that this reading is correct? If so, the fix is included in the latest push. (Side note: |
… raw-JSON topics (#3442) PublishRawJson had the same UseInterop overload-resolution bug fixed for ReceiveRawJson in GH-3407: the (e, _) lambda bound to the Action customization overload, discarded the JsonOnlyMapper, and left the default KafkaEnvelopeMapper stamping Wolverine protocol headers on every raw-JSON record since 5.0. JsonOnlyMapper now also writes the message-type header for sender liveness pings so raw-JSON listeners can keep recognizing them per the GH-2838 guard. Adds mapper-registration and wire-format regression tests plus an upgrade warning in the Kafka docs. Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
Summary
Set Envelope.SentAt from the native Kafka record timestamp when available.
Copy Kafka record headers into Wolverine envelope headers without overwriting existing values.
Why
Using Kafka’s record timestamp gives consumers an accurate transport-level reference point for measuring end-to-end message latency. It also distinguishes message production or broker-append time from business-event and consumer-processing timestamps.
Exposing Kafka headers through the Wolverine envelope makes broker metadata readily available to handlers and middleware without requiring direct access to Kafka-specific APIs.