diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 4df1a4a..5bbee69 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -184,7 +184,20 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { // without the other is refused by the server with a message that does not say which half is // missing. if c.Queue != "" || c.Push { - want.DeliverSubject = "_DELIVER." + c.Name + // **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146). + // A push consumer delivers onto an ordinary subject, and everything subscribed to that + // subject gets a copy. The controller holds a consumer called `controller` on CONTROL and + // another called `controller` on EVENTS, and both were given `_DELIVER.controller` — so the + // one process, holding both subscriptions, acted on every message twice. It enrolled a + // joining machine twice from one request, minting a second credential that replaced the one + // the machine had just been given; the same doubling applied to every report and every + // event the controller follows. + // + // The stream is in the name because the pair is what identifies a consumer — the server + // scopes a durable's name to its stream, and this subject is the only place that scoping + // was dropped. Already within what the controller may subscribe (`_DELIVER.controller.>`), + // so no permission moves. + want.DeliverSubject = DeliverSubjectFor(c) } switch _, err := j.js.ConsumerInfo(c.Stream, c.Name); { diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 7d5a913..8e2c8dc 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -269,7 +269,13 @@ func PermissionsFor(p Principal) (Permissions, error) { "mesh.control." + p.Node + ".>", "$JS.API.CONSUMER.INFO.NODES." + p.Node, } - sub = []string{"mesh.node." + p.Node + ".declare", "_DELIVER." + p.Node} + // The deliver subject carries the stream as well as the consumer's name, so what a + // subscriber is permitted has to carry it too (novox/hq 04-ISSUES/146). The bare name + // stays: an existing consumer keeps delivering where it always did until the controller's + // next assertion moves it, and a permission that only allowed the new shape would refuse + // every node in the mesh for exactly as long as that took. + sub = []string{"mesh.node." + p.Node + ".declare", + "_DELIVER." + p.Node, "_DELIVER." + p.Node + ".>"} case KindModule: // 1. Its own namespace: it publishes its events there and serves its tools there. Nothing @@ -323,7 +329,7 @@ func PermissionsFor(p Principal) (Permissions, error) { // take work over the new bus was refused the asking (2026-09-28). worker := "SEAT_" + upperSnake(s.Name) + "_worker" stream := seatStreamName(s.Name) - sub = append(sub, "_DELIVER."+worker) + sub = append(sub, "_DELIVER."+worker, "_DELIVER."+worker+".>") pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") for _, a := range s.Accepts { sub = append(sub, seatSubject(s, "accept", a)) diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 6408db5..f6dcd2d 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -94,6 +94,24 @@ func MeshStreams() []Stream { } } +// DeliverSubjectFor is where a push consumer's messages land. +// +// **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146). A push +// consumer delivers onto an ordinary subject, and everything subscribed to that subject gets a +// copy. The controller holds a consumer called `controller` on CONTROL and another called +// `controller` on EVENTS; while both were given `_DELIVER.controller`, the one process holding +// both subscriptions acted on every message twice — a joining machine was enrolled twice from one +// request, and the second enrolment minted a credential that replaced the one the machine had just +// been handed. Every report and every followed event doubled the same way, silently: nothing is +// redelivered, no count is wrong, the work simply happens twice. +// +// The stream belongs in it because the pair is what identifies a consumer — the server scopes a +// durable's name to its stream, and this subject was the one place that scoping was dropped. It +// stays inside what a controller may already subscribe (`_DELIVER.controller.>`). +func DeliverSubjectFor(c Consumer) string { + return "_DELIVER." + c.Name + "." + c.Stream +} + // An Asserter is the part of a JetStream connection stream assertion needs. Narrow on purpose: it // keeps this testable without a server, and keeps the client library out of everything that only // wants to know what the streams are. diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index ca450a3..9a49289 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -233,3 +233,30 @@ func containsStep(steps []string, want string) bool { } return false } + +// **Two consumers may share a name, and must not share a delivery subject** (novox/hq +// 04-ISSUES/146). +// +// A push consumer delivers onto an ordinary subject and everything subscribed to it gets a copy. +// The controller holds a consumer called `controller` on CONTROL and another called `controller` on +// EVENTS; while both were given `_DELIVER.controller`, the one process holding both subscriptions +// acted on every message twice — a joining machine enrolled twice from one request, with the second +// enrolment minting a credential that replaced the one the machine had just been handed. +// +// Checked here rather than against a server because it is a property of what the mesh asks for, and +// because the failure it produces is silent: every count is right, nothing is redelivered, and the +// work simply happens twice. +func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) { + seen := map[string]string{} + for _, c := range MeshConsumers() { + if !c.Push && c.Queue == "" { + continue + } + subject := DeliverSubjectFor(c) + if other, taken := seen[subject]; taken { + t.Errorf("%s on %s and %s deliver onto %s, so whoever holds both acts on every "+ + "message twice", c.Name, c.Stream, other, subject) + } + seen[subject] = c.Name + " on " + c.Stream + } +} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 03cf35a..faebd3b 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -34,11 +34,11 @@ accounts { } } { user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: { publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] } - subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] } + subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] } } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } - subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } + subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {