package link import ( "context" "crypto/sha256" "encoding/hex" "encoding/json" "errors" "fmt" "github.com/novox/mesh-controller/internal/envfile" "github.com/novox/mesh-controller/internal/inventory" "log" "os" "time" amqp "github.com/rabbitmq/amqp091-go" ) // AMQPVar is the control plane's own connection to the broker. const AMQPVar = "MESH_BROKER_AMQP" // Enroller is what the control plane does with an enrolment request. // // An interface so the serving loop can be tested against a real broker 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) } // Server consumes what nodes say. // Listener is what the control plane 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 } type Server struct { conn *amqp.Connection channel *amqp.Channel enroller Enroller listener Listener recorder Recorder log *log.Logger upgrader Upgrader replayer Replayer // Messages the store could not take right now, held unacknowledged and tried again on a // ticker, by subject (novox/hq issues 082, 083). again is the ticker's interval, zero meaning // TryAgainAfter; giveUp is how long one is kept, zero meaning GiveUpAfter. again time.Duration giveUp time.Duration parked map[string]*held } // held is one message the store could not take, kept to be tried again. type held struct { delivery amqp.Delivery retry func(context.Context, amqp.Delivery) what string first time.Time } // 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") // TryAgainAfter is how often messages the store could not take are tried again. A store comes // back in seconds, and a report a few seconds late is still current. const TryAgainAfter = 2 * time.Second // 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 // Prefetch is how many messages the broker hands the control plane before it has settled them. // More than one because a message the store could not take is held, unsettled, while the loop goes // on answering others — an enrolment above all, which a host is waiting on (novox/hq issue 083). // Bounded, because what is held is also what the broker has not kept on its own disk as pending. const Prefetch = 64 // PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to // answer — an enrolment above all — while others wait for the store. const PrefetchHeadroom = 8 // Records tells the server where to keep build results. // // Set after Connect rather than passed to it, because a control plane 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 } // Upgrader is what the control plane 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 requeued: an upgrade the control plane could not act on is not one it will // act on by being handed 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 } // Follows says what to do about upgrades, and binds the queue they arrive on. // // **Not bound unless something is listening.** A durable queue bound to every upgrade with no // consumer fills up quietly, and the first symptom is a broker out of disk rather than anything // about modules. func (s *Server) Follows(u Upgrader) error { if _, err := s.channel.QueueDeclare(UpgradeQueue, true, false, false, false, nil); err != nil { return fmt.Errorf("cannot declare the %s queue: %w", UpgradeQueue, err) } if err := s.channel.QueueBind(UpgradeQueue, KeyModuleUpgraded, EventsExchange, false, nil); err != nil { return fmt.Errorf("cannot bind %s to %s/%s: %w", UpgradeQueue, EventsExchange, KeyModuleUpgraded, err) } s.upgrader = u return nil } // Answers binds the queue a catalogue's catch-up request arrives on. // // **Not bound unless something is listening**, for the same reason upgrades are not: a durable // queue with no consumer fills quietly and the first symptom is a broker out of disk. func (s *Server) Answers(r Replayer) error { if _, err := s.channel.QueueDeclare(CatchUpQueue, true, false, false, false, nil); err != nil { return fmt.Errorf("cannot declare the %s queue: %w", CatchUpQueue, err) } if err := s.channel.QueueBind(CatchUpQueue, KeyCatchingUp, EventsExchange, false, nil); err != nil { return fmt.Errorf("cannot bind %s to %s/%s: %w", CatchUpQueue, EventsExchange, KeyCatchingUp, err) } s.replayer = r return nil } // Connect opens the control plane's own connection to the broker. // // 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 control plane's own broker 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 control plane 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 control plane 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 control plane 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 broker // 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{conn: conn, channel: channel, enroller: enroller, listener: listener, log: log.New(os.Stdout, "", log.LstdFlags)}, nil } // Channel is the control plane's channel, for sending declarations. func (s *Server) Channel() *amqp.Channel { return s.channel } func (s *Server) Close() { if s.channel != nil { _ = s.channel.Close() } if s.conn != nil { _ = s.conn.Close() } } // Serve consumes until the context ends. // // One consumer, deliberately: with two, the broker would round-robin between them and each would // receive half of what it expects — a fault this project has already had, between a module's // daemon and its capability server. func (s *Server) Serve(ctx context.Context) error { // A bounded prefetch rather than one. The loop still takes messages one at a time; what the // prefetch buys is that a message the store could not take can be held while the loop goes on // to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything held // goes back to the broker if the control plane stops, because nothing held is acknowledged. if err := s.channel.Qos(Prefetch, 0, false); err != nil { return err } deliveries, err := s.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane", false, false, false, false, nil) if err != nil { return err } // The upgrade queue, when something is listening for them. A second queue rather than a // second consumer on the first: two consumers on one queue split its messages between them, // which is the fault the comment above exists about. Two queues share nothing. var upgrades <-chan amqp.Delivery if s.upgrader != nil { upgrades, err = s.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades", false, false, false, false, nil) if err != nil { return err } } // Its own queue and its own consumer, for the reason above: two consumers on one queue split // its messages, and a catch-up request going to whichever half was not listening is a gap that // looks like a working mesh. var catchups <-chan amqp.Delivery if s.replayer != nil { catchups, err = s.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup", false, false, false, false, nil) if err != nil { return err } } closed := s.conn.NotifyClose(make(chan *amqp.Error, 1)) s.log.Printf("consuming %s, bound to %s/{%s,%s,%s,%s}", ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt) if s.upgrader != nil { s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded) } again := s.again if again == 0 { again = TryAgainAfter } ticker := time.NewTicker(again) defer ticker.Stop() for { select { case <-ctx.Done(): return nil case <-ticker.C: if ctx.Err() != nil { return nil } s.retryHeld(ctx) case delivery, ok := <-catchups: if !ok { if catchups != nil { return errors.New("the broker stopped delivering catch-up requests") } continue } s.catchingUp(ctx, delivery) case delivery, ok := <-upgrades: // A nil channel blocks for ever, so this case simply never fires when nothing is // listening for upgrades. Closed is different, and means the broker stopped. if !ok { if upgrades != nil { return errors.New("the broker stopped delivering upgrades") } continue } s.upgraded(ctx, delivery) case reason := <-closed: // Said rather than returned quietly. A control plane whose broker connection dropped // is a mesh where nothing can be told anything, and the reason is the first thing // anybody will want. return fmt.Errorf("the broker connection closed: %v", reason) case delivery, ok := <-deliveries: if !ok { return errors.New("the broker stopped delivering") } s.handle(ctx, delivery) } } } func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) { switch delivery.RoutingKey { case KeyEnrol: s.handleEnrol(ctx, delivery) case KeyReport: s.handleReport(ctx, delivery) case KeyAlive: s.handleAlive(delivery) case KeyBuilt: s.handleBuilt(ctx, delivery) default: // Rejected without requeue: a message nothing understands will not be understood on the // next attempt either, and requeuing it would spin. s.log.Printf("refusing a message with routing key %q", delivery.RoutingKey) _ = delivery.Reject(false) } } // handleAlive 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. func (s *Server) handleAlive(delivery amqp.Delivery) { var alive Alive if err := json.Unmarshal(delivery.Body, &alive); err != nil || alive.Node == "" { _ = delivery.Reject(false) 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) } } _ = delivery.Ack(false) } // handleReport 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) handleReport(ctx context.Context, delivery amqp.Delivery) { var report Report if err := json.Unmarshal(delivery.Body, &report); err != nil { s.log.Printf("a report could not be read: %v", err) _ = delivery.Reject(false) return } // A node's newer report supersedes one of its older reports still held: the older is its // past, and recorded after the newer it would overwrite what the node is doing now. subject := "report " + report.Node s.supersede(subject, delivery) if s.listener != nil { if err := s.listener.Heard(context.Background(), report); err != nil { // Held, not acknowledged, 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). what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) if s.tryLater(ctx, delivery, subject, what, err, s.handle) { return } // 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)) } s.settled(subject, delivery) _ = delivery.Ack(false) } // tryLater holds a message the store could not take right now, to be tried again on the ticker, // and says whether it did (novox/hq issues 082, 083). // // "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error the // inventory reads as an outage. Anything else is an answer, and is left to the caller to settle. // Held means unacknowledged and set aside: the loop goes on to the next message, so an enrolment a // host is waiting on is answered while a report waits for the store. One message is held at most // giveUpAfter; past it, it is let go with a line saying it was lost, and the caller settles it. func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, subject, what string, err error, retry func(context.Context, amqp.Delivery)) bool { // Shutting down: nothing is settled. Unsettled, the broker hands the message to whatever // consumes next — a cancelled context is not an answer about the message (issue 083, review). if ctx.Err() != nil { return true } if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) { return false } if s.parked == nil { s.parked = map[string]*held{} } h, ok := s.parked[subject] if !ok || h.delivery.DeliveryTag != delivery.DeliveryTag { // Held no further than the prefetch leaves room: past it, the broker would hand the loop // nothing new — enrolments included — until something held was let go. if !ok && len(s.parked) >= Prefetch-PrefetchHeadroom { s.log.Printf("LOST %s: %d messages are already held for the store, and holding more "+ "would stop the queue: %v", what, len(s.parked), err) return false } h = &held{delivery: delivery, retry: retry, what: what, first: time.Now()} s.parked[subject] = h s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err) return true } if time.Since(h.first) >= s.giveUpAfter() { delete(s.parked, subject) s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err) return false } return true } // supersede drops a message held for a subject when a newer one for it arrives: the older is // acknowledged, because acting on it after the newer would undo the newer. func (s *Server) supersede(subject string, newer amqp.Delivery) { h, ok := s.parked[subject] if !ok || h.delivery.DeliveryTag == newer.DeliveryTag { return } delete(s.parked, subject) s.log.Printf("set aside %s: a newer one arrived", h.what) _ = h.delivery.Ack(false) } // settled forgets a message once it has been handled either way. func (s *Server) settled(subject string, delivery amqp.Delivery) { if h, ok := s.parked[subject]; ok && h.delivery.DeliveryTag == delivery.DeliveryTag { delete(s.parked, subject) } } // digest names a message by its content. func digest(body []byte) string { sum := sha256.Sum256(body) return hex.EncodeToString(sum[:]) } // retryHeld tries every held message again. Each handler holds it again, settles it, or lets it // go past the bound. func (s *Server) retryHeld(ctx context.Context) { for _, h := range s.snapshot() { if ctx.Err() != nil { return } h.retry(ctx, h.delivery) } } func (s *Server) snapshot() []*held { out := make([]*held, 0, len(s.parked)) for _, h := range s.parked { out = append(out, h) } return out } func (s *Server) giveUpAfter() time.Duration { if s.giveUp == 0 { return GiveUpAfter } return s.giveUp } func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { reply := EnrolReply{Refusal: "that token cannot be used"} var request EnrolRequest if err := json.Unmarshal(delivery.Body, &request); err != nil { s.log.Printf("an enrolment request could not be read: %v", err) } else { request.Redelivered = delivery.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 broker handed this request over again after // the control plane 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) } } s.reply(ctx, delivery, reply) // Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the // request on the broker 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). _ = delivery.Ack(false) } func (s *Server) reply(ctx context.Context, delivery amqp.Delivery, reply EnrolReply) { if delivery.ReplyTo == "" { s.log.Print("an enrolment request named no reply queue, so nothing can be told the answer") return } body, err := json.Marshal(reply) if err != nil { s.log.Printf("cannot encode a reply: %v", err) return } timeout, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() if err := s.channel.PublishWithContext(timeout, "", delivery.ReplyTo, false, false, amqp.Publishing{ ContentType: "application/json", CorrelationId: delivery.CorrelationId, Body: body, }); err != nil { s.log.Printf("cannot reply to %s: %v", delivery.ReplyTo, err) } } // handleBuilt 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) handleBuilt(ctx context.Context, delivery amqp.Delivery) { var result BuildResult if err := json.Unmarshal(delivery.Body, &result); err != nil { s.log.Printf("a build result could not be read: %v", err) _ = delivery.Reject(false) return } if s.recorder == nil { // Nothing to keep it in. Rejected rather than dropped silently, so the broker's own // counters show something arriving that nothing handles. s.log.Printf("a build result arrived and this control plane keeps none") _ = delivery.Reject(false) return } // Each build result its own subject: none supersedes another, and recording one twice is // harmless — the build is kept by its id. subject := "build " + digest(delivery.Body) s.supersede(subject, delivery) if err := s.recorder.Built(ctx, result); err != nil { // Held while the store cannot take it: a build result lost here is never announced (083). if s.tryLater(ctx, delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) { return } s.log.Printf("cannot keep a build result from %s: %v", result.On, err) s.settled(subject, delivery) _ = delivery.Reject(false) return } s.settled(subject, delivery) 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) } _ = delivery.Ack(false) } // catchingUp answers a catalogue that has just started and may have missed builds. // // Acknowledged 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 acknowledged 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, delivery amqp.Delivery) { const subject = "catch-up" s.supersede(subject, delivery) holding := false defer func() { if !holding { s.settled(subject, delivery) _ = delivery.Ack(false) } }() if s.replayer == nil { s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay") return } announcements, err := s.replayer.Announceable(ctx) if err != nil { if s.tryLater(ctx, delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) { holding = true return } s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err) return } sent := 0 for _, a := range announcements { a.Replay = true if err := EmitEvent(ctx, s.channel, 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) return } sent++ } s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent) } // upgraded hands one announcement to whatever is following them. // // **Acknowledged whatever happens, but one thing.** A failure here is usually the control plane // being unable to act on an upgrade — a machine that cannot be resolved, a broker 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) upgraded(ctx context.Context, delivery amqp.Delivery) { var u Upgraded _ = json.Unmarshal(delivery.Body, &u) subject := "upgrade " + u.Module s.supersede(subject, delivery) holding := false defer func() { if !holding { s.settled(subject, delivery) _ = delivery.Ack(false) } }() if err := json.Unmarshal(delivery.Body, &u); err != nil { s.log.Printf("an upgrade announcement could not be read: %v", err) return } if u.Module == "" { s.log.Printf("an upgrade announcement named no module; ignored") return } if err := s.upgrader.Upgraded(ctx, u); err != nil { // Shutting down is not an answer about the announcement: left for the broker. if ctx.Err() != nil { holding = true return } if errors.Is(err, ErrTryAgain) && s.tryLater(ctx, delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) { holding = true return } s.log.Printf("%s moved to %s and the mesh could not act on it: %v", u.Module, short(u.Commit), err) } } // short is a commit as people read it. func short(commit string) string { if len(commit) > 8 { return commit[:8] } return commit }