Files
mesh-host/internal/link/run.go
T
jochen 6e1bcdbde4 Say again what was last applied when the mesh asks (hq to-be 45 Phase 3, the report verb)
The controller's healer H1 answers a send that went unreported by asking the
machine first: a report lost on its way (issue 264) needs no second send. The
node-engine now hears mesh.node.<self>.ask.report on core NATS and enqueues a
reconcile whose account is said whether or not it is news; a delivery waiting
meanwhile is applied and reported instead. The answer is an ordinary report on
its own subject, so the node publishes nothing new and answers nobody's inbox.

The genesis lock grants the controller the healer-acted seat event, which it
now composes; the genesis test in mesh-controller holds the two equal.
2026-10-06 14:05:02 +02:00

427 lines
18 KiB
Go

package link
import (
"context"
"crypto/ed25519"
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
"time"
)
// ErrForged is what a node returns for a declaration whose signature is not the mesh's.
//
// Its own error, and it must never be confused with a malformed message. novox/hq ADR 0004
// requires a host to tell *this is not from the mesh I joined* apart from *this is malformed*:
// the first means somebody is trying, the second means something is broken.
var ErrForged = errors.New("this declaration was not signed by the mesh this node joined")
// AliveEvery is how often a node says it is there.
//
// Often enough that "no word for five minutes" means something, rarely enough that a hundred
// nodes are not a hundred messages a second. The mesh reads absence rather than presence, so what
// matters is the interval being known and steady.
const AliveEvery = 60 * time.Second
// Membership is what a node needs to reach its mesh again, held by the caller.
type Membership struct {
Node string
Broker string
Fingerprint string
Password string
Signer ed25519.PublicKey
// Transport is which bus this membership was minted for (hearing.go): the mesh's own, and a
// membership that names another is one this host cannot dial with.
Transport string
}
// Applier is what the host does with a declaration that has been proved to come from the mesh.
//
// It receives the signature as well as the declaration, so the host can keep both: what it was
// told is kept signed and verified again when it is read back, which means the file on disk is
// trusted for the same reason the message was rather than for being local.
type Applier func(ctx context.Context, declaration, signature []byte) Report
// Announce is how the link says what is happening, so a node running unattended leaves an
// account of it. Nil is allowed and means say nothing.
type Announce func(string)
// Hold keeps this node in the mesh, reconnecting for as long as it is asked to.
//
// Disconnection is an ordinary situation and not a failure (novox/hq ADR 0004), so this does not
// give up. A laptop shut for a week comes back and reconnects; it does not come back needing
// somebody to start it again.
//
// The backoff exists because the two common reasons differ in how long they last: a broker
// restarting is back in seconds, and a machine that has moved to a network with no route may be
// hours. Retrying every second for hours is a node shouting into nothing; waiting a minute after
// a broker blip is a node that is needlessly late. So it starts fast and slows down, and resets
// once a connection has actually held.
// Roused is a channel that says the machine has reason to believe its link is stale — it woke
// from suspend, or its network changed.
//
// **The machine knows before any timeout does.** A suspended laptop's connection is dead the
// moment it wakes, and heartbeats find that out in twenty or thirty seconds; for that time the
// node believes it is in the mesh and is not, which is the one state this design says must never
// be indistinguishable from being connected. Nothing new listens on the node to arrange it — the
// signal a service manager already sends is enough (novox/hq ADR 0004).
//
// Nil is allowed and means nothing ever rouses it, which is every machine that does not suspend.
type Roused <-chan struct{}
func Hold(ctx context.Context, m Membership, queue *Queue, say Announce, timeout time.Duration) error {
return HoldRoused(ctx, m, queue, say, timeout, 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)
}
// HoldRoused is Hold, told when the machine has reason to think its link is stale.
//
// **What arrives is enqueued, never applied here** (novox/hq to-be 45 §6). The queue's one worker
// applies, outlives every link, and says its reports on whichever link is open; the caller runs it
// (Queue.Run) for as long as this holds.
func HoldRoused(ctx context.Context, m Membership, queue *Queue, say Announce,
timeout time.Duration, roused Roused) error {
return holdWith(ctx, func(ctx context.Context) error {
return Run(ctx, m, queue, say, timeout)
}, say, roused)
}
// attempt is one try at holding the link open, returning when it ends for any reason.
//
// Named so the loop below can be driven without a broker. What the loop decides — when to wait,
// how long, what being roused does — is the part with the reasoning in it, and it was reachable
// only through a real connection before.
type attempt func(context.Context) error
func holdWith(ctx context.Context, run attempt, say Announce, roused Roused) error {
const (
first = 2 * time.Second
most = 2 * time.Minute
// A connection that lasted this long counts as having worked, so the next failure starts
// from the bottom again. Without it a node that reconnects and immediately drops climbs
// to the maximum and stays there, long after whatever caused it went away.
settled = 30 * time.Second
)
wait := first
for {
began := time.Now()
// The link runs under a context this loop can cancel, so being roused ends the current
// attempt rather than only shortening the wait after it.
//
// **That is the whole of it.** After a resume the socket looks perfectly healthy from
// inside this process — there is no error and no close, because nothing has tried to
// send anything. It is heartbeats that eventually discover it, twenty or thirty seconds
// later. A machine that knows it just woke does not have to wait to be told.
//
// A rouse that turns out to be spurious costs one reconnect, which is cheap and
// idempotent: the node redeclares its queue and anything unacknowledged is redelivered.
// The alternative costs half a minute of believing it is in a mesh it has left.
trying, done := context.WithCancel(ctx)
if roused != nil {
go func() {
select {
case <-trying.Done():
case <-roused:
say("woken or moved — dropping the link and opening it again")
done()
}
}()
}
err := run(trying)
done()
if ctx.Err() != nil {
return nil
}
if time.Since(began) > settled {
wait = first
}
switch {
case errors.Is(err, ErrWrongCertificate):
// Said in full every time rather than folded into a retry count. This does not mean
// the network is down; it means what answered is not the mesh this node joined, and
// no amount of waiting fixes it. The node keeps running what it was last told, which
// is the right thing to do while somebody works out what happened.
say("the broker is not the one this node joined: " + err.Error())
say("this will not fix itself. This node keeps running what it was last told.")
case err != nil:
say(fmt.Sprintf("disconnected: %v — trying again in %s", err, wait))
default:
say(fmt.Sprintf("the link closed — trying again in %s", wait))
}
select {
case <-ctx.Done():
return nil
case <-roused:
// And it does not serve out a wait computed for a broker that was restarting, either.
//
// The backoff is not *reset* by this. Being roused says the machine changed, not that
// whatever was refusing the connection has stopped — a laptop woken repeatedly on a
// network with no route would otherwise retry at full speed for as long as somebody
// keeps opening the lid.
say("woken or moved — trying again now")
case <-time.After(wait):
}
if wait *= 2; wait > most {
wait = most
}
}
}
// Run holds the link open once, applying what arrives and reporting what happened.
//
// 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, queue *Queue, say Announce, timeout time.Duration) error {
link, err := Open(ctx, m, timeout)
if err != nil {
return err
}
defer link.Close()
return serve(ctx, link, m, queue, say, timeout)
}
// 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, queue *Queue, say Announce,
timeout time.Duration) error {
if say == nil {
say = func(string) {}
}
// 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
// still broken when it means the opposite.
say("in the mesh, hearing what this node should be")
// A word every so often, so the mesh can tell a node that is quiet from one that is gone.
// Cheap on purpose: it carries a name and nothing else, because anything more would be a
// report, and reports are rare where this is constant.
beat := time.NewTicker(AliveEvery)
defer beat.Stop()
publishAlive(ctx, link, m, say, timeout)
// **What the last apply did, if the mesh never heard it**, is said before anything newly
// delivered can be: the queue says it as the link is handed to it (novox/hq issue 264). And the
// link is not let go while the queue's worker is mid-act: an apply that ends the link — the host
// standing aside for a successor it delivered (novox/hq ADR 0141) — is still reported on it,
// before the deferred close.
queue.attach(ctx, link)
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.
if ctx.Err() != nil {
return nil
}
select {
case <-ctx.Done():
return nil
case <-beat.C:
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
// its own says nothing about why, and why is the whole of what an operator wants.
select {
case reason := <-link.Lost():
return reason
default:
return errors.New("the mesh stopped sending this node declarations")
}
}
// Whatever else is already waiting goes with it. **Enqueued, not applied**: the queue
// applies the newest of everything delivered and not yet taken, and reports the rest as
// set aside (novox/hq to-be 45 §6). The link goes on beating meanwhile.
queue.Deliver(gather(declarations, declaration, drainWindow)...)
}
}
}
// short is a declaration's digest as a person reads it in a log line.
func short(digest string) string {
if len(digest) > 12 {
return digest[:12]
}
return digest
}
// drainDepth is how many declarations the link holds unread; drainWindow is how long it waits for
// another to follow the one it has. Both small: a push is rare and a backlog is the exception this
// exists for, not the shape of ordinary traffic.
const (
drainDepth = 16
drainWindow = 750 * time.Millisecond
)
// gather takes what is already waiting behind `first`, in the order it arrived. It waits `window`
// for a straggler after each arrival and no longer: a declaration in flight from the mesh arrives
// within that; one that does not is the next push. Which of them is applied is the queue's to decide
// (pick), across everything delivered and not yet taken.
//
// **Its job narrows once declarations are state rather than messages, and does not disappear.**
// On the bus a declaration is last-per-subject (novox/hq design 29 §4), so a node that was away
// receives exactly the current one instead of a queue of superseded ones — the catch-up half is the
// stream's. And a declaration's order decides which is newest, where this window only infers it from
// arrival time, which is the wire-level answer to novox/hq issue 107.
//
// What remains is the live case: three pushes in quick succession to a *connected*, idle node are
// three deliveries, whatever the stream later retains, and gathering them is what makes them one
// apply rather than two. So this is narrowed, not deleted — and saying which half goes is worth more
// than a note that it "can probably be removed", which is how a load-bearing window gets deleted by
// somebody in a hurry.
func gather(arriving <-chan Declaration, first Declaration, window time.Duration) []Declaration {
batch := []Declaration{first}
for {
select {
case next, ok := <-arriving:
if !ok {
return batch
}
batch = append(batch, next)
case <-time.After(window):
return batch
}
}
}
// handleBody is the whole of deciding whether to trust a message, separated from the broker so it
// can be tested as the security check it is rather than as message plumbing.
func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) Report {
var signed Signed
if err := json.Unmarshal(body, &signed); err != nil {
return Report{Node: m.Node, Refused: "this message is not a declaration: " + err.Error()}
}
// Before anything is read out of it, let alone applied. The host applies whatever the link
// delivers, so this check is the difference between the mesh changing this machine and
// anybody changing it.
if !ed25519.Verify(m.Signer, signed.Declaration, signed.Signature) {
return Report{Node: m.Node, Refused: ErrForged.Error()}
}
return apply(ctx, signed.Declaration, signed.Signature)
}
// publishReport tells the mesh what this node did, and says whether the broker took it.
// Publish sends one report on this node's own connection and returns: the one-shot path for a
// report a command makes rather than the running host — a rekey (novox/hq ADR 0105). The same
// account, the same pinned certificate and the same exchange as the running host's reports.
func Publish(ctx context.Context, m Membership, report Report, timeout time.Duration) error {
link, err := Open(ctx, m, timeout)
if err != nil {
return err
}
defer link.Close()
var said string
if !publishReport(ctx, link, m, report, func(s string) { said = s }, timeout) {
return errors.New(said)
}
return nil
}
func publishReport(ctx context.Context, bus Bus, m Membership, report Report,
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 false
}
publish, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
// Said rather than swallowed. A report that fails to publish leaves the mesh believing this
// node never answered, while the node believes it did — and the two would go on disagreeing
// with nothing anywhere saying so. That shape of fault is the one this project keeps finding.
if err := bus.Report(publish, m.Node, 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.
func publishAlive(ctx context.Context, bus Bus, m Membership, say Announce,
timeout time.Duration) {
body, err := json.Marshal(Alive{Node: m.Node, IntervalSeconds: int(AliveEvery / time.Second)})
if err != nil {
return
}
publish, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
// Not mandatory, unlike a report. Losing one is nothing: the next is a minute away, and the
// mesh is reading a gap rather than counting arrivals. Insisting on delivery would turn a
// harmless miss into a logged failure every minute.
if err := bus.Alive(publish, m.Node, body); err != nil {
say("could not tell the mesh this node is here: " + err.Error())
}
}
// heldNote is what this apply did NOT do, for the line that says what it did.
//
// **A count that does not add up is the only symptom a held resource had** (novox/hq 04-ISSUES/125).
// An adopted node keeps what it found until its module is taken (ADR 0100), and that is correct — but
// it was recorded only in the node's own state file. On the edge cut-over the mesh sent 346 resources,
// the journal said it applied 330, and nothing anywhere said which sixteen or why. Reading it took
// opening state.json by hand; not reading it took every public name on the machine down, because the
// operator had four green surfaces and a discrepancy nobody could interpret.
//
// So the line that reports the apply carries it. Grouped by module and ordered by name, because the
// sentence an operator needs is "route-proxy is assigned and not taken", and the module is the thing
// they can act on — `take` is the verb, and it takes a module.
func heldNote(held []Held) string {
if len(held) == 0 {
return ""
}
byModule := map[string]int{}
for _, h := range held {
byModule[h.Module]++
}
names := make([]string, 0, len(byModule))
for name := range byModule {
names = append(names, name)
}
sort.Strings(names)
parts := make([]string, 0, len(names))
for _, name := range names {
parts = append(parts, fmt.Sprintf("%s: %d", name, byModule[name]))
}
return fmt.Sprintf(", %d held until their module is taken (%s)",
len(held), strings.Join(parts, ", "))
}