diff --git a/examples/foundation-first-node-nats.lock b/examples/foundation-first-node-nats.lock index 78bfaa1..769b128 100644 --- a/examples/foundation-first-node-nats.lock +++ b/examples/foundation-first-node-nats.lock @@ -161,7 +161,7 @@ "type": "file", "path": "/var/lib/mesh-bus-conf/accounts.conf", "mode": "0600", - "content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"$KV.mesh-controller_calls.>\", \"$KV.mesh-controller_hand-acts.>\", \"$KV.mesh-controller_conditions.>\", \"$KV.mesh-controller_condition-history.>\", \"$KV.mesh-controller_lease.>\", \"$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>\", \"$KV.SEAT_NODE_BUILD_AGENT_cancelled.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\", \"mesh.seat.mesh-controller.event.condition-raised\", \"mesh.seat.mesh-controller.event.condition-changed\", \"mesh.seat.mesh-controller.event.condition-cleared\", \"mesh.seat.mesh-controller.event.doctor-heartbeat\", \"mesh.seat.mesh-controller.event.secret-replaced\", \"$SRV.INFO\", \"mesh.seat.node-intrusion-prevention.tool.banned.*\"] }\n subscribe: { allow: [\"$JS.API.>\", \"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>\", \"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n" + "content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"$KV.mesh-controller_calls.>\", \"$KV.mesh-controller_hand-acts.>\", \"$KV.mesh-controller_conditions.>\", \"$KV.mesh-controller_condition-history.>\", \"$KV.mesh-controller_lease.>\", \"$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>\", \"$KV.SEAT_NODE_BUILD_AGENT_cancelled.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\", \"mesh.seat.mesh-controller.event.condition-raised\", \"mesh.seat.mesh-controller.event.condition-changed\", \"mesh.seat.mesh-controller.event.condition-cleared\", \"mesh.seat.mesh-controller.event.doctor-heartbeat\", \"mesh.seat.mesh-controller.event.secret-replaced\", \"mesh.seat.mesh-controller.event.healer-acted\", \"$SRV.INFO\", \"mesh.seat.node-intrusion-prevention.tool.banned.*\"] }\n subscribe: { allow: [\"$JS.API.>\", \"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>\", \"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n" }, { "id": "broker", diff --git a/internal/link/asked_test.go b/internal/link/asked_test.go new file mode 100644 index 0000000..8211c4b --- /dev/null +++ b/internal/link/asked_test.go @@ -0,0 +1,130 @@ +package link + +import ( + "context" + "strings" + "testing" + "time" +) + +// The `report` verb (novox/hq to-be 45 §6), what healer H1 asks a machine whose send went unreported: +// the node says again what it last applied, as an ordinary report through the one queue. + +// **Asked, the reconcile's account is said even when it is not news** — the mesh asked because it does +// not know, and silence would be the very fault it is asking about. +func TestAskedToReportSaysWhatWasLastAppliedEvenWhenNotNews(t *testing.T) { + m, _ := aMember(t) + l := newQuietLink() + var lines []string + q := aQueue(m, (&applies{}).apply, &keptInMemory{}) + q.Say = func(s string) { lines = append(lines, s) } + q.Reconcile = func(context.Context) (Report, bool) { + return Report{Declared: "d-kept", Order: Order{Epoch: 7, Sequence: 12}, Applied: []string{"a"}}, true + } + q.News = func(Report) bool { return false } + q.attach(context.Background(), l) + ctx, stop := context.WithCancel(context.Background()) + worker := runWorker(ctx, q) + + q.ReconcileDue() // an ordinary reconcile: not news, not said + time.Sleep(50 * time.Millisecond) + if n := len(l.said()); n != 0 { + t.Fatalf("an ordinary reconcile that is not news said %d report(s)", n) + } + q.ReportAsked() + eventually(t, "asked, the node said nothing", func() bool { return len(l.said()) == 1 }) + stop() + <-worker + r := l.said()[0] + if r.Declared != "d-kept" || r.Sequence != 12 || r.Epoch != 7 || r.ReportSequence == 0 || r.Node != m.Node { + t.Fatalf("the account said is not the one kept, in order: %+v", r) + } + if !strings.Contains(strings.Join(lines, "\n"), "asked by the mesh") { + t.Errorf("being asked is not said in the journal: %q", lines) + } +} + +// **A delivery waiting answers the question better**: it is applied, once, and its report is the answer. +func TestAskedWhileADeliveryWaitsIsThatDeliverysReport(t *testing.T) { + m, key := aMember(t) + l := newQuietLink() + a := &applies{} + q := aQueue(m, a.apply, &keptInMemory{}) + reconciled := 0 + q.Reconcile = func(context.Context) (Report, bool) { reconciled++; return Report{}, true } + q.attach(context.Background(), l) + q.Deliver(ordered(t, key, "new", 7, 13)) + q.ReportAsked() + ctx, stop := context.WithCancel(context.Background()) + worker := runWorker(ctx, q) + eventually(t, "the delivery was never reported", func() bool { return len(l.said()) == 1 }) + time.Sleep(50 * time.Millisecond) + stop() + <-worker + if reconciled != 0 || len(l.said()) != 1 || l.said()[0].Sequence != 13 { + t.Fatalf("reconciled %d, said %+v", reconciled, l.said()) + } +} + +// askingLink is a quiet link the mesh can ask. +type askingLink struct { + *quietLink + asks chan struct{} +} + +func (l askingLink) AskedToReport() <-chan struct{} { return l.asks } + +// **The link hands the question to the queue**: asked on the link, the account is said on it. +func TestALinkThatIsAskedHandsItToTheQueue(t *testing.T) { + m, _ := aMember(t) + l := askingLink{quietLink: newQuietLink(), asks: make(chan struct{}, 1)} + q := aQueue(m, (&applies{}).apply, &keptInMemory{}) + q.Reconcile = func(context.Context) (Report, bool) { return Report{Declared: "d-kept"}, true } + ctx, stop := context.WithCancel(context.Background()) + defer stop() + worker := runWorker(ctx, q) + served := make(chan error, 1) + go func() { served <- serve(ctx, l, m, q, nil, time.Second) }() + l.asks <- struct{}{} + eventually(t, "asked on the link, nothing was said", func() bool { + for _, r := range l.said() { + if r.Declared == "d-kept" { + return true + } + } + return false + }) + stop() + <-served + <-worker +} + +// **Against a real server**: the question on the node's own subject reaches it, and another's does not. +func TestNatsTheMeshAsksANodeToReport(t *testing.T) { + conn, js := aBus(t) + l := &natsLink{conn: conn, js: js, node: "asked"} + l.hearAsks() + defer l.Close() + if err := conn.Publish(AskReportSubject("another"), []byte(`{}`)); err != nil { + t.Fatal(err) + } + if err := conn.Publish(AskReportSubject("asked"), []byte(`{"by":"the controller's healer H1"}`)); err != nil { + t.Fatal(err) + } + if err := conn.Flush(); err != nil { + t.Fatal(err) + } + select { + case <-l.AskedToReport(): + case <-time.After(5 * time.Second): + t.Fatal("the question on the node's own subject never reached it") + } + select { + case <-l.AskedToReport(): + t.Fatal("asked twice: another node's question reached this one") + case <-time.After(200 * time.Millisecond): + } + if AskReportSubject("asked") != "mesh.node.asked.ask.report" { + t.Fatalf("the subject %q is not the one the controller asks on and grants", AskReportSubject("asked")) + } +} diff --git a/internal/link/hearing.go b/internal/link/hearing.go index 173ea12..e1f4503 100644 --- a/internal/link/hearing.go +++ b/internal/link/hearing.go @@ -38,6 +38,13 @@ type Link interface { Close() } +// Asked is a link that hears the mesh asking this node to say again what it last applied (novox/hq +// to-be 45 §6, the `report` verb). Optional: a link without it is never asked, and the mesh then sends +// the current declaration again instead. +type Asked interface { + AskedToReport() <-chan struct{} +} + // Declaration is one thing the mesh told this node to be. // // **Handled, once — after the report is published.** A node that dies between applying and diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go index 0bad8e9..3207772 100644 --- a/internal/link/hearing_nats.go +++ b/internal/link/hearing_nats.go @@ -38,6 +38,11 @@ const EnrolSubject = "mesh.control.enrol" // for. func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" } +// AskReportSubject is where the mesh asks this node to say again what it last applied (novox/hq to-be 45 +// §6): its own, on core NATS and off any stream. The answer is an ordinary report on its report subject, +// the one thing it already says — nobody's inbox is answered. +func AskReportSubject(node string) string { return "mesh.node." + node + ".ask.report" } + // natsURL is a bus address as the client wants it. A membership records host and port, because that // is what genesis sealed into it and what the other transport takes; the scheme is this transport's // own business. @@ -53,6 +58,8 @@ type natsLink struct { conn *nats.Conn js nats.JetStreamContext sub *nats.Subscription + asking *nats.Subscription + asked chan struct{} node string arrived chan Declaration lost chan error @@ -106,6 +113,7 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e conn: conn, js: js, node: m.Node, arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1), + asked: make(chan struct{}, 1), } // Bound to the consumer the controller made for this node, named after the node because that is @@ -124,6 +132,8 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e } l.sub = sub + l.hearAsks() + conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) { select { case l.lost <- fmt.Errorf("the link dropped: %w", err): @@ -163,13 +173,36 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e return l, nil } +// hearAsks listens for the mesh asking what this node last applied (the `report` verb). One question +// waiting is enough: a second asked before the first is taken is the same question. A bus whose user +// list is older than the grant refuses the subscription asynchronously — said by the client, and the +// mesh then sends again instead of asking; nothing here depends on being asked. +func (l *natsLink) hearAsks() { + if l.asked == nil { + l.asked = make(chan struct{}, 1) + } + asking, err := l.conn.Subscribe(AskReportSubject(l.node), func(*nats.Msg) { + select { + case l.asked <- struct{}{}: + default: + } + }) + if err == nil { + l.asking = asking + } +} + func (l *natsLink) Declarations() <-chan Declaration { return l.arrived } +func (l *natsLink) AskedToReport() <-chan struct{} { return l.asked } func (l *natsLink) Lost() <-chan error { return l.lost } func (l *natsLink) Close() { if l.sub != nil { _ = l.sub.Unsubscribe() } + if l.asking != nil { + _ = l.asking.Unsubscribe() + } if l.conn != nil { l.conn.Close() } diff --git a/internal/link/queue.go b/internal/link/queue.go index 51c0bfc..e50899e 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -64,8 +64,11 @@ type Queue struct { delivering bool waiting []Declaration // delivered and not yet taken, in the order they arrived due bool // a reconcile asked for - bus Bus // the link open now; nil while there is none - unasked *Report // a reconcile's report not yet said, made while there was no link + // asked is the mesh asking this node to say what it last applied (the `report` verb, to-be 45 §6): + // the reconcile it enqueues says its report whether or not it is news. + asked bool + bus Bus // the link open now; nil while there is none + unasked *Report // a reconcile's report not yet said, made while there was no link // saying is held while a report is made and said, so the order they leave in is the order they // are made in — whoever says them, the worker or a link opening. @@ -116,6 +119,19 @@ func (q *Queue) ReconcileDue() { q.signal() } +// ReportAsked is the mesh asking this node to say again what it last applied (novox/hq to-be 45 §6, the +// `report` verb its healer H1 asks when a send went unreported). **An enqueue, like everything else**: a +// reconcile, whose report — the account of the declaration this node keeps, the last it applied — is +// said whether or not it is news. A delivery waiting meanwhile is applied instead and reported, which +// answers the question better. Asked again before the worker takes it is the same question. +func (q *Queue) ReportAsked() { + q.init() + q.mu.Lock() + q.due, q.asked = true, true + q.mu.Unlock() + q.signal() +} + func (q *Queue) signal() { select { case q.wake <- struct{}{}: @@ -147,8 +163,8 @@ func (q *Queue) step(ctx context.Context) bool { q.mu.Unlock() return false } - batch, due := q.waiting, q.due - q.waiting, q.due = nil, false + batch, due, asked := q.waiting, q.due, q.asked + q.waiting, q.due, q.asked = nil, false, false q.delivering = len(batch) > 0 q.mu.Unlock() defer func() { @@ -159,7 +175,7 @@ func (q *Queue) step(ctx context.Context) bool { }() if len(batch) == 0 { - q.reconcile(ctx) + q.reconcile(ctx, asked) return true } @@ -210,14 +226,27 @@ func (q *Queue) step(ctx context.Context) bool { return true } -func (q *Queue) reconcile(ctx context.Context) { +func (q *Queue) reconcile(ctx context.Context, asked bool) { if q.Reconcile == nil { + if asked { + q.Say("asked to say what this node last applied, and it holds nothing to reconcile with") + } return } report, ok := q.Reconcile(ctx) - if !ok || q.News == nil || !q.News(report) { + if !ok { + if asked { + q.Say("asked to say what this node last applied, and the reconcile could not say it") + } return } + if !asked && (q.News == nil || !q.News(report)) { + return + } + if asked { + q.Say(fmt.Sprintf("asked by the mesh, saying again what this node last applied: declaration %s (%s)", + short(report.Declared), report.Order.Words())) + } q.tell(ctx, report, unasked) } diff --git a/internal/link/run.go b/internal/link/run.go index dfaeedb..0d615e4 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -234,6 +234,11 @@ func serve(ctx context.Context, link Link, m Membership, queue *Queue, say Annou defer queue.detach(link) declarations := link.Declarations() + // And the mesh asking what this node last applied (to-be 45 §6), on a link that hears the question. + var asks <-chan struct{} + if a, ok := link.(Asked); ok { + asks = a.AskedToReport() + } for { // Asked to stop — by standing aside, say — nothing further is taken in. @@ -247,6 +252,8 @@ func serve(ctx context.Context, link Link, m Membership, queue *Queue, say Annou publishAlive(ctx, link, m, say, timeout) case reason := <-link.Lost(): return reason + case <-asks: + queue.ReportAsked() case declaration, ok := <-declarations: if !ok { // The link's own reason, when it has managed to say one: "stopped delivering" on