The loop the whole thing exists for: told, apply, report. `run` holds one outbound connection open and consumes the node's own queue. Every declaration is verified against the control plane's signing key before a byte of it is read as an instruction -- not once at connect, every time. The transport being pinned is a different question from the instruction being genuine, and pinning only the first would make the second transitive: a compromised broker could forge declarations, and this host applies whatever the link delivers. Malformed and forged are reported differently, because ADR 0004 requires a host to tell "this is not from the mesh I joined" from "this is broken". One means somebody is trying and the other means something needs fixing. A node now keeps what it needs to come back on its own: the broker's address and fingerprint, the signing key it believes, and its own broker password -- which the mesh issues at enrolment to replace the token's secret, so the one-time thing stays one-time and the credential it holds for years is not the one that was pasted into a terminal. Verified in the lab end to end. The node enrolled, held its link, received a signed declaration and applied it -- the file is on the machine with the right contents, and the host's own record lists both resources. That run also found issue 010, which is recorded in novox/hq: the declaration removed every container on the machine, including the control plane that sent it. Correct reconciliation, shared store, and the first thing that happens.
141 lines
4.8 KiB
Go
141 lines
4.8 KiB
Go
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})
|
|
}
|