package link import ( "context" "crypto/ed25519" "encoding/json" "errors" "fmt" "log" "os" "strings" "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, secret string, public ed25519.PublicKey, profile map[string]any) (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 } type Server struct { conn *amqp.Connection channel *amqp.Channel enroller Enroller listener Listener log *log.Logger } // Connect opens the control plane's own connection to the broker. func Connect(enroller Enroller, listener Listener) (*Server, error) { url := strings.TrimSpace(os.Getenv(AMQPVar)) 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) } 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} { 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 { // Prefetch of one. The control plane writes to a database per message, and a burst of // enrolments delivered all at once would be held in memory rather than left on the broker, // which is the one place they survive a restart. if err := s.channel.Qos(1, 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 } closed := s.conn.NotifyClose(make(chan *amqp.Error, 1)) s.log.Printf("consuming %s, bound to %s/{%s,%s}", ControlQueue, Exchange, KeyEnrol, KeyReport) for { select { case <-ctx.Done(): return nil 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(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) } } // 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(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 } if s.listener != nil { if err := s.listener.Heard(context.Background(), report); 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)) } _ = delivery.Ack(false) } 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 { accepted, err := s.enroller.Enrol(ctx, request.Secret, request.PublicKey, request.Profile) if 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) } else { 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. Enrolment is idempotent // only in the sense that the token is spent — a redelivery gets the refusal, which is // correct and visible, where a lost request is neither. _ = 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) } }