The bus on NATS: both transports behind seams, and the rollout switch #87

Merged
jschoubben merged 40 commits from feat/nats-genesis into main 2026-09-27 17:36:41 +00:00
2 changed files with 75 additions and 0 deletions
Showing only changes of commit e65b3950cc - Show all commits
+39
View File
@@ -146,6 +146,45 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
}, true
}
// NodeConsumer is the durable consumer a node reads its own declaration through.
//
// **Derived from a node existing, and created by the controller, because a host cannot create it.**
// A host's account may subscribe its own declaration subject and publish its own ack subject, and
// reaches no part of the JetStream API — which is correct (the controller is the only writer of
// consumer definitions, design 25 §3) and means the consumer must be waiting before the host binds
// to it. Named after the node, because the node's ack grant is `$JS.ACK.NODES.<node>.>` and a
// consumer named anything else is one the host cannot acknowledge a delivery from.
//
// **No max-deliver, and a long ack wait.** A declaration is settled only after the node has applied
// it and reported, which is minutes on a machine pulling images; and a declaration the mesh cannot
// get a node to accept is not one to dead-letter, because the stream keeps only the newest per node
// anyway — so there is exactly one message per node to redeliver, for as long as that node is away.
func NodeConsumer(node string) Consumer {
return Consumer{
Name: node,
Stream: "NODES",
Filters: []string{"mesh.node." + node + ".declare"},
Push: true,
AckWaitSeconds: 300,
Why: "how " + node + " hears what it should be; last-per-subject, so a node that was away " +
"gets exactly the current declaration and nothing older",
}
}
// AssertNodeConsumers brings every known node's declaration consumer into being.
//
// Asserted on start as well as created at enrolment, for the reason the streams are: a mesh raised
// from a restored backup, or one whose bus was recreated, has node records and no consumers, and a
// node whose consumer is missing hears nothing while everything else about it looks correct.
func AssertNodeConsumers(e Ensurer, nodes []string) error {
for _, n := range nodes {
if err := e.EnsureConsumer(NodeConsumer(n)); err != nil {
return fmt.Errorf("asserting how %s hears its declaration: %w", n, err)
}
}
return nil
}
// AllOverlaps reports subject filters claimed by more than one stream, across the mesh's own and
// every derived one.
//
+36
View File
@@ -117,3 +117,39 @@ func TestASeatsStreamIsNamedAfterTheSeat(t *testing.T) {
t.Fatalf("unexpected stream name %q", name)
}
}
// A node hears its declaration through a consumer only the controller can make.
//
// The three things that would each break it silently: a name other than the node's is one the host
// cannot acknowledge a delivery from, because its ack grant is derived from the node's name; a
// filter other than its own declaration subject is a node reading another's; and a pull consumer is
// one the host cannot bind a channel to without creating something, which it has no authority for.
func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) {
c := NodeConsumer("anchor")
if c.Name != "anchor" {
t.Fatalf("named %q, so the node cannot ack from it: its grant is $JS.ACK.NODES.anchor.>", c.Name)
}
if c.Stream != "NODES" {
t.Fatalf("on stream %q rather than the one declarations live in", c.Stream)
}
if len(c.Filters) != 1 || c.Filters[0] != "mesh.node.anchor.declare" {
t.Fatalf("filters %v, which is not this node's own declaration and nothing else", c.Filters)
}
if !c.Push {
t.Fatal("pulled, which a host cannot do: pulling needs the JetStream API and a host reaches none of it")
}
if c.MaxDeliver != 0 {
t.Fatalf("max-deliver %d: a declaration a node has not taken yet is not one to dead-letter, "+
"because the stream holds exactly one per node", c.MaxDeliver)
}
// And the grant the node actually gets has to match, or none of the above matters.
perms, err := PermissionsFor(Principal{Kind: KindNode, Node: "anchor"})
if err != nil {
t.Fatal(err)
}
// Without the ack grant every declaration a node receives is redelivered for ever; without the
// subscribe grant its consumer delivers to nobody.
has(t, perms.Publish, "$JS.ACK.NODES."+c.Name+".>")
has(t, perms.Subscribe, c.Filters[0])
}