package link import ( "context" "encoding/json" "errors" "fmt" "net/url" "time" amqp "github.com/rabbitmq/amqp091-go" ) // The wire format shared with the control plane, which defines it separately because this binary // requires nothing present and does not import it. A test on each side asserts the field names. const ( Exchange = "mesh" KeyEnrol = "enrol" ) // QueueFor is the queue this node consumes from — the only one its account may read. func QueueFor(node string) string { return "node." + node } // EnrolRequest is what this node says when joining. type EnrolRequest struct { Node string `json:"node"` Secret string `json:"secret"` PublicKey []byte `json:"public_key"` Profile map[string]any `json:"profile,omitempty"` } // EnrolReply is what the mesh says back. type EnrolReply struct { Accepted bool `json:"accepted"` Node string `json:"node,omitempty"` Queue string `json:"queue,omitempty"` // What this node keeps so it can come back on its own. Without these a restart would need a // person with a new token, which would make disconnection a crisis rather than the ordinary // situation novox/hq ADR 0004 says it is. Password string `json:"password,omitempty"` Broker string `json:"broker,omitempty"` Fingerprint string `json:"fingerprint,omitempty"` Signer []byte `json:"signer,omitempty"` Refusal string `json:"refusal,omitempty"` } // ErrRefused is what a node gets when the mesh will not have it. var ErrRefused = errors.New("the mesh refused this enrolment") // Enrol presents this node's key and its one-time secret, and waits to be told it is known. // // The broker has already authenticated this connection: the account was created when the token // was issued and the secret is its password. So this is not how the node gets in — it is what it // says once it is in, and the secret travels again because the control plane must not have to ask // the broker who connected. func Enrol(ctx context.Context, address, pin, node, secret string, public []byte, profile map[string]any, timeout time.Duration) (EnrolReply, error) { config, err := PinnedConfig(pin) if err != nil { return EnrolReply{}, err } // The account name is the node's, and the password is the token's secret. Escaped because a // name or secret containing a colon or an at-sign would otherwise change which host this // connects to — a credential silently redirecting a connection is the worst shape this could // take. dsn := fmt.Sprintf("amqps://%s:%s@%s/", url.QueryEscape(node), url.QueryEscape(secret), address) conn, err := amqp.DialConfig(dsn, amqp.Config{ TLSClientConfig: config, Dial: amqp.DefaultDial(timeout), }) if err != nil { if errors.Is(err, ErrWrongCertificate) { return EnrolReply{}, err } // Not quoted back: the DSN carries the one-time secret. return EnrolReply{}, fmt.Errorf("cannot reach the broker at %s as %s: %w", address, node, err) } defer conn.Close() channel, err := conn.Channel() if err != nil { return EnrolReply{}, err } defer channel.Close() // This node's own queue, which its account is scoped to and nothing else may read. queue, err := channel.QueueDeclare(QueueFor(node), true, false, false, false, nil) if err != nil { return EnrolReply{}, fmt.Errorf( "cannot declare this node's queue %s: %w", QueueFor(node), err) } replies, err := channel.Consume(queue.Name, "", true, false, false, false, nil) if err != nil { return EnrolReply{}, err } request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, Profile: profile} body, err := json.Marshal(request) if err != nil { return EnrolReply{}, err } correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano()) publish, cancel := context.WithTimeout(ctx, timeout) defer cancel() if err := channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false, amqp.Publishing{ ContentType: "application/json", CorrelationId: correlation, ReplyTo: queue.Name, Body: body, }); err != nil { return EnrolReply{}, fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err) } // Waited for rather than assumed. A published message that nothing answers means the control // plane is not running, and a node that carried on regardless would believe it had joined a // mesh that has never heard of it. deadline := time.NewTimer(timeout) defer deadline.Stop() closed := conn.NotifyClose(make(chan *amqp.Error, 1)) for { select { case <-ctx.Done(): return EnrolReply{}, ctx.Err() case reason := <-closed: return EnrolReply{}, fmt.Errorf("the broker closed the connection: %v", reason) case <-deadline.C: return EnrolReply{}, fmt.Errorf( "the broker accepted this node's connection and nothing answered within %s. The "+ "mesh's broker is running and its control plane is not", timeout) case delivery, ok := <-replies: if !ok { return EnrolReply{}, errors.New("the broker stopped delivering") } // Anything else on this queue is not the answer to this question. if delivery.CorrelationId != correlation { continue } var reply EnrolReply if err := json.Unmarshal(delivery.Body, &reply); err != nil { return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err) } if !reply.Accepted { return reply, fmt.Errorf("%w: %s", ErrRefused, reply.Refusal) } return reply, nil } } }