From f62ee0bc75dab6931e24a3e455e5ce8f9620bb7e Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 00:30:15 +0200 Subject: [PATCH] Say an apply's report even when the apply ends the link A host that delivered its own successor stood aside before the report of that apply was published, so it failed with 'context canceled' and the release plan waited for a report that never came (novox/hq issue 264). Reports are now published on a context the stand-aside does not cancel, and the last apply's report is kept until the broker takes it and said again on the next link, so a crash between apply and report is covered too. --- cmd/mesh-host/main.go | 46 ++++++++- cmd/mesh-host/unsaid_test.go | 64 ++++++++++++ internal/link/run.go | 81 +++++++++++++-- internal/link/unsaid_test.go | 191 +++++++++++++++++++++++++++++++++++ internal/store/unsaid.go | 76 ++++++++++++++ 5 files changed, 447 insertions(+), 11 deletions(-) create mode 100644 cmd/mesh-host/unsaid_test.go create mode 100644 internal/link/unsaid_test.go create mode 100644 internal/store/unsaid.go diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 5056c69..51c572d 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -1068,7 +1068,7 @@ func runLink(ctx context.Context, opts options) error { // Asked with the version this host is RUNNING, read from where it sits — not the link-time // stamp, which every delivered host carries as "development build". Asked with the stamp, // a delivered host never matched the newest delivered version, so it stood aside on every - // push for ever, and standing aside cancels the report, so the mesh never heard from it + // push for ever, and standing aside then cancelled the report, so the mesh never heard from it // again (novox/hq 04-ISSUES/163). switch next, waiting, err := upgrade.Successor(upgrade.VersionsDir(""), runningVersion()); { case err != nil: @@ -1117,7 +1117,7 @@ func runLink(ctx context.Context, opts options) error { Password: mine.Membership.Password, Transport: mine.Membership.Transport, Signer: mine.Membership.Signer, - }, applier, say, opts.timeout, rousedBySignal(ctx), outbox) + }, applier, say, opts.timeout, rousedBySignal(ctx), outbox, unsaidBeside(opts.state)) // **Cleanly**, or the launcher counts standing aside as a crash and rolls the new host back // before it has run once. The context this returns on was cancelled deliberately, so its error // is not a fault to report. @@ -1127,6 +1127,48 @@ func runLink(ctx context.Context, opts options) error { return held } +// unsaidFile keeps the last apply's report beside the node's state until the mesh has taken it +// (novox/hq issue 264), so a report lost when this host stood aside for its successor — or died +// between applying and publishing — is said by whichever host runs next, once it is linked. +type unsaidFile struct{ path string } + +func unsaidBeside(statePath string) unsaidFile { return unsaidFile{path: store.UnsaidPath(statePath)} } + +func (u unsaidFile) Keep(r link.Report) error { + // A report naming no declaration is one the mesh cannot match to what it sent, and keeping it + // would only re-say an older apply after a newer one was refused. + if r.Declared == "" { + return store.ClearUnsaid(u.path) + } + body, err := json.Marshal(r) + if err != nil { + return err + } + return store.SaveUnsaid(u.path, body) +} + +func (u unsaidFile) Said(declared string) error { + kept, ok, err := u.Pending() + if err != nil || !ok || kept.Declared != declared { + return err + } + return store.ClearUnsaid(u.path) +} + +func (u unsaidFile) Pending() (link.Report, bool, error) { + raw, err := store.ReadUnsaid(u.path) + if err != nil || raw == nil { + return link.Report{}, false, err + } + var r link.Report + if err := json.Unmarshal(raw, &r); err != nil { + // Unreadable is forgotten rather than kept for ever: it can never be said. + _ = store.ClearUnsaid(u.path) + return link.Report{}, false, fmt.Errorf("the kept report at %s is unreadable, and is dropped: %w", u.path, err) + } + return r, r.Declared != "", nil +} + // adoptionWatch remembers what the node last said about what it holds and its firewall, so a // reconcile speaks unasked only when that changed. type adoptionWatch struct { diff --git a/cmd/mesh-host/unsaid_test.go b/cmd/mesh-host/unsaid_test.go new file mode 100644 index 0000000..4c63317 --- /dev/null +++ b/cmd/mesh-host/unsaid_test.go @@ -0,0 +1,64 @@ +package main + +import ( + "os" + "path/filepath" + "reflect" + "testing" + + "github.com/novox/mesh-host/internal/link" + "github.com/novox/mesh-host/internal/store" +) + +// The last apply's report, kept beside the state until the mesh takes it (novox/hq issue 264). + +// What a host that died between applying and reporting kept is read back by the next one exactly — +// the digest and the outcome, never a new account of the machine. +func TestTheKeptReportSurvivesTheHost(t *testing.T) { + state := filepath.Join(t.TempDir(), "state.json") + made := link.Report{Declared: "sha256:d1", Applied: []string{"a"}, Failed: map[string]string{"b": "no room"}, + Host: "v2"} + if err := unsaidBeside(state).Keep(made); err != nil { + t.Fatal(err) + } + + kept, ok, err := unsaidBeside(state).Pending() + if err != nil || !ok { + t.Fatalf("nothing kept: %v", err) + } + if !reflect.DeepEqual(kept, made) { + t.Fatalf("kept %+v, made %+v", kept, made) + } + + // Taken about another declaration, it stays; taken about this one, it goes. + if err := unsaidBeside(state).Said("sha256:other"); err != nil { + t.Fatal(err) + } + if _, ok, _ := unsaidBeside(state).Pending(); !ok { + t.Fatal("a report about another declaration cleared this one") + } + if err := unsaidBeside(state).Said("sha256:d1"); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(store.UnsaidPath(state)); !os.IsNotExist(err) { + t.Fatalf("the report was taken and is still kept: %v", err) + } +} + +// A node that never applied anything has nothing to say again; and an apply refused before its +// declaration was read clears an older report, which would otherwise be re-said after it. +func TestNothingKeptWhenNothingWasApplied(t *testing.T) { + state := filepath.Join(t.TempDir(), "state.json") + if _, ok, err := unsaidBeside(state).Pending(); ok || err != nil { + t.Fatalf("a node that applied nothing has a report to say: %v %v", ok, err) + } + if err := unsaidBeside(state).Keep(link.Report{Declared: "sha256:d1"}); err != nil { + t.Fatal(err) + } + if err := unsaidBeside(state).Keep(link.Report{Refused: "forged"}); err != nil { + t.Fatal(err) + } + if _, ok, _ := unsaidBeside(state).Pending(); ok { + t.Fatal("an older report is kept after a later declaration was refused") + } +} diff --git a/internal/link/run.go b/internal/link/run.go index ae1f28f..8832183 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -72,7 +72,26 @@ type Announce func(string) type Roused <-chan struct{} func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error { - return HoldRoused(ctx, m, apply, say, timeout, nil, nil) + return HoldRoused(ctx, m, apply, say, timeout, nil, nil, nil) +} + +// Unsaid keeps the report of the last apply until the mesh has taken it, so a report lost between +// applying and publishing — the host standing aside for a successor, a crash, a power cut — is said +// again the next time this node is linked (novox/hq issue 264). Without it the declaration has been +// settled or applied, the machine is what the mesh said, and the mesh waits for a report that the +// node believes is gone. +// +// **Exactly the report the apply made, never a new one.** Re-applying the kept declaration would +// describe the machine as it is now; what the mesh waits for is what was done with what it sent. +// Nil is allowed and keeps nothing. +type Unsaid interface { + // Keep holds the report of an apply until it is said. A report naming no declaration — refused + // before it was read — replaces nothing worth saying, and clears what was kept. + Keep(Report) error + // Said forgets the kept report once the broker has taken the report about the same declaration. + Said(declared string) error + // Pending is the report kept and not yet said, if any. + Pending() (Report, bool, error) } // Outbox carries reports the node has to say without having been sent anything — what a @@ -91,10 +110,10 @@ type Unasked struct { // HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed // reports to publish between deliveries. func HoldRoused(ctx context.Context, m Membership, apply Applier, say Announce, - timeout time.Duration, roused Roused, outbox Outbox) error { + timeout time.Duration, roused Roused, outbox Outbox, unsaid Unsaid) error { return holdWith(ctx, func(ctx context.Context) error { - return Run(ctx, m, apply, say, timeout, outbox) + return Run(ctx, m, apply, say, timeout, outbox, unsaid) }, say, roused) } @@ -188,15 +207,29 @@ func holdWith(ctx context.Context, run attempt, say Announce, roused Roused) err // Outbound only, and nothing listens on this machine. Returns when the link ends, for any reason; // Hold is what decides whether to open it again. func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration, - outbox Outbox) error { - if say == nil { - say = func(string) {} - } + outbox Outbox, unsaid Unsaid) error { link, err := Open(ctx, m, timeout) if err != nil { return err } defer link.Close() + return serve(ctx, link, m, apply, say, timeout, outbox, unsaid) +} + +// serve is Run on a link already open: separated so what the node says, and when, can be tested +// without a broker. +func serve(ctx context.Context, link Link, m Membership, apply Applier, say Announce, + timeout time.Duration, outbox Outbox, unsaid Unsaid) error { + if say == nil { + say = func(string) {} + } + // **A report about something done is said even when the link is being let go.** An apply can + // end the link itself — the host stands aside for a successor it delivered (novox/hq ADR 0141) + // — and a report published on the cancelled context was refused before it left: "applied, and + // could not tell the mesh: reporting: context canceled", while the release plan waited for it + // (novox/hq issue 264). Not cancelled with the link, still bounded by the timeout, and sent + // before the deferred close lets go of the connection. + reporting := context.WithoutCancel(ctx) // Said, because it is the event anybody watching actually wants. Without it a node logs every // failure and nothing on success, so a log full of "trying again" and then silence reads as @@ -210,9 +243,30 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout defer beat.Stop() publishAlive(ctx, link, m, say, timeout) + // **What the last apply did, if the mesh never heard it** — said before anything newly + // delivered is applied, so it can never land after, and read as newer than, a later report. + if unsaid != nil { + switch kept, ok, err := unsaid.Pending(); { + case err != nil: + say("cannot read the report this node kept unsaid: " + err.Error()) + case ok: + say("saying again what the last apply did: its report never reached the mesh") + if publishReport(reporting, link, m, kept, say, timeout) { + if err := unsaid.Said(kept.Declared); err != nil { + say("said the kept report, and cannot forget it: " + err.Error()) + } + } + } + } + declarations := link.Declarations() for { + // Asked to stop — by standing aside, say — nothing further is applied, whichever of the + // cases below the select would otherwise have picked. + if ctx.Err() != nil { + return nil + } select { case <-ctx.Done(): return nil @@ -243,11 +297,16 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout declaration, superseded := newest(declarations, declaration, drainWindow) for _, old := range superseded { say("set aside a declaration: a newer one arrived with it") - publishReport(ctx, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()), + publishReport(reporting, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()), Superseded: declaredIn(declaration.Body())}, say, timeout) _ = old.Handled() } report := handleBody(ctx, m, declaration.Body(), apply) + if unsaid != nil { + if err := unsaid.Keep(report); err != nil { + say("cannot keep this apply's report until it is said: " + err.Error()) + } + } switch { case report.Refused != "": say("refused a declaration: " + report.Refused) @@ -258,7 +317,11 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout say(fmt.Sprintf("applied %d resource(s)%s", len(report.Applied), heldNote(report.Held))) } - publishReport(ctx, link, m, report, say, timeout) + if publishReport(reporting, link, m, report, say, timeout) && unsaid != nil { + if err := unsaid.Said(report.Declared); err != nil { + say("said this apply's report, and cannot forget it: " + err.Error()) + } + } // Settled after the report is published. A node that dies between applying and // reporting leaves the declaration with the mesh and applies it again on return, // which is safe because applying is reconciliation — it converges rather than diff --git a/internal/link/unsaid_test.go b/internal/link/unsaid_test.go new file mode 100644 index 0000000..d42da39 --- /dev/null +++ b/internal/link/unsaid_test.go @@ -0,0 +1,191 @@ +package link + +import ( + "context" + "crypto/ed25519" + "encoding/json" + "errors" + "reflect" + "sync" + "testing" + "time" +) + +// A report of an apply reaches the mesh even when the apply ended the link, and one lost anyway is +// said again the next time the node is linked (novox/hq issue 264). + +// quietLink is a link with no broker behind it. A report is refused once the context it is +// published on is cancelled, as the bus's is ("reporting: context canceled"), or while refusing. +type quietLink struct { + mu sync.Mutex + reports []Report + refusing bool + delivered chan Declaration + lost chan error +} + +func newQuietLink(arriving ...Declaration) *quietLink { + l := &quietLink{delivered: make(chan Declaration, len(arriving)), lost: make(chan error, 1)} + for _, d := range arriving { + l.delivered <- d + } + return l +} + +func (l *quietLink) Report(ctx context.Context, _ string, body []byte) error { + if err := ctx.Err(); err != nil { + return err + } + l.mu.Lock() + defer l.mu.Unlock() + if l.refusing { + return errors.New("refused") + } + var r Report + if err := json.Unmarshal(body, &r); err != nil { + return err + } + l.reports = append(l.reports, r) + return nil +} +func (l *quietLink) Alive(context.Context, string, []byte) error { return nil } +func (l *quietLink) Declarations() <-chan Declaration { return l.delivered } +func (l *quietLink) Lost() <-chan error { return l.lost } +func (l *quietLink) Close() {} + +func (l *quietLink) said() []Report { + l.mu.Lock() + defer l.mu.Unlock() + return append([]Report(nil), l.reports...) +} + +// keptInMemory is Unsaid without a disk. +type keptInMemory struct { + report *Report +} + +func (k *keptInMemory) Keep(r Report) error { + if r.Declared == "" { + k.report = nil + return nil + } + k.report = &r + return nil +} +func (k *keptInMemory) Said(declared string) error { + if k.report != nil && k.report.Declared == declared { + k.report = nil + } + return nil +} +func (k *keptInMemory) Pending() (Report, bool, error) { + if k.report == nil { + return Report{}, false, nil + } + return *k.report, true, nil +} + +func aMember(t *testing.T) (Membership, ed25519.PrivateKey) { + t.Helper() + public, private, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + return Membership{Node: "n1", Signer: public}, private +} + +// Measured on the anchor: a declaration carried a new controller and a new host, the host applied +// it, stood aside for its successor — which cancels the link — and the report of that apply failed +// with "context canceled". The plan waited on it until somebody pushed by hand. +func TestAnApplyThatEndsTheLinkIsStillReported(t *testing.T) { + m, key := aMember(t) + declaration := &said{body: signedBy(t, key, []byte(`{"declaration":1}`))} + l := newQuietLink(declaration) + kept := &keptInMemory{} + + ctx, standAside := context.WithCancel(context.Background()) + defer standAside() + apply := func(context.Context, []byte, []byte) Report { + standAside() // the host delivered its successor, and stands aside for it + return Report{Declared: "d1", Applied: []string{"a"}} + } + + done := make(chan error, 1) + go func() { done <- serve(ctx, l, m, apply, nil, time.Second, nil, kept) }() + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(5 * time.Second): + t.Fatal("the link did not end after the apply stood aside") + } + + reports := l.said() + if len(reports) != 1 || reports[0].Declared != "d1" || !reflect.DeepEqual(reports[0].Applied, []string{"a"}) { + t.Fatalf("the apply's report did not reach the mesh: %+v", reports) + } + if !declaration.handled { + t.Fatal("the declaration was not settled after its report") + } + if _, ok, _ := kept.Pending(); ok { + t.Fatal("a report the mesh took is still kept to be said again") + } +} + +// A report that did not reach the mesh — the host died between applying and publishing, or the +// broker refused it — is said on the next link, before anything else, and exactly as it was made. +func TestAReportLostAfterTheApplyIsSaidOnTheNextLink(t *testing.T) { + m, key := aMember(t) + made := Report{Declared: "d1", Applied: []string{"a"}, Failed: map[string]string{"b": "no room"}} + kept := &keptInMemory{} + + // The apply happens and its report is lost. + first := newQuietLink(&said{body: signedBy(t, key, []byte(`{"declaration":1}`))}) + first.refusing = true + ctx, stop := context.WithCancel(context.Background()) + apply := func(context.Context, []byte, []byte) Report { stop(); return made } + if err := serve(ctx, first, m, apply, nil, time.Second, nil, kept); err != nil { + t.Fatal(err) + } + if len(first.said()) != 0 { + t.Fatal("the broker refused the report and it counted as said") + } + + // The next host links. It is stopped at once, so all it does is what it does on linking. + next := newQuietLink() + gone, cancel := context.WithCancel(context.Background()) + cancel() + never := func(context.Context, []byte, []byte) Report { + t.Fatal("nothing was delivered, and something was applied") + return Report{} + } + if err := serve(gone, next, m, never, nil, time.Second, nil, kept); err != nil { + t.Fatal(err) + } + reports := next.said() + made.Node = m.Node + if len(reports) != 1 || !reflect.DeepEqual(reports[0], made) { + t.Fatalf("the lost report was not said again as it was made: %+v", reports) + } + if _, ok, _ := kept.Pending(); ok { + t.Fatal("a report said again is still kept") + } +} + +// A node that never applied anything, or whose last report was taken, says nothing on linking. +func TestNothingIsSaidAgainWhenNothingWasLost(t *testing.T) { + m, _ := aMember(t) + l := newQuietLink() + gone, cancel := context.WithCancel(context.Background()) + cancel() + if err := serve(gone, l, m, nil, nil, time.Second, nil, &keptInMemory{}); err != nil { + t.Fatal(err) + } + if err := serve(gone, l, m, nil, nil, time.Second, nil, nil); err != nil { + t.Fatal(err) + } + if reports := l.said(); len(reports) != 0 { + t.Fatalf("a node with nothing lost reported %+v", reports) + } +} diff --git a/internal/store/unsaid.go b/internal/store/unsaid.go new file mode 100644 index 0000000..4412b19 --- /dev/null +++ b/internal/store/unsaid.go @@ -0,0 +1,76 @@ +package store + +import ( + "errors" + "os" + "path/filepath" +) + +// The report of the last apply of a declaration from the mesh, kept until the mesh has taken it. +// +// **An apply the mesh never heard about is one it waits on for ever.** The controller's release +// plan waits for a machine to report the exact declaration it was sent; a report lost between the +// apply and the publish — a host standing aside for its successor, a crash, a power cut — leaves the +// machine applied and the plan waiting until somebody pushes by hand (novox/hq issue 264). The +// declaration itself is kept (declared.go), and re-applying it would say something new about the +// machine; what was lost is what *that* apply did, so that is what is kept, exactly as it was made, +// and said again once the host is linked. +// +// Opaque bytes here: the report is the link's word, and the store keeps it without reading it. + +// UnsaidName is where it lives, beside the state. +const UnsaidName = "unsaid.json" + +// UnsaidPath is where the unsaid report lives, given where the state lives. +func UnsaidPath(statePath string) string { + return filepath.Join(filepath.Dir(statePath), UnsaidName) +} + +// SaveUnsaid keeps a report the mesh has not taken yet, replacing whatever was kept before: only +// the last apply's report is worth saying again, because the mesh compares against what it sent last. +func SaveUnsaid(path string, report []byte) error { + if len(report) == 0 { + return errors.New("refusing to keep an empty report") + } + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + return err + } + tmp, err := os.CreateTemp(filepath.Dir(path), ".unsaid-*") + if err != nil { + return err + } + defer os.Remove(tmp.Name()) + if err := tmp.Chmod(0o600); err != nil { + tmp.Close() + return err + } + if _, err := tmp.Write(report); err != nil { + tmp.Close() + return err + } + if err := tmp.Sync(); err != nil { + tmp.Close() + return err + } + if err := tmp.Close(); err != nil { + return err + } + return os.Rename(tmp.Name(), path) +} + +// ReadUnsaid is the report kept unsaid, or nil when there is none — the ordinary case. +func ReadUnsaid(path string) ([]byte, error) { + raw, err := os.ReadFile(path) + if errors.Is(err, os.ErrNotExist) { + return nil, nil + } + return raw, err +} + +// ClearUnsaid forgets the kept report, once the mesh has taken it. Absent is already clear. +func ClearUnsaid(path string) error { + if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) { + return err + } + return nil +}