Files
mesh-controller/internal/link/serve.go
T
jschoubben fe5988c536 The filter constrains what arrives from outside, and names no network
The forward chain blocked everything passing through the machine and then allowed
the machine's own containers back by naming their address ranges: 172.16.0.0/12 and
192.168.128.0/17 fixed here, the rest recorded per machine by 0043. Every way of
keeping that list correct fails — a constant describes one machine, a recorded range
goes stale in silence and cannot tell a network the mesh made from one a predecessor
left behind, and generating it from the modules would put half the rule set on the
machine.

The mesh has no position on a container reaching outward: that is not a port opened
to anybody. So both chains are written around the links traffic arrives on. What did
not arrive from outside is accepted in one line; what did meets the declared rules.
The tunnel is named beside the outward links rather than treated as inside, or a port
nothing declares would be reachable from every machine in the mesh.

A machine that has not reported an outward link is sent no filter and keeps the one
it has, refused where a person reads it rather than as a rule set that will not load.

Removes the two constants, `node networks`, and the column behind it. novox/hq ADR
0140, superseding 0137 and 0139.
2026-09-29 01:01:29 +02:00

613 lines
24 KiB
Go

package link
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"github.com/novox/mesh-controller/internal/broker"
"log"
"os"
"time"
)
// AMQPVar is the controller's own connection to the bus the mesh runs on today.
const AMQPVar = "MESH_BROKER_AMQP"
// Enroller is what the controller does with an enrolment request.
//
// An interface so the serving loop can be tested against a real bus without a database, and so the
// two concerns — moving messages, and deciding — stay apart.
type Enroller interface {
// Enrol spends the token, records the key, and reports the node's name. The error is
// returned to the node as a refusal; it must be the same for every reason a token can fail.
Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, error)
}
// 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 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.
//
// Separate from Listener because they are different things arriving: a report is a node saying
// what it did with a declaration, and a build result is a machine saying what came of some work.
// One interface carrying both would mean an implementation of one having to say something about
// the other.
type Recorder interface {
Built(ctx context.Context, result BuildResult) error
}
// Upgrader is what the controller does when the catalogue says a module moved.
//
// An interface for the same reason Enroller is one: deciding what an upgrade means for the
// machines running it is a different concern from noticing that one was announced, and only the
// first needs a database.
type Upgrader interface {
// Upgraded is told which module moved and between which commits. An error is logged and the
// message is not handed back: an upgrade the controller could not act on is not one it will
// act on by being given the same message again, and a poison message on a durable queue would
// stop every upgrade behind it — except the store unreachable for the moment, which is asked
// again for a bounded time (novox/hq issue 083).
Upgraded(ctx context.Context, u Upgraded) error
// SourceMoved is a merge on the forge: build what that source produces, bases first.
SourceMoved(ctx context.Context, m SourceMoved) error
}
// Server acts on what nodes and modules say.
//
// **It holds no transport.** What arrives comes through Inbound and what it publishes goes through
// Bus, so this file is the controller's *decisions* about messages and nothing about wires. The
// connection below is the bus the mesh runs on today, kept because the command line publishes over
// the same one until the rollout (ADR 0116).
type Server struct {
inbound Inbound
bus Bus
js *broker.JetStream
enroller Enroller
listener Listener
recorder Recorder
upgrader Upgrader
replayer Replayer
log *log.Logger
// giveUp is how long one message is held for the store; zero means GiveUpAfter.
giveUp time.Duration
}
// ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and
// asking again, as when the store is restarting (novox/hq issue 082).
var ErrTryAgain = errors.New("not now, try again")
// GiveUpAfter bounds how long one message is kept trying. A store that has not come back in this
// long is not restarting, and the message is let go with a line saying it was lost.
const GiveUpAfter = 2 * time.Minute
// Records tells the server where to keep build results.
//
// Set after Connect rather than passed to it, because a controller that only publishes — the
// `build` command, which waits for its own answer — needs a connection and no recorder, and making
// it supply one would have it construct something it never uses.
func (s *Server) Records(r Recorder) { s.recorder = r }
// Follows says what to do about upgrades, and asks for them to be delivered.
func (s *Server) Follows(u Upgrader) error {
if err := s.inbound.Also(KindSourceMoved); err != nil {
return err
}
if err := s.inbound.Also(KindModuleMoved); err != nil {
return err
}
s.upgrader = u
return nil
}
// Answers says what to do about a catalogue's catch-up request, and asks for them to be delivered.
func (s *Server) Answers(r Replayer) error {
if err := s.inbound.Also(KindCatchUp); err != nil {
return err
}
s.replayer = r
return nil
}
// Connect opens the controller's own connection to the bus.
//
// On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq
// 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may
// have moved since.
// Connect is ConnectNats: the mesh has one bus (novox/hq ADR 0131, design 28 task 5.5). Kept as
// the name callers know; the transport it opened before is gone with the bus it spoke to.
func Connect(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
return ConnectNats(js, enroller, listener)
}
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
func (s *Server) Close() {
if s.inbound != nil {
s.inbound.Close()
}
if s.js != nil {
s.js.Close()
}
}
// Serve acts on what arrives until the context ends.
func (s *Server) Serve(ctx context.Context) error {
s.log.Printf("consuming what nodes say: %s, %s, %s, %s",
KindEnrolment, KindReport, KindHeartbeat, KindBuilt)
if s.upgrader != nil {
s.log.Printf("following %s", KindModuleMoved)
}
if s.replayer != nil {
s.log.Printf("answering %s", KindCatchUp)
}
return s.inbound.Receive(ctx, s.act)
}
// act is one message, whichever bus it came over.
func (s *Server) act(ctx context.Context, m Control) {
switch m.Kind() {
case KindEnrolment:
s.enrolling(ctx, m)
case KindReport:
s.reported(ctx, m)
case KindHeartbeat:
s.heartbeat(m)
case KindBuilt:
s.wasBuilt(ctx, m)
case KindModuleMoved:
s.moved(ctx, m)
case KindSourceMoved:
s.sourceMoved(ctx, m)
case KindCatchUp:
s.catchingUp(ctx, m)
default:
// Dropped: a message nothing understands will not be understood on the next attempt
// either, and asking for it again would spin.
s.log.Printf("refusing a message the mesh has no name for: %q", m.Kind())
_ = m.Drop()
}
}
// decide asks the window what to do about one message, and holds it when that is the answer —
// because holding is the one verdict that means the same thing everywhere.
//
// The other three come back, because "settled without acting" is a build result dropped and an
// upgrade taken, and only the handler knows which its message is. Hold means the message is the
// bus's problem now and **the caller must not settle it**; that is also what a cancelled context
// gets, because shutting down is not an answer about a message (novox/hq issue 083, on review).
//
// declaredIn is the declaration this message is about, empty when it is about none; outstanding is
// what the mesh last sent that node, empty when it has sent none or could not be asked.
func (s *Server) decide(ctx context.Context, m Control, what, declaredIn, outstanding string,
err error) Verdict {
if ctx.Err() != nil {
return Hold
}
window := StoreWindow{GiveUpAfter: s.giveUpAfter()}
switch v := window.Decide(err, declaredIn, outstanding, m.HeldFor()); v {
case Hold:
// Said once, on the first hold. Every redelivery saying it again would fill the log with
// one store restart.
first := m.HeldFor() == 0
if held := m.Hold(RedeliverAfter(m.HeldFor())); held != nil {
s.log.Printf("LOST %s: it could not be held for the store (%v): %v", what, held, err)
return GiveUp
}
if first {
s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err)
}
return Hold
case Stale:
s.log.Printf("set aside %s: the mesh has moved past that declaration", what)
return Stale
case GiveUp:
s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err)
return GiveUp
default:
return v
}
}
func (s *Server) giveUpAfter() time.Duration {
if s.giveUp == 0 {
return GiveUpAfter
}
return s.giveUp
}
// outstanding is the digest of the declaration the mesh last sent a node.
//
// Nothing is guessed when it cannot be answered: a listener that keeps no record of what was sent
// gives a report nothing to be stale against, and a store that cannot be read will hold the report
// anyway, so there is nothing for the check to decide.
func (s *Server) outstanding(ctx context.Context, node string) string {
if node == "" {
return ""
}
asks, ok := s.listener.(Outstanding)
if !ok {
return ""
}
digest, err := asks.Outstanding(ctx, node)
if err != nil {
return ""
}
return digest
}
// heartbeat records that a node was heard from, and nothing else.
//
// Deliberately silent: a node saying it is there every minute would fill the log with the ordinary
// case, and a log where the ordinary case is loud is a log nobody reads. Not held for the store
// either — the next heartbeat is a minute away, and a heartbeat kept for two minutes to be written
// late says nothing the one after it will not say better.
func (s *Server) heartbeat(m Control) {
var alive Alive
if err := json.Unmarshal(m.Body(), &alive); err != nil || alive.Node == "" {
_ = m.Drop()
return
}
if s.listener != 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)
}
}
_ = m.Took()
}
// reported records what a node says it did.
//
// A node states; nothing here writes anything the node claimed about itself beyond that it was
// heard from. What it applied is its own account of its own machine, and the mesh keeps the last
// one as a copy for recovery rather than as a source (novox/hq 09-the-node-lifecycle).
func (s *Server) reported(ctx context.Context, m Control) {
var report Report
if err := json.Unmarshal(m.Body(), &report); err != nil {
s.log.Printf("a report could not be read: %v", err)
_ = m.Drop()
return
}
what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared)
if s.listener != nil {
// **Asked before the store, not after** (design 25 §3). A report about a declaration the
// mesh has moved past would otherwise wait out a restarting store to be written and then
// overwrite what the node is doing now — and on the bus being built, where the holding is
// the server's, it comes back after the newer was applied whatever the controller does.
outstanding := s.outstanding(ctx, report.Node)
declaredIn := staleAgainst(report)
if Superseded(declaredIn, outstanding) {
s.log.Printf("set aside %s: the mesh has moved past that declaration", what)
_ = m.Took()
return
}
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,
// and a report lost here is a node the mesh never hears from again — the store
// restarting under the adoption that node just applied lost exactly that (issue 082).
return
case Stale, GiveUp:
_ = m.Took()
return
}
if err != nil {
// Said rather than swallowed. A report the mesh heard and failed to write down is a
// 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 {
case report.Refused != "":
s.log.Printf("%s refused a declaration: %s", report.Node, report.Refused)
case len(report.Failed) > 0:
s.log.Printf("%s applied %d and failed: %v", report.Node, len(report.Applied), report.Failed)
default:
s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied))
}
_ = m.Took()
}
// staleAgainst is the declaration a report may be judged stale against, and it is empty for a
// report that is not only an account of an apply.
//
// **A report carries two different things, and only one of them is about a declaration.** What the
// node applied is; what the machine *is* — the tunnel it took over, the ports its own bundle holds,
// what an adopted node found and is keeping, a node moving its overlay key — is not. Those reach
// the mesh on a report because a report is the message a node sends, and nowhere else: a rekey
// dropped as stale is a node whose overlay key never moves, and no retry is coming, because the
// node said it once.
//
// So staleness is asked only of a report that is purely an apply's account. The rest is acted on
// whenever it arrives, which is the behaviour the mesh has had all along.
func staleAgainst(report Report) string {
if report.Rekey != nil || report.Tunnel != nil || len(report.Held) > 0 ||
report.Firewall != "" || len(report.Reachable) > 0 || len(report.Carried) > 0 ||
len(report.Outward) > 0 {
return ""
}
return report.Declared
}
func (s *Server) enrolling(ctx context.Context, m Control) {
reply := EnrolReply{Refusal: "that token cannot be used"}
var request EnrolRequest
if err := json.Unmarshal(m.Body(), &request); err != nil {
s.log.Printf("an enrolment request could not be read: %v", err)
} else {
request.Redelivered = m.Redelivered()
accepted, err := s.enroller.Enrol(ctx, request)
switch {
case errors.Is(err, ErrTryAgain):
// Not a refusal: nothing was spent, and the same request asked again will be
// answered. Replied at once rather than held, so the node — which is waiting on this
// answer — decides when to ask, and the queue behind it moves (issue 083).
reply = EnrolReply{TryAgain: true, Refusal: "the mesh cannot answer right now; ask again"}
s.log.Printf("asked %q to enrol again shortly: %v", request.Node, err)
case err != nil && request.Redelivered:
// Said as what it most likely is: the bus handed this request over again after the
// controller stopped mid-answer, and an enrolment already spent is not finished a
// second time. The node may need a new token.
s.log.Printf("refusing a redelivered enrolment for %q — it may have finished before "+
"the control plane stopped, and if the node did not get its answer it needs a new "+
"token: %v", request.Node, err)
case err != nil:
// Logged in full here, where an operator can see it; sent back as one refusal, so
// that somebody guessing learns nothing from which reason came back.
s.log.Printf("refusing enrolment for %q: %v", request.Node, err)
default:
reply = accepted
s.log.Printf("enrolled %s", accepted.Node)
}
}
if body, err := json.Marshal(reply); err != nil {
s.log.Printf("cannot encode a reply: %v", err)
} else {
answer, cancel := context.WithTimeout(ctx, 10*time.Second)
if err := m.Answer(answer, body); err != nil {
s.log.Printf("cannot answer an enrolment: %v", err)
}
cancel()
}
// Settled after the answer is sent, so a controller that dies mid-answer leaves the request on
// the bus rather than having consumed it silently. Asked again by the same presenter, an
// enrolment finishes: the token is held for its key and spent last (issue 083).
_ = m.Took()
}
// wasBuilt keeps what a builder said, whichever way it went.
//
// This is for results nobody was waiting for. A build asked for with `build` is answered directly
// to the asker; one triggered any other way is published here, and without this it would be
// reported into the void — which is the same as not reporting it.
func (s *Server) wasBuilt(ctx context.Context, m Control) {
var result BuildResult
if err := json.Unmarshal(m.Body(), &result); err != nil {
s.log.Printf("a build result could not be read: %v", err)
_ = m.Drop()
return
}
if s.recorder == nil {
// Nothing to keep it in. Dropped rather than swallowed, so the bus's own counters show
// something arriving that nothing handles.
s.log.Printf("a build result arrived and this control plane keeps none")
_ = m.Drop()
return
}
err := s.recorder.Built(ctx, result)
switch s.decide(ctx, m, fmt.Sprintf("a build result from %s", result.On), "", "", err) {
case Hold:
// Held while the store cannot take it: a build result lost here is never announced (083).
return
case Stale, GiveUp:
_ = m.Drop()
return
}
if err != nil {
s.log.Printf("cannot keep a build result from %s: %v", result.On, err)
_ = m.Drop()
return
}
switch {
case result.Failed != "":
s.log.Printf("%s could not build %s", result.On, result.Repository)
default:
s.log.Printf("%s built %s from %s", result.On, result.Repository, result.Commit)
}
_ = m.Took()
}
// catchingUp answers a catalogue that has just started and may have missed builds.
//
// Settled after the work. A replay that fails for a reason other than the store is not one that
// succeeds by being handed the same request again, so that is settled and said; but a store that
// could not be read right now is held and asked again, bounded, rather than the request lost until
// the catalogue next restarts (issue 083). One request stands for all: a newer one supersedes one
// still held.
func (s *Server) catchingUp(ctx context.Context, m Control) {
if s.replayer == nil {
s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay")
_ = m.Took()
return
}
announcements, err := s.replayer.Announceable(ctx)
switch s.decide(ctx, m, "a catalogue's request to catch up", "", "", err) {
case Hold:
return
case Stale, GiveUp:
_ = m.Took()
return
}
if err != nil {
s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err)
_ = m.Took()
return
}
sent := 0
for _, a := range announcements {
a.Replay = true
// 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).
body, err := json.Marshal(a)
if err != nil {
// A body that cannot be written is this program's fault, not the bus's, and publishing
// an empty one would put a fact on the mesh that says nothing.
s.log.Printf("cannot re-announce %s at %s: %v", a.Module, short(a.Commit), err)
_ = m.Took()
return
}
if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, KeyBuiltBefore, body); 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",
a.Module, short(a.Commit), err)
_ = m.Took()
return
}
sent++
}
s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent)
_ = m.Took()
}
// moved hands one announcement to whatever is following them.
//
// **Settled whatever happens, but one thing.** A failure here is usually the controller being
// unable to act on an upgrade — a machine that cannot be resolved, a bus that will not take a
// declaration — and none of those get better by being handed the same message again. The one
// exception is the upgrader saying the store could not be read for the moment (ErrTryAgain): that
// is held and asked again, bounded (novox/hq issue 083). Only the upgrader's word counts here, not
// an error that merely looks like an outage — a push that timed out on the second machine is not
// asked again, or the first would be pushed every few seconds for two minutes. A newer move of the
// same module supersedes one still held.
func (s *Server) moved(ctx context.Context, m Control) {
var u Upgraded
if err := json.Unmarshal(m.Body(), &u); err != nil {
s.log.Printf("an upgrade announcement could not be read: %v", err)
_ = m.Took()
return
}
if u.Module == "" {
s.log.Printf("an upgrade announcement named no module; ignored")
_ = m.Took()
return
}
err := s.upgrader.Upgraded(ctx, u)
// Only "not now" is worth holding. Anything else is an answer, and the window would read a
// timed-out push as an outage and ask for it again.
holdable := err
if !errors.Is(err, ErrTryAgain) {
holdable = nil
}
what := fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit))
switch s.decide(ctx, m, what, "", "", holdable) {
case Hold:
return
case Stale, GiveUp:
_ = m.Took()
return
}
if err != nil {
s.log.Printf("%s moved to %s and the mesh could not act on it: %v",
u.Module, short(u.Commit), err)
}
_ = m.Took()
}
// digest names a message by its content.
func digest(body []byte) string {
sum := sha256.Sum256(body)
return hex.EncodeToString(sum[:])
}
// short is a commit as people read it.
func short(commit string) string {
if len(commit) > 8 {
return commit[:8]
}
return commit
}
// sourceMoved acts on the forge's announcement of a merge. Taken whatever happens: a build that
// fails is reported by the build itself, and re-delivering the merge would only re-fail it.
func (s *Server) sourceMoved(ctx context.Context, m Control) {
var moved SourceMoved
if err := json.Unmarshal(m.Body(), &moved); err != nil {
s.log.Printf("a merge announcement could not be read: %v", err)
_ = m.Took()
return
}
if moved.Commit == "" || moved.Repo == "" {
s.log.Printf("a merge announcement named no repository or no commit; ignored")
_ = m.Took()
return
}
if err := s.upgrader.SourceMoved(ctx, moved); err != nil {
s.log.Printf("%s/%s moved to %.8s and the mesh could not act on it: %v", moved.Owner, moved.Repo, moved.Commit, err)
}
_ = 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)
}
}