The bus on NATS: both transports behind seams, and the rollout switch #87

Merged
jschoubben merged 40 commits from feat/nats-genesis into main 2026-09-27 17:36:41 +00:00
2 changed files with 190 additions and 0 deletions
Showing only changes of commit d0a9abcb1c - Show all commits
+106
View File
@@ -0,0 +1,106 @@
package link
import (
"errors"
"time"
"github.com/novox/mesh-controller/internal/inventory"
)
// The store window, as a decision rather than a mechanism.
//
// The guarantee (novox/hq ADR 0083): a push the controller cannot record because its store is
// restarting is **held and retried**, not dropped and not falsely acknowledged. On the bus the
// mesh runs on today that is done by keeping the delivery unacknowledged in memory and settling
// it later. On the bus being built it is a `nak` with a delay: the server holds it and redelivers,
// so the controller keeps no list of parked messages and a controller that restarts mid-window
// loses nothing it was holding.
//
// **The decision is the same either way, and the mechanism is not the interesting part.** What is
// interesting is that moving the holding into the server introduces a problem the in-memory
// version did not have, and the answer was already in the message.
// Verdict is what to do with one control message.
type Verdict int
const (
// Take it: apply, then acknowledge.
Take Verdict = iota
// Hold it: the store cannot record this yet. Nak with a delay and let the server redeliver.
Hold
// Stale: this is about a declaration the node has already moved past, and applying it would
// undo what came after. Acknowledge without acting — redelivering forever is worse.
Stale
// GiveUp: the store has not come back in time. Settle it and say so, loudly.
GiveUp
)
// StoreWindow decides. Pure, so the guarantee is testable without a bus, a store or a clock.
type StoreWindow struct {
// GiveUpAfter is how long one message may be held before it is let go with a line saying so.
GiveUpAfter time.Duration
}
// Decide answers for one delivery.
//
// - err is what the store said, or nil.
// - declaredIn is the digest of the declaration this message is about, empty when it is not
// about one (an enrolment, a build result).
// - outstanding is the digest the mesh last sent that node, empty when it has sent none.
// - heldFor is how long this message has already been held; zero on first delivery.
func (w StoreWindow) Decide(err error, declaredIn, outstanding string, heldFor time.Duration) Verdict {
// **Staleness is checked before the store, not after.** A redelivery that lost its race is
// not worth waiting on a store for, and asking the store first would mean a message about a
// superseded declaration holding a slot in the window that a current one needs.
if declaredIn != "" && outstanding != "" && declaredIn != outstanding {
return Stale
}
if err == nil {
return Take
}
if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) {
// Not the store being away: a refusal is an answer, and holding it would turn a message
// the mesh understood into one it retries forever.
return Take
}
if heldFor >= w.GiveUpAfter {
return GiveUp
}
return Hold
}
// RedeliverAfter is how long the server should hold a naked message before trying again.
//
// Backed off, and bounded. A store restarting is back in seconds; a store that is gone is not
// helped by being asked every second, and the delay is what keeps a window of held messages from
// becoming a spin.
func RedeliverAfter(heldFor time.Duration) time.Duration {
switch {
case heldFor < 5*time.Second:
return time.Second
case heldFor < 30*time.Second:
return 5 * time.Second
default:
return 15 * time.Second
}
}
// The problem holding-in-the-server introduces, and why the answer was already in the message.
//
// Holding a delivery in memory let the controller do something a server cannot: when a newer
// report for the same node arrived, it dropped the older one, "because acting on it after the
// newer would undo the newer". A `nak`ed message is the server's, and the server will redeliver
// it whatever else has happened in the meantime — so the older report comes back *after* the
// newer was applied, and applying it would undo exactly what that comment describes.
//
// **A report already says which declaration it is about.** `Declared` is the digest of the exact
// bytes the mesh sent, and it exists because an earlier attempt to order reports by time lost the
// race it invited — an apply that started under the previous declaration finishes after the next
// is sent, and the report reads as newer than the send. Clocks cannot answer *which*; the digest
// is the answer itself.
//
// So supersession stops being a thing the controller remembers and becomes a thing it checks: a
// report whose digest is not the one outstanding for that node is stale, and is acknowledged
// without being acted on. Which is the same shape as a node refusing a superseded declaration by
// sequence (novox/hq issue 107) — ordering settled by what the message says, not by when it
// happened to arrive.
+84
View File
@@ -0,0 +1,84 @@
package link
import (
"errors"
"testing"
"time"
)
func window() StoreWindow { return StoreWindow{GiveUpAfter: 2 * time.Minute} }
// The guarantee itself (novox/hq ADR 0083): a push the store cannot record is held, not dropped
// and not falsely acknowledged.
func TestAMessageTheStoreCannotTakeYetIsHeld(t *testing.T) {
if got := window().Decide(ErrTryAgain, "", "", 0); got != Hold {
t.Fatalf("got %v; a push the store could not record was not held", got)
}
}
// A refusal is an answer. Holding it would turn a message the mesh understood into one it
// retries forever.
func TestARefusalIsNotHeld(t *testing.T) {
if got := window().Decide(errors.New("that node does not exist"), "", "", 0); got != Take {
t.Fatalf("got %v; a refusal was mistaken for the store being away", got)
}
}
// Held has a limit, and past it the message is settled rather than held for ever.
func TestAStoreThatNeverComesBackEndsTheHold(t *testing.T) {
if got := window().Decide(ErrTryAgain, "", "", 3*time.Minute); got != GiveUp {
t.Fatalf("got %v; the window has no end", got)
}
if got := window().Decide(ErrTryAgain, "", "", time.Minute); got != Hold {
t.Fatalf("got %v; the window ended early", got)
}
}
// **The problem that holding in the server introduces.** A naked message is redelivered whatever
// else happened meanwhile, so a report about a superseded declaration comes back after the newer
// one was applied — and applying it would undo the newer.
func TestAReportAboutASupersededDeclarationIsNotApplied(t *testing.T) {
got := window().Decide(nil, "digest-of-the-old-one", "digest-of-the-current-one", 0)
if got != Stale {
t.Fatalf("got %v; a redelivery that lost its race would have undone what came after", got)
}
}
func TestAReportAboutTheOutstandingDeclarationIsApplied(t *testing.T) {
if got := window().Decide(nil, "same", "same", 0); got != Take {
t.Fatalf("got %v; a current report was discarded", got)
}
}
// A message that is about no declaration — an enrolment, a build result — is never stale: there
// is nothing for it to be out of date with.
func TestAMessageAboutNoDeclarationIsNeverStale(t *testing.T) {
if got := window().Decide(nil, "", "whatever-is-outstanding", 0); got != Take {
t.Fatalf("got %v; an enrolment was treated as a stale report", got)
}
if got := window().Decide(nil, "a-digest", "", 0); got != Take {
t.Fatalf("got %v; a report was called stale against a node that was sent nothing", got)
}
}
// Staleness is decided before the store is waited on: a redelivery that lost its race must not
// hold a slot in the window that a current message needs.
func TestAStaleMessageIsNotHeldForTheStore(t *testing.T) {
if got := window().Decide(ErrTryAgain, "old", "current", 0); got != Stale {
t.Fatalf("got %v; a superseded message was held for a store it would never be applied to", got)
}
}
// The delay backs off: a store that is gone is not helped by being asked every second, and the
// delay is what keeps a window of held messages from becoming a spin.
func TestRedeliveryBacksOff(t *testing.T) {
first := RedeliverAfter(0)
later := RedeliverAfter(10 * time.Second)
last := RedeliverAfter(time.Minute)
if !(first < later && later < last) {
t.Fatalf("delays do not back off: %v %v %v", first, later, last)
}
if last > 30*time.Second {
t.Fatalf("a held message waits %v between attempts, which is longer than a store restart", last)
}
}