Files
mesh-controller/internal/link/serve.go
T
jschoubben 3907ea0db0 The control plane serves and pushes on the bus it is told to
The seams were there and nothing chose a side: serve, push, ask and build all
opened the old bus's connection and declared over its channel, whatever
MESH_BUS_NATS said. So the switch moved every host and left the control plane
unable to follow — "this control plane has no MESH_BROKER_AMQP" with the new bus
named and standing (2026-09-28). That was task 4.3 of design 28, still open.

One place now decides: connectLink reads the switch, refuses both buses named at
once, raises the new bus's streams and this controller's consumers when it is
handed the inventory, and opens the link over whichever bus it is on. Every
caller that sent a declaration or asked a tool through the old channel goes
through the server's bus instead, which the new transport has and the channel is
not. OverNats is that outbound: a declaration is a JetStream publish into the
node's own subject, an event is announced on the subject its name derives to, a
tool is request and reply on the module's tool subject.
2026-09-28 01:29:44 +02:00

606 lines
22 KiB
Go

package link
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"github.com/novox/mesh-controller/internal/broker"
"log"
"os"
"time"
amqp "github.com/rabbitmq/amqp091-go"
"github.com/novox/mesh-controller/internal/envfile"
)
// 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(ctx context.Context, report Report) 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
}
// 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
conn *amqp.Connection
channel *amqp.Channel
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(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.
func Connect(enroller Enroller, listener Listener) (*Server, error) {
url, err := envfile.Placed(AMQPVar)
if err != nil {
return nil, err
}
if url == "" {
return nil, fmt.Errorf(
"this control plane has no %s, so it cannot reach its broker. Nodes talk to it over "+
"the broker and nowhere else, so without this it can hold records and answer "+
"nothing", AMQPVar)
}
conn, err := amqp.Dial(url)
if err != nil {
// Not quoted back: the URL carries the controller's own bus password.
return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err)
}
channel, err := conn.Channel()
if err != nil {
conn.Close()
return nil, err
}
// Declared here rather than assumed. The controller is the only thing that may create them — a
// node's account can write to this exchange and read its own queue, and configure nothing
// else, so a node arriving before the controller has ever run finds nothing and says so,
// rather than quietly creating a topology nobody designed.
if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err)
}
// The events exchange too. The controller is not the only publisher on it — modules announce
// onto it with their own accounts — but it is the only thing permitted to create it, for the
// same reason it is the only thing permitted to create the direct one.
if err := channel.ExchangeDeclare(EventsExchange, "topic", true, false, false, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s exchange: %w", EventsExchange, err)
}
if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err)
}
// Every key a node may publish. Binding one and forgetting another is a message the bus
// accepts, finds no queue for, and drops — the publisher sees success and the consumer sees
// nothing. That is exactly what happened to reports: `report` was left unbound while `enrol`
// worked, so nodes announced what they had applied into a void for an afternoon.
for _, key := range []string{KeyEnrol, KeyReport, KeyAlive, KeyBuilt} {
if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err)
}
}
return &Server{
inbound: Current(conn, channel),
bus: OverCurrent{Channel: channel},
conn: conn,
channel: channel,
enroller: enroller,
listener: listener,
log: newLog(),
}, nil
}
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
// Channel is the controller's channel, for the command line's own publishing.
func (s *Server) Channel() *amqp.Channel { return s.channel }
func (s *Server) Close() {
if s.inbound != nil {
s.inbound.Close()
}
if s.channel != nil {
_ = s.channel.Close()
}
if s.conn != nil {
_ = s.conn.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 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
}
m.About("report " + report.Node)
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
}
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)
}
}
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 {
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
}
// Each build result its own subject: none supersedes another, and recording one twice is
// harmless — the build is kept by its id.
m.About("build " + digest(m.Body()))
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) {
m.About("catch-up")
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
if err := EmitEvent(ctx, s.bus, KeyModuleBuilt, "control-plane", "", a); 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
}
m.About("upgrade " + u.Module)
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
}