diff --git a/internal/broker/derived.go b/internal/broker/derived.go index 861b12c..97d3174 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -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..>` 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. // diff --git a/internal/broker/derived_test.go b/internal/broker/derived_test.go index bafe7d0..b4529d4 100644 --- a/internal/broker/derived_test.go +++ b/internal/broker/derived_test.go @@ -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]) +}