From 112cbd294de437d73a500391f7d4b5e8a84f59d4 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 21:49:55 +0200 Subject: [PATCH] Per-subject caps on EVENTS, and a comment corrected against the server MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit NATS refuses overlapping streams rather than double-storing, which is the opposite of what the Overlaps comment claimed. The check still earns its place — it names both streams at composition rather than one at apply — and the refusal is what rules out a shared stream beside per-module ones. --- internal/broker/streams.go | 35 +++++++++++++++++++++++++---------- 1 file changed, 25 insertions(+), 10 deletions(-) diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 239a118..b7be332 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -37,8 +37,13 @@ type Stream struct { Name string Subjects []string Retention Retention - // MaxAge in seconds, zero for unbounded. + // MaxAge in seconds, zero for unbounded. Per stream — JetStream has no per-subject age, + // which is why differing retention between modules would mean a stream each. MaxAge int + // MaxMsgsPerSubject caps each subject independently, so one noisy emitter cannot push + // another's events out of a shared stream. Verified: with a cap of 3, ten messages on one + // subject and one on another leave four in the stream, not three. + MaxMsgsPerSubject int // Why is carried into the assertion so an operator reading the server's own state finds the // reason there, rather than only in a repository they may not have. Why string @@ -78,11 +83,14 @@ func MeshStreams() []Stream { Why: "at least once, one builder at a time; a builder that dies mid-build has its message redelivered", }, { - Name: "EVENTS", - Subjects: []string{"mesh.mod.*.event.>"}, - Retention: RetentionLimits, - MaxAge: 7 * 24 * 60 * 60, - Why: "a subscriber that was down catches up; tool traffic under the same prefix is excluded by the event token", + Name: "EVENTS", + Subjects: []string{"mesh.mod.*.event.>"}, + Retention: RetentionLimits, + MaxAge: 7 * 24 * 60 * 60, + MaxMsgsPerSubject: 10000, + Why: "a subscriber that was down catches up; tool traffic under the same prefix is " + + "excluded by the event token; per-subject caps keep a noisy emitter from " + + "evicting a quiet one without splitting the stream", }, } } @@ -113,10 +121,17 @@ func AssertMeshStreams(a Asserter) error { // Overlaps reports subject filters claimed by more than one stream. // -// Two streams matching one subject is not a warning in NATS; it is accepted, and the message is -// stored twice under two retentions. For the mesh that would mean a declaration kept as both -// state and a work queue, acknowledged in one and lingering in the other — a divergence nothing -// reports and nobody would think to look for. So it is refused here, where the set is written. +// **Corrected against the server**: an earlier version of this comment said NATS accepts +// overlapping streams and stores the message twice. It does not — it refuses the second stream +// with "subjects overlap with an existing stream" (verified against nats-server 2.10). The check +// still earns its place, for a different reason: the server's refusal arrives when the controller +// is applying, naming one stream, at a moment when the mesh is half-configured. This one arrives +// where the set is written, names both, and cannot reach a running mesh. +// +// It also decides a design question. Because overlap is refused rather than merged, a shared +// EVENTS stream and a per-module stream cannot coexist — the module's would be refused — so +// "one stream for most, its own for a module that wants different retention" is not an option +// the server allows. It is all of one or all of the other. func Overlaps() []string { seen := map[string]string{} var clashes []string