From 1ebad3786cd3935dd9540fe1947e77b20bbf642b Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 28 Sep 2026 16:07:18 +0200 Subject: [PATCH] The mesh says what it applied, and the replay has an address it may use MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pipeline was observable from a merge to an artifact and went dark where it touched a machine: a node's report is control traffic only the control plane reads, so nothing said which version a machine runs, or that it refused to (novox/hq ADR 0134). The control plane now states both under the seat it holds — a role's events belong to the role and keep their address when the holder is replaced — and only when the report is news, because a machine reconciles every minute and a fact per report would be a fact per minute per machine. Whether a report is news is the store's answer: it holds the previous one, so the listener returns it and the server states the fact. That also gives the catch-up replay a subject the controller may publish: it was published as a module's event from a module called "control-plane", which does not exist, so the controller's own account refused it and every catalogue that asked what it missed was answered with nothing. --- cmd/mesh-controller/adopting_test.go | 2 +- internal/broker/nats.go | 8 ++ internal/broker/states_agreement_test.go | 30 ++++++ internal/broker/streams.go | 11 +++ internal/broker/testdata/composed.conf | 2 +- internal/inventory/doing_test.go | 24 ++--- internal/inventory/nodes.go | 47 +++++---- internal/link/bus.go | 33 +++++++ internal/link/enrolment.go | 29 +++--- internal/link/events.go | 33 +++++++ internal/link/heard_test.go | 12 +-- internal/link/receive_nats_test.go | 117 +++++++++++++++++++++-- internal/link/rekey_test.go | 8 +- internal/link/report_retry_test.go | 4 +- internal/link/serve.go | 70 +++++++++++++- internal/link/stale_report_test.go | 6 +- 16 files changed, 366 insertions(+), 70 deletions(-) create mode 100644 internal/broker/states_agreement_test.go diff --git a/cmd/mesh-controller/adopting_test.go b/cmd/mesh-controller/adopting_test.go index 3eddda5..d023ab6 100644 --- a/cmd/mesh-controller/adopting_test.go +++ b/cmd/mesh-controller/adopting_test.go @@ -88,7 +88,7 @@ func reportsReaching(t *testing.T, open *stores, reachable []link.Reach, held .. if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil { t.Fatal(err) } - if err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{ + if _, err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{ Node: "anchor", Applied: []string{"hello-web.x"}, Declared: digestOf(body), Firewall: "ufw", Held: held, Reachable: reachable, }); err != nil { diff --git a/internal/broker/nats.go b/internal/broker/nats.go index d69cd01..7d5a913 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -181,6 +181,14 @@ func PermissionsFor(p Principal) (Permissions, error) { for _, seat := range meshSeatsTheControllerUses { pub = append(pub, "mesh.seat."+seat+".accept.>") } + // **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own + // facts under the seat it holds, because a role's events belong to the role and keep their + // address while the holder is replaced. Named one by one rather than as a whole namespace: + // least authority, and a fact nothing states is authority nobody uses. + for _, event := range ControllerStates { + pub = append(pub, seatEventSubject(ControllerSeat, event)) + } + // Every module's tools: **the control plane is the way in** (novox/hq ADR 0095). A person // or an agent asks through it and every question passes one process where an audit // belongs — so it, alone among principals, may call any tool by name. The first `ask` on diff --git a/internal/broker/states_agreement_test.go b/internal/broker/states_agreement_test.go new file mode 100644 index 0000000..5d64bc3 --- /dev/null +++ b/internal/broker/states_agreement_test.go @@ -0,0 +1,30 @@ +package broker_test + +import ( + "slices" + "testing" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/link" +) + +// The facts the control plane states are named twice — in the grant that permits them and in the code +// that states them — because `link` imports `broker` and the dependency cannot go the other way. So a +// test keeps them agreeing: a subject the grant omits is refused at the moment the mesh has something +// to say, and one the grant adds that nothing states is authority nobody uses. +// +// An external test package, because it may import both while neither imports the other. +func TestTheFactsTheGrantPermitsAreTheFactsTheMeshStates(t *testing.T) { + if broker.ControllerSeat != link.MeshControllerSeat { + t.Fatalf("the grant is written for the %q seat and the mesh states its facts under %q", + broker.ControllerSeat, link.MeshControllerSeat) + } + for _, event := range []string{link.KeyApplied, link.KeyRefused, link.KeyBuiltBefore} { + if !slices.Contains(broker.ControllerStates, event) { + t.Errorf("the mesh states %q and its account may not publish it", event) + } + } + if len(broker.ControllerStates) != 3 { + t.Errorf("the grant permits %v, which is more than the mesh states", broker.ControllerStates) + } +} diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 9c6dcf7..6408db5 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -160,6 +160,17 @@ func Overlaps() []string { // ack subject is derived from (nats.go: `$JS.ACK..controller.>`). const ControllerName = "controller" +// ControllerSeat is the role the control plane holds, and ControllerStates are the facts it states +// under it (novox/hq ADR 0134). +// +// **Written here as well as in `link`, and a test keeps them agreeing.** `link` imports `broker`, so +// `broker` cannot import `link`; a grant naming a subject the controller never publishes is authority +// nobody uses, and a controller publishing one the grant omits is refused at the moment it has +// something to say. +const ControllerSeat = "mesh-controller" + +var ControllerStates = []string{"applied", "refused", "built-before"} + // ControllerFollows are the events the controller reacts to: the catalogue saying a module's // current version moved, and a catalogue that has just started saying it may have missed builds. // diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index e6275f8..03cf35a 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused"] } subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "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"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/inventory/doing_test.go b/internal/inventory/doing_test.go index 55a3552..8ce8f0e 100644 --- a/internal/inventory/doing_test.go +++ b/internal/inventory/doing_test.go @@ -30,12 +30,12 @@ func TestRefusedAndFailedAreDifferentSituations(t *testing.T) { refuser := nodeNamed(t, inv, "refuser") failer := nodeNamed(t, inv, "failer") - if err := inv.RecordDoing(ctx, refuser, Doing{ + if _, err := inv.RecordDoing(ctx, refuser, Doing{ Outcome: OutcomeRefused, Refused: "resource \"x\": a file needs a path", }); err != nil { t.Fatal(err) } - if err := inv.RecordDoing(ctx, failer, Doing{ + if _, err := inv.RecordDoing(ctx, failer, Doing{ Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}, Applied: 4, @@ -70,7 +70,7 @@ func TestAMachineDoingWhatItWasToldIsNotOnTheList(t *testing.T) { inv := fresh(t) ctx := context.Background() id := nodeNamed(t, inv, "fine") - if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 6}); err != nil { + if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 6}); err != nil { t.Fatal(err) } wrong, err := inv.NotDoingWhatTheyWereTold(ctx) @@ -97,12 +97,12 @@ func TestTheLastReportReplacesTheOneBefore(t *testing.T) { inv := fresh(t) ctx := context.Background() id := nodeNamed(t, inv, "recovered") - if err := inv.RecordDoing(ctx, id, Doing{ + if _, err := inv.RecordDoing(ctx, id, Doing{ Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "a", Error: "no"}}, }); err != nil { t.Fatal(err) } - if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 3}); err != nil { + if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 3}); err != nil { t.Fatal(err) } wrong, err := inv.NotDoingWhatTheyWereTold(ctx) @@ -141,7 +141,7 @@ func TestWhatANodeSaidGoesWhenTheNodeDoes(t *testing.T) { inv := fresh(t) ctx := context.Background() id := nodeNamed(t, inv, "leaving") - if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeFailed}); err != nil { + if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeFailed}); err != nil { t.Fatal(err) } if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil { @@ -257,7 +257,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) { id := nodeNamed(t, inv, "looping") same := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image"}}} - if err := inv.RecordDoing(ctx, id, same); err != nil { + if _, err := inv.RecordDoing(ctx, id, same); err != nil { t.Fatal(err) } first, _, err := inv.DoingOf(ctx, "looping") @@ -269,7 +269,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) { } for range StuckAfter - 1 { - if err := inv.RecordDoing(ctx, id, same); err != nil { + if _, err := inv.RecordDoing(ctx, id, same); err != nil { t.Fatal(err) } } @@ -287,7 +287,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) { // The same resource failing with different words — a duration, a counter — is still the same // failure: it is the resource that loops, not the sentence. reworded := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image (after 31s)"}}} - if err := inv.RecordDoing(ctx, id, reworded); err != nil { + if _, err := inv.RecordDoing(ctx, id, reworded); err != nil { t.Fatal(err) } still, _, err := inv.DoingOf(ctx, "looping") @@ -300,7 +300,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) { // A different failure is a new situation, not a longer one. other := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}} - if err := inv.RecordDoing(ctx, id, other); err != nil { + if _, err := inv.RecordDoing(ctx, id, other); err != nil { t.Fatal(err) } changed, _, err := inv.DoingOf(ctx, "looping") @@ -312,7 +312,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) { } // And a clean apply clears it: the machine is doing what it was told, since nothing. - if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 2}); err != nil { + if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 2}); err != nil { t.Fatal(err) } fine, _, err := inv.DoingOf(ctx, "looping") @@ -324,7 +324,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) { } // The list of what is wrong carries the count, so `status` can say it. - if err := inv.RecordDoing(ctx, id, same); err != nil { + if _, err := inv.RecordDoing(ctx, id, same); err != nil { t.Fatal(err) } wrong, err := inv.NotDoingWhatTheyWereTold(ctx) diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index 5dfb5e6..b0d52b8 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -670,31 +670,39 @@ func sameFailure(a, b Doing) bool { // a clean apply clears both (novox/hq 04-ISSUES/065). The previous row is read first and the // comparison made here, so "the same" is a rule this package states rather than a jsonb equality // that would restart the count on a changed word in an error. -func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error { +// **And whether this report was news**, which is what makes a fact about it worth stating (novox/hq +// ADR 0134). A machine reconciles continuously and reports each time; the same outcome about the same +// declaration is the same state said again, and a fact per report would be a fact per minute per +// machine that tells nobody anything. Read here because the previous row is read here anyway. +func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) (news bool, err error) { failed, err := json.Marshal(d.Failed) if err != nil { - return err + return false, err + } + var before Doing + var beforeFailed []byte + found := i.store.Pool().QueryRow(ctx, + `select outcome, refused, failed, failing_since, failures, coalesce(declared,'') + from node_report where node = $1`, + node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times, + &before.Declared) + switch { + case errors.Is(found, pgx.ErrNoRows): + news = true + case found != nil: + return false, found + default: + if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil { + return false, err + } + news = before.Outcome != d.Outcome || before.Declared != d.Declared || !sameFailure(before, d) } var since *time.Time times := 0 if d.Outcome != OutcomeApplied { - var before Doing - var beforeFailed []byte - err := i.store.Pool().QueryRow(ctx, - `select outcome, refused, failed, failing_since, failures from node_report where node = $1`, - node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times) - switch { - case errors.Is(err, pgx.ErrNoRows): - case err != nil: - return err - default: - if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil { - return err - } - } now := time.Now() since, times = &now, 1 - if err == nil && sameFailure(before, d) && before.Since != nil { + if found == nil && sameFailure(before, d) && before.Since != nil { since, times = before.Since, before.Times+1 } } @@ -707,7 +715,10 @@ func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error declared = excluded.declared, failing_since = excluded.failing_since, failures = excluded.failures`, node, d.Outcome, d.Refused, failed, d.Applied, d.Declared, since, times) - return err + if err != nil { + return false, err + } + return news, nil } // NotDoingWhatTheyWereTold is every machine whose last report was not a clean apply. diff --git a/internal/link/bus.go b/internal/link/bus.go index dbe4f46..5640c79 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -30,6 +30,11 @@ type Bus interface { // (design 29 §4, the *state* shape). PublishDeclaration(ctx context.Context, node string, body []byte) error + // PublishSeatEvent states a fact under a role's own name, for the holder of that role. A + // module's event is addressed to the module; a role's is addressed to the role, so it keeps + // meaning when the holder changes (novox/hq ADR 0121, ADR 0129). + PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error + // AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves // must say so **at once** rather than after the whole wait: the difference between "that // module is down" and "that tool is slow" is the first thing a person asking wants. @@ -81,6 +86,13 @@ func EventSubject(source, key string) string { return "mesh.mod." + source + ".event." + key } +// SeatEventSubject is where a role's own event lands. Derived from the role, never from its holder: +// a fact about the build machine or about the control plane keeps its address when the module holding +// that role is replaced (novox/hq ADR 0121, ADR 0129). +func SeatEventSubject(seat, event string) string { + return "mesh.seat." + seat + ".event." + event +} + // DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so // a node that was away gets exactly the current one and a replayed older one is refused by // sequence — the wire-level answer to novox/hq issue 107. @@ -112,6 +124,27 @@ func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, bo return nil } +// PublishSeatEvent states a role's own fact. Same envelope as a module's event and a different +// address: the source header is the role, because that is what the fact is about. +func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error { + id, err := eventID() + if err != nil { + return err + } + h := nats.Header{} + h.Set("x-event-id", id) + h.Set("x-source", seat) + _, err = b.JS.PublishMsg(&nats.Msg{ + Subject: SeatEventSubject(seat, event), + Header: h, + Data: body, + }, nats.MsgId(id), nats.Context(ctx)) + if err != nil { + return fmt.Errorf("stating %s of the %s seat: %w", event, seat, err) + } + return nil +} + func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error { _, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx)) if err != nil { diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index 615a60e..3193d71 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -265,7 +265,7 @@ func (e Enrolment) Outstanding(ctx context.Context, node string) (string, error) return e.Inventory.Outstanding(ctx, node) } -func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { +func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err error) { // A store that could not be asked right now is said as such, so the report is kept for // another attempt rather than acknowledged and lost (novox/hq issue 082). defer func() { @@ -274,11 +274,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { } }() if report.Node == "" { - return errors.New("a report named no node") + return false, errors.New("a report named no node") } node, err := e.Inventory.NodeByName(ctx, report.Node) if err != nil { - return err + return false, err } // What an adopted node holds, which firewall it found, and what is reachable on it (novox/hq @@ -298,7 +298,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { Port: r.Port, By: r.By, Published: r.Published, ContainerPort: r.ContainerPort}) } if err := e.Inventory.RecordAdoption(ctx, node.ID, held, report.Firewall, reachable); err != nil { - return err + return false, err } } // What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it. @@ -308,7 +308,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { Peers: report.Tunnel.Peers, State: report.Tunnel.State, Note: report.Tunnel.Note, Kept: report.Tunnel.Kept, }); err != nil { - return err + return false, err } } // A node taking a found tunnel's key after enrolment (novox/hq ADR 0105). Verified against the @@ -317,9 +317,9 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { // does not verify or is stale — a refusal, not "not now", so the node hears why. if report.Rekey != nil { if err := e.rekey(ctx, node, *report.Rekey); err != nil { - return err + return false, err } - return e.Inventory.Seen(ctx, node.ID) + return false, e.Inventory.Seen(ctx, node.ID) } // A bare word that a node is there is not an account of what the machine did or holds: it @@ -334,7 +334,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { if report.Superseded != "" { log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, report.Superseded) } - return e.Inventory.Seen(ctx, node.ID) + return false, e.Inventory.Seen(ctx, node.ID) } // What it did is kept whichever way it went. Until this, a refusal or a failure moved // last_seen and the reason went to a log line, so "which machine is not doing what it was @@ -361,17 +361,20 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { // on top of it (novox/hq ADR 0038). Kept even when the declaration was refused: what the // machine carries is true regardless of what it thought of the last thing it was sent. if err := e.Inventory.RecordCarried(ctx, report.Node, report.Carried); err != nil { - return err + return false, err } - if err := e.Inventory.RecordDoing(ctx, node.ID, doing); err != nil { - return err + // **Whether this is news** is the store's answer: it holds the previous report, and a machine + // that reconciles every minute says the same thing until something changes (novox/hq ADR 0134). + news, err = e.Inventory.RecordDoing(ctx, node.ID, doing) + if err != nil { + return false, err } // A refusal, a failure, or a bare word that the node is there — none of them is an account of // what the machine holds, so each moves last_seen and nothing else. Recording a partial list // as though it were the whole would tell a rebuilding node to remove what it still has. if report.Refused != "" || len(report.Failed) > 0 || report.Applied == nil { - return e.Inventory.Seen(ctx, node.ID) + return news, e.Inventory.Seen(ctx, node.ID) } - return e.Inventory.RecordOwned(ctx, node.ID, report.Applied) + return news, e.Inventory.RecordOwned(ctx, node.ID, report.Applied) } diff --git a/internal/link/events.go b/internal/link/events.go index 9a14c04..cfcfe18 100644 --- a/internal/link/events.go +++ b/internal/link/events.go @@ -49,6 +49,39 @@ func eventID() (string, error) { return hex.EncodeToString(raw), nil } +// MeshControllerSeat is the role the control plane holds, and therefore where its own facts live: a +// role's events belong to the role, not to whichever container is holding it today (novox/hq ADR 0121, +// ADR 0129). It is what makes them addressable while the control plane itself is being replaced. +const MeshControllerSeat = "mesh-controller" + +// The facts the mesh states about its own work (novox/hq ADR 0134). +const ( + // KeyApplied: a machine now runs what it was sent. + KeyApplied = "applied" + // KeyRefused: a machine did not take what it was sent, and why. + KeyRefused = "refused" + // KeyBuiltBefore: a build the mesh already held, for a catalogue that asked what it missed. Not + // `built` — that is the build machine's, said as it happens, and a replay is neither. + KeyBuiltBefore = "built-before" +) + +// Applied is what a machine now runs, as the mesh states it. +type Applied struct { + Node string `json:"node"` + Declared string `json:"declared,omitempty"` + // Resources is how many the machine applied, not which: the list is the machine's own account + // of itself and belongs in the records, not in a fact every listener has to read past. + Resources int `json:"resources"` +} + +// Refused is a machine that would not take what it was sent. +type Refused struct { + Node string `json:"node"` + Declared string `json:"declared,omitempty"` + Refused string `json:"refused,omitempty"` + Failed map[string]string `json:"failed,omitempty"` +} + // KeyModuleBuilt is what the builder announces when it has built something. The catalogue places // it in the module graph; nothing else need care. const KeyModuleBuilt = "module.builder.built" diff --git a/internal/link/heard_test.go b/internal/link/heard_test.go index 3a89b21..6174c4b 100644 --- a/internal/link/heard_test.go +++ b/internal/link/heard_test.go @@ -22,7 +22,7 @@ func heardFrom(t *testing.T, report link.Report) (*inventory.Inventory, inventor if _, err := inv.AddNode(ctx, report.Node); err != nil { t.Fatal(err) } - if err := (link.Enrolment{Inventory: inv}).Heard(ctx, report); err != nil { + if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, report); err != nil { t.Fatal(err) } doing, said, err := inv.DoingOf(ctx, report.Node) @@ -103,7 +103,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T) if err := inv.RecordSent(ctx, node.ID, digest); err != nil { t.Fatal(err) } - if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{ + if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{ Node: "anchor", Applied: []string{"a", "b"}, Declared: digest, Carried: []int{5432}, }); err != nil { t.Fatal(err) @@ -126,7 +126,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T) } // Now the node says only that it is there, as it does every minute. - if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil { + if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil { t.Fatal(err) } if !currentOf("anchor") { @@ -162,7 +162,7 @@ func TestAFailureDoesNotBecomeTheAccountOfWhatTheMachineHolds(t *testing.T) { if err := inv.RecordOwned(ctx, node.ID, []string{"one", "two", "three"}); err != nil { t.Fatal(err) } - if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{ + if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{ Node: "workstation", Applied: []string{"one"}, Failed: map[string]string{"two": "no"}, }); err != nil { t.Fatal(err) @@ -200,13 +200,13 @@ func TestWhatAnAdoptedNodeHoldsIsKeptAndAnAliveWordDoesNotWipeIt(t *testing.T) { } check("after the report") - if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil { + if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil { t.Fatal(err) } check("after an alive word") // A reconcile report carrying only adoption is recorded, though it applied nothing. - if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor", + if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor", Firewall: "ufw"}); err != nil { t.Fatal(err) } diff --git a/internal/link/receive_nats_test.go b/internal/link/receive_nats_test.go index ad799e1..6f2b779 100644 --- a/internal/link/receive_nats_test.go +++ b/internal/link/receive_nats_test.go @@ -86,14 +86,15 @@ type counted struct { heard []Report } -func (c *counted) Heard(_ context.Context, r Report) error { +func (c *counted) Heard(_ context.Context, r Report) (bool, error) { c.mu.Lock() defer c.mu.Unlock() if c.err != nil { - return c.err + return false, c.err } c.heard = append(c.heard, r) - return nil + // News, so what the mesh states about a report is exercised wherever a report is. + return true, nil } func (c *counted) refusing(err error) { @@ -207,11 +208,11 @@ type sentAndHeardSafely struct { heard []Report } -func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) error { +func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) (bool, error) { s.mu.Lock() defer s.mu.Unlock() s.heard = append(s.heard, r) - return nil + return true, nil } func (s *sentAndHeardSafely) Outstanding(context.Context, string) (string, error) { @@ -451,4 +452,108 @@ func TestNatsWorkSlowerThanTheWindowIsNotHandedOverAgain(t *testing.T) { // slowly is a listener that runs whatever it was given. type slowly struct{ work func() } -func (s slowly) Heard(context.Context, Report) error { s.work(); return nil } +func (s slowly) Heard(context.Context, Report) (bool, error) { s.work(); return true, nil } + +// **The mesh says what it applied** (novox/hq ADR 0134), under the seat the control plane holds — and +// says nothing when a report is the same state said again, which is what a machine reconciling every +// minute sends. +func TestNatsTheMeshSaysWhatAMachineApplied(t *testing.T) { + js := aBus(t) + heard := make(chan *nats.Msg, 4) + sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) { + heard <- m + }) + if err != nil { + t.Fatal(err) + } + defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection + + _, stop := servingOn(t, js, &counted{}) + defer stop() + + // A report that changed something: the store says it was news. + body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store", "broker"}}) + if err != nil { + t.Fatal(err) + } + if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil { + t.Fatal(err) + } + select { + case m := <-heard: + if m.Subject != SeatEventSubject(MeshControllerSeat, KeyApplied) { + t.Fatalf("the mesh stated %q", m.Subject) + } + var said Applied + if err := json.Unmarshal(m.Data, &said); err != nil { + t.Fatal(err) + } + if said.Node != "anchor" || said.Declared != "d1" || said.Resources != 2 { + t.Fatalf("it said %+v", said) + } + case <-time.After(10 * time.Second): + t.Fatal("the mesh said nothing about a machine that now runs something else") + } + + // A refusal is its own fact, with the reason in it rather than only in a log. + refusal, err := json.Marshal(Report{Node: "anchor", Declared: "d2", + Failed: map[string]string{"gitea.server": "no such image"}}) + if err != nil { + t.Fatal(err) + } + if _, err := js.Context().Publish(ReportSubject("anchor"), refusal); err != nil { + t.Fatal(err) + } + select { + case m := <-heard: + if m.Subject != SeatEventSubject(MeshControllerSeat, KeyRefused) { + t.Fatalf("a refusal was stated as %q", m.Subject) + } + var said Refused + if err := json.Unmarshal(m.Data, &said); err != nil { + t.Fatal(err) + } + if said.Failed["gitea.server"] == "" { + t.Fatalf("the refusal does not say which resource or why: %+v", said) + } + case <-time.After(10 * time.Second): + t.Fatal("the mesh said nothing about a machine that refused what it was sent") + } +} + +// And a report that is not news is not a fact. A machine reconciles every minute; a fact per report +// would be a fact per minute per machine, which is a stream nobody reads. +func TestNatsAReportThatIsNotNewsIsNotStated(t *testing.T) { + js := aBus(t) + heard := make(chan *nats.Msg, 4) + sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) { + heard <- m + }) + if err != nil { + t.Fatal(err) + } + defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection + + // A store that records the report and says it was nothing new — which is what the mesh's own + // store says about a machine repeating itself. + _, stop := servingOn(t, js, sameAgain{}) + defer stop() + + body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}}) + if err != nil { + t.Fatal(err) + } + if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil { + t.Fatal(err) + } + select { + case m := <-heard: + t.Fatalf("the mesh stated %q about a machine that changed nothing", m.Subject) + case <-time.After(3 * time.Second): + } +} + +// sameAgain records a report and says it was the same state said again. +type sameAgain struct{} + +func (sameAgain) Heard(context.Context, Report) (bool, error) { return false, nil } diff --git a/internal/link/rekey_test.go b/internal/link/rekey_test.go index 95d5cc7..fc7be80 100644 --- a/internal/link/rekey_test.go +++ b/internal/link/rekey_test.go @@ -60,7 +60,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) { rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()} rekey.Proof = ed25519.Sign(private, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel())) - if err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}); err != nil { + if _, err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}); err != nil { t.Fatal(err) } placed, err := e.Inventory.Overlays(ctx) @@ -77,7 +77,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) { _ = hub // Replayed, it is stale: the previous key it names is no longer the node's. - err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}) + _, err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}) if err == nil || !strings.Contains(err.Error(), "previous overlay key") { t.Fatalf("a replayed rekey was accepted: %v", err) } @@ -93,7 +93,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) { rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()} rekey.Proof = ed25519.Sign(stranger, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel())) - err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}) + _, err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}) if err == nil || !strings.Contains(err.Error(), "not signed by anchor's identity key") { t.Fatalf("a rekey signed by a stranger was accepted: %v", err) } @@ -111,7 +111,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) { other := theTunnel() other.Port = 51820 moved.Proof = ed25519.Sign(mustPrivate(t, e, "anchor"), link.RekeyProof("anchor", ownKey, tunnelKey, other)) - if err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: moved}); err == nil { + if _, err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: moved}); err == nil { t.Fatal("a proof over another tunnel was accepted") } } diff --git a/internal/link/report_retry_test.go b/internal/link/report_retry_test.go index 6b46a33..b996b1e 100644 --- a/internal/link/report_retry_test.go +++ b/internal/link/report_retry_test.go @@ -9,12 +9,12 @@ import ( type heardWith struct{ err error } -func (h heardWith) Heard(context.Context, Report) error { return h.err } +func (h heardWith) Heard(context.Context, Report) (bool, error) { return h.err == nil, h.err } // switchable answers with whatever it is set to — the store away, then back. type switchable struct{ err error } -func (h *switchable) Heard(context.Context, Report) error { return h.err } +func (h *switchable) Heard(context.Context, Report) (bool, error) { return h.err == nil, h.err } func aReport(node, declared string) Report { return Report{Node: node, Declared: declared, Applied: []string{"store"}} diff --git a/internal/link/serve.go b/internal/link/serve.go index 63f4e7c..ddc13fc 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -29,7 +29,12 @@ type Enroller interface { // Listener is what the controller does with a report. Separate from Enroller so the two can be // given independently, and so a server that only sends declarations needs neither. type Listener interface { - Heard(ctx context.Context, report Report) error + // Heard records what a node said, and says whether it was **news** — a machine that now runs + // something else, or refuses something it did not refuse before. A machine reconciles + // continuously and reports each time, so what is news is the store's answer rather than the + // bus's: only this side has the previous report to compare with. What the mesh states about it + // is the server's (novox/hq ADR 0134). + Heard(ctx context.Context, report Report) (news bool, err error) } // Recorder keeps what builders say. @@ -257,7 +262,7 @@ func (s *Server) heartbeat(m Control) { return } if s.listener != nil { - if err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil { + if _, err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil { s.log.Printf("could not record that %s is here: %v", alive.Node, err) } } @@ -291,7 +296,7 @@ func (s *Server) reported(ctx context.Context, m Control) { return } - err := s.listener.Heard(context.Background(), report) + news, err := s.listener.Heard(context.Background(), report) switch s.decide(ctx, m, what, declaredIn, outstanding, err) { case Hold: // Held, not settled, while the store cannot take it: the node reports an apply once, @@ -307,6 +312,13 @@ func (s *Server) reported(ctx context.Context, m Control) { // node whose recovery copy is silently older than it looks. s.log.Printf("could not record %s's report: %v", report.Node, err) } + // **And the mesh says what it did** (novox/hq ADR 0134). Only when the report was news: a + // machine reports every convergence, and a fact per report would be a fact per minute per + // machine saying nothing. Stated after it is recorded, so nothing is announced that the + // mesh does not hold. + if err == nil && news { + s.saysWhatItDid(ctx, report) + } } switch { @@ -460,7 +472,11 @@ func (s *Server) catchingUp(ctx context.Context, m Control) { sent := 0 for _, a := range announcements { a.Replay = true - if err := EmitEvent(ctx, s.bus, KeyModuleBuilt, "control-plane", "", a); err != nil { + // Under the control plane's own seat (novox/hq ADR 0134). It used to be published as a + // module's event from a module called "control-plane", which does not exist — so the + // controller's own account refused it, every catalogue that asked what it missed was + // answered with nothing, and its graph kept the gap (found 2026-09-28). + if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, KeyBuiltBefore, replayed(a)); err != nil { // Said and abandoned rather than retried: the catalogue asks again every time it // starts, and half a graph delivered twice is no better than half delivered once. s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v", @@ -551,3 +567,49 @@ func (s *Server) sourceMoved(ctx context.Context, m Control) { } _ = m.Took() } + +// saysWhatItDid states what a machine now runs, or what it would not take, as a fact on the bus +// (novox/hq ADR 0134). +// +// **The control plane speaks, as the holder of its seat.** A node's report is control traffic only +// this process may read, so the chain from a merge to a machine went dark exactly where it touched +// one: nothing said which version a machine runs, or that it refused to. The facts are second-hand +// on purpose — one emitter, one ordering — and a machine that cannot reach the bus produces none, so +// absence is not health. +// +// A failure to state a fact is logged and nothing else: the report is recorded, which is the part +// that must not be lost, and the next change says the same thing again. +func (s *Server) saysWhatItDid(ctx context.Context, report Report) { + if s.bus == nil { + return + } + event, body := KeyApplied, any(Applied{ + Node: report.Node, Declared: report.Declared, Resources: len(report.Applied), + }) + if report.Refused != "" || len(report.Failed) > 0 { + event, body = KeyRefused, Refused{ + Node: report.Node, Declared: report.Declared, + Refused: report.Refused, Failed: report.Failed, + } + } + raw, err := json.Marshal(body) + if err != nil { + s.log.Printf("could not say what %s did: %v", report.Node, err) + return + } + if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, event, raw); err != nil { + s.log.Printf("could not say that %s %s: %v", report.Node, event, err) + } +} + +// replayed is one announcement as the control plane states it. The same body the build machine's +// outcome carries, because what the catalogue does with it is the same. +func replayed(a Announcement) []byte { + raw, err := json.Marshal(a) + if err != nil { + // A body that cannot be marshalled is a programming error, not a bus failure, and an empty + // one is refused by the reader rather than silently taken as an announcement of nothing. + return nil + } + return raw +} diff --git a/internal/link/stale_report_test.go b/internal/link/stale_report_test.go index 2e0b949..dcf69fb 100644 --- a/internal/link/stale_report_test.go +++ b/internal/link/stale_report_test.go @@ -23,12 +23,12 @@ type sentAndHeard struct { err error } -func (s *sentAndHeard) Heard(_ context.Context, r Report) error { +func (s *sentAndHeard) Heard(_ context.Context, r Report) (bool, error) { if s.err != nil { - return s.err + return false, s.err } s.heard = append(s.heard, r) - return nil + return true, nil } func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil }