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 }