Files
mesh-controller/internal/link/serve.go
T
jschoubben aecac5bda2 One bus: the AMQP transport is gone from the controller
The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The old
transport's consume loop, build request, tool ask, management API and account scoping are
deleted, and the bus switch with them; the controller connects to the broker seat and to
nothing else. The store-window tests keep their assertions on a bus-less fake, and the tests
that only made sense for the old transport's in-memory holding go with it.
2026-09-28 03:36:16 +02:00

554 lines
21 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(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
// 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
}
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
}
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
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
}
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()
}