diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 62c0477..681eb3c 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -638,11 +638,11 @@ func runLink(ctx context.Context, opts options) error { // (novox/hq ADR 0004). // Reports a reconcile has to make unasked — what an adopted node holds changed, or its // firewall did — go out over the link when it is up (novox/hq ADR 0100). - outbox := make(chan link.Report, 1) + outbox := make(chan link.Unasked, 1) watch := &adoptionWatch{} applier = watch.noting(applier) go holdTheMachine(ctx, opts, mine, say, sched, func(r link.Report) { - if !watch.changed(r) { + if !watch.differs(r) { return } select { @@ -650,7 +650,13 @@ func runLink(ctx context.Context, opts options) error { // An older one nobody has published yet; this one says everything it did. default: } - outbox <- r + // Counted as said only once the broker has taken it: queued and lost — the link down, the + // publish refused — the change would never be said again (novox/hq ADR 0100). + outbox <- link.Unasked{Report: r, Done: func(published bool) { + if published { + watch.said(r) + } + }} }) return link.HoldRoused(ctx, link.Membership{ @@ -686,15 +692,27 @@ func adoptionFingerprint(r link.Report) string { return strings.Join(parts, "\n") } -// changed records a report and says whether it differs from the last one that went out. -func (w *adoptionWatch) changed(r link.Report) bool { +// differs says whether a report says anything the last one that went out did not. It records +// nothing: what was said is what reached the mesh, not what was written down to send. +func (w *adoptionWatch) differs(r link.Report) bool { w.mu.Lock() defer w.mu.Unlock() - now := adoptionFingerprint(r) - if now == w.last { + return adoptionFingerprint(r) != w.last +} + +// said records a report the mesh has actually been told. +func (w *adoptionWatch) said(r link.Report) { + w.mu.Lock() + defer w.mu.Unlock() + w.last = adoptionFingerprint(r) +} + +// changed is differs and said together, for a report published as it is made. +func (w *adoptionWatch) changed(r link.Report) bool { + if !w.differs(r) { return false } - w.last = now + w.said(r) return true } diff --git a/cmd/mesh-host/main_test.go b/cmd/mesh-host/main_test.go index 4747647..4caab47 100644 --- a/cmd/mesh-host/main_test.go +++ b/cmd/mesh-host/main_test.go @@ -244,3 +244,21 @@ func TestOnlyOneApplyRunsAtATime(t *testing.T) { t.Errorf("the apply recorded %d resource(s): %v", len(known.Resources), loadErr) } } + +// Defends novox/hq ADR 0100: a change is counted as said only once the mesh has been told. Queued +// and lost — the link down when the reconcile spoke — it must be said again. +func TestAChangeThatNeverReachedTheMeshIsSaidAgain(t *testing.T) { + w := &adoptionWatch{} + held := link.Report{Firewall: "ufw", Held: []link.Held{{ID: "hello-web.page", Changed: "rewritten"}}} + if !w.differs(held) { + t.Fatal("the first report of a change was not new") + } + // The link was down: nothing published it, so nothing says it was said. + if !w.differs(held) { + t.Error("a change that never reached the mesh was counted as said") + } + w.said(held) + if w.differs(held) { + t.Error("a change the mesh was told was said again") + } +} diff --git a/internal/link/run.go b/internal/link/run.go index 5668af4..1379369 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -76,7 +76,15 @@ func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeou // Outbox carries reports the node has to say without having been sent anything — what a // reconcile found changed on an adopted node (novox/hq ADR 0100). Published while the link is up; // a report made while it is down waits in the channel for the next one. Nil is allowed. -type Outbox <-chan Report +type Outbox <-chan Unasked + +// Unasked is one such report, with the way to say whether it reached the mesh. Done is called +// with true only when the broker took it — a node that marked a change said because it queued it +// would never say it again, and the mesh would go on believing nothing changed. +type Unasked struct { + Report Report + Done func(published bool) +} // HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed // reports to publish between deliveries. @@ -261,10 +269,13 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout return nil case <-beat.C: publishAlive(ctx, channel, m, say, timeout) - case report := <-outbox: + case unasked := <-outbox: // Said without having been asked: a reconcile found what an adopted node holds, or // its firewall, changed since it last said. - publishReport(ctx, channel, m, report, say, timeout) + published := publishReport(ctx, channel, m, unasked.Report, say, timeout) + if unasked.Done != nil { + unasked.Done(published) + } case reason := <-closed: return fmt.Errorf("the link closed: %v", reason) case delivery, ok := <-deliveries: @@ -365,13 +376,14 @@ func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) R return apply(ctx, signed.Declaration, signed.Signature) } +// publishReport tells the mesh what this node did, and says whether the broker took it. func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report, - say Announce, timeout time.Duration) { + say Announce, timeout time.Duration) bool { report.Node = m.Node body, err := json.Marshal(report) if err != nil { say("cannot encode this node's own report: " + err.Error()) - return + return false } publish, cancel := context.WithTimeout(ctx, timeout) defer cancel() @@ -382,7 +394,9 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep if err := channel.PublishWithContext(publish, Exchange, KeyReport, true, false, amqp.Publishing{ContentType: "application/json", Body: body}); err != nil { say(fmt.Sprintf("applied, and could not tell the mesh: %v", err)) + return false } + return true } // publishAlive says this node is here, and nothing else.