package link import ( "context" "crypto/ed25519" "encoding/json" "errors" "fmt" "net/url" "time" amqp "github.com/rabbitmq/amqp091-go" ) // ErrForged is what a node returns for a declaration whose signature is not the mesh's. // // Its own error, and it must never be confused with a malformed message. novox/hq ADR 0004 // requires a host to tell *this is not from the mesh I joined* apart from *this is malformed*: // the first means somebody is trying, the second means something is broken. var ErrForged = errors.New("this declaration was not signed by the mesh this node joined") // Membership is what a node needs to reach its mesh again, held by the caller. type Membership struct { Node string Broker string Fingerprint string Password string Signer ed25519.PublicKey } // Applier is what the host does with a declaration that has been proved to come from the mesh. type Applier func(ctx context.Context, declaration []byte) Report // Run holds the link open, applying what arrives and reporting what happened. // // Outbound only, and nothing listens on this machine. The connection is the node's presence in // the mesh: while it is up the node is enrolled, and while it is down the node is disconnected — // which is an ordinary situation and not a failure, so this returns rather than panicking and // leaves restarting to whatever supervises it. func Run(ctx context.Context, m Membership, apply Applier, timeout time.Duration) error { config, err := PinnedConfig(m.Fingerprint) if err != nil { return err } dsn := fmt.Sprintf("amqps://%s:%s@%s/", url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker) conn, err := amqp.DialConfig(dsn, amqp.Config{ TLSClientConfig: config, Dial: amqp.DefaultDial(timeout), // Kept short so a node that has silently lost its route notices, rather than holding a // connection the broker forgot about and believing it is still in the mesh. Heartbeat: 10 * time.Second, }) if err != nil { if errors.Is(err, ErrWrongCertificate) { return err } return fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err) } defer conn.Close() channel, err := conn.Channel() if err != nil { return err } defer channel.Close() queue := QueueFor(m.Node) if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil { return fmt.Errorf("cannot declare this node's queue %s: %w", queue, err) } // One at a time. A declaration is applied to a machine, and applying two at once would race // on the same filesystem — so the broker holds the next one until this one is finished, // where it survives a restart. if err := channel.Qos(1, 0, false); err != nil { return err } deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil) if err != nil { return err } closed := conn.NotifyClose(make(chan *amqp.Error, 1)) for { select { case <-ctx.Done(): return nil case reason := <-closed: return fmt.Errorf("the link closed: %v", reason) case delivery, ok := <-deliveries: if !ok { return errors.New("the broker stopped delivering") } report := handle(ctx, m, apply, delivery) publishReport(ctx, channel, m, report, timeout) // Acknowledged after the report is published. A node that dies between applying and // reporting leaves the declaration on the broker and applies it again on return, // which is safe because applying is reconciliation — it converges rather than // repeating. _ = delivery.Ack(false) } } } func handle(ctx context.Context, m Membership, apply Applier, delivery amqp.Delivery) Report { return handleBody(ctx, m, delivery.Body, apply) } // handleBody is the whole of deciding whether to trust a message, separated from the broker so it // can be tested as the security check it is rather than as message plumbing. func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) Report { var signed Signed if err := json.Unmarshal(body, &signed); err != nil { return Report{Node: m.Node, Refused: "this message is not a declaration: " + err.Error()} } // Before anything is read out of it, let alone applied. The host applies whatever the link // delivers, so this check is the difference between the mesh changing this machine and // anybody changing it. if !ed25519.Verify(m.Signer, signed.Declaration, signed.Signature) { return Report{Node: m.Node, Refused: ErrForged.Error()} } return apply(ctx, signed.Declaration) } func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report, timeout time.Duration) { report.Node = m.Node body, err := json.Marshal(report) if err != nil { return } publish, cancel := context.WithTimeout(ctx, timeout) defer cancel() _ = channel.PublishWithContext(publish, Exchange, KeyReport, false, false, amqp.Publishing{ContentType: "application/json", Body: body}) }