Review of the ADR 0105 build (hq ADR 0105). The takeover stopped the found unit and then found out whether the mesh's interface would do; a start that failed left the machine with no tunnel at all. Now nothing is stopped until the declared interface listens on the found port at the found address and the key file it names holds the found key — the refusal names the remedy — and a mesh interface that fails to start after the takeover has the found unit started again, with the account saying so. The account has three states (not taken, taken, down) and is given on every takeover, failure included. An interface raised by hand is looked at again for a moment and then refused naming `wg-quick down`. A found unit started again by hand beside the mesh's is said, not stopped: on the hub it cannot hold the port, and on a spoke two interfaces with one key would fight. `mesh-host overlay take --tunnel <iface>` is the path for a node that enrolled before the mesh knew to take a tunnel over: the found key becomes its overlay key — identity, sealing and serving keys untouched, so nothing sealed to the node is remade — and the mesh is told with a rekey signed by the identity key, over the key left, the key taken and the tunnel. Told first, written second, so a run again puts right whichever half did not happen.
452 lines
18 KiB
Go
452 lines
18 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")
|
|
|
|
// AliveEvery is how often a node says it is there.
|
|
//
|
|
// Often enough that "no word for five minutes" means something, rarely enough that a hundred
|
|
// nodes are not a hundred messages a second. The mesh reads absence rather than presence, so what
|
|
// matters is the interval being known and steady.
|
|
const AliveEvery = 60 * time.Second
|
|
|
|
// 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.
|
|
//
|
|
// It receives the signature as well as the declaration, so the host can keep both: what it was
|
|
// told is kept signed and verified again when it is read back, which means the file on disk is
|
|
// trusted for the same reason the message was rather than for being local.
|
|
type Applier func(ctx context.Context, declaration, signature []byte) Report
|
|
|
|
// Announce is how the link says what is happening, so a node running unattended leaves an
|
|
// account of it. Nil is allowed and means say nothing.
|
|
type Announce func(string)
|
|
|
|
// Hold keeps this node in the mesh, reconnecting for as long as it is asked to.
|
|
//
|
|
// Disconnection is an ordinary situation and not a failure (novox/hq ADR 0004), so this does not
|
|
// give up. A laptop shut for a week comes back and reconnects; it does not come back needing
|
|
// somebody to start it again.
|
|
//
|
|
// The backoff exists because the two common reasons differ in how long they last: a broker
|
|
// restarting is back in seconds, and a machine that has moved to a network with no route may be
|
|
// hours. Retrying every second for hours is a node shouting into nothing; waiting a minute after
|
|
// a broker blip is a node that is needlessly late. So it starts fast and slows down, and resets
|
|
// once a connection has actually held.
|
|
// Roused is a channel that says the machine has reason to believe its link is stale — it woke
|
|
// from suspend, or its network changed.
|
|
//
|
|
// **The machine knows before any timeout does.** A suspended laptop's connection is dead the
|
|
// moment it wakes, and heartbeats find that out in twenty or thirty seconds; for that time the
|
|
// node believes it is in the mesh and is not, which is the one state this design says must never
|
|
// be indistinguishable from being connected. Nothing new listens on the node to arrange it — the
|
|
// signal a service manager already sends is enough (novox/hq ADR 0004).
|
|
//
|
|
// Nil is allowed and means nothing ever rouses it, which is every machine that does not suspend.
|
|
type Roused <-chan struct{}
|
|
|
|
func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error {
|
|
return HoldRoused(ctx, m, apply, say, timeout, nil, nil)
|
|
}
|
|
|
|
// Outbox carries reports the node has to say without having been sent anything — what a
|
|
// reconcile found changed on an adopted node (novox/hq ADR 0100). Published while the link is up;
|
|
// a report made while it is down waits in the channel for the next one. Nil is allowed.
|
|
type Outbox <-chan Unasked
|
|
|
|
// Unasked is one such report, with the way to say whether it reached the mesh. Done is called
|
|
// with true only when the broker took it — a node that marked a change said because it queued it
|
|
// would never say it again, and the mesh would go on believing nothing changed.
|
|
type Unasked struct {
|
|
Report Report
|
|
Done func(published bool)
|
|
}
|
|
|
|
// HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed
|
|
// reports to publish between deliveries.
|
|
func HoldRoused(ctx context.Context, m Membership, apply Applier, say Announce,
|
|
timeout time.Duration, roused Roused, outbox Outbox) error {
|
|
|
|
return holdWith(ctx, func(ctx context.Context) error {
|
|
return Run(ctx, m, apply, say, timeout, outbox)
|
|
}, say, roused)
|
|
}
|
|
|
|
// attempt is one try at holding the link open, returning when it ends for any reason.
|
|
//
|
|
// Named so the loop below can be driven without a broker. What the loop decides — when to wait,
|
|
// how long, what being roused does — is the part with the reasoning in it, and it was reachable
|
|
// only through a real connection before.
|
|
type attempt func(context.Context) error
|
|
|
|
func holdWith(ctx context.Context, run attempt, say Announce, roused Roused) error {
|
|
const (
|
|
first = 2 * time.Second
|
|
most = 2 * time.Minute
|
|
// A connection that lasted this long counts as having worked, so the next failure starts
|
|
// from the bottom again. Without it a node that reconnects and immediately drops climbs
|
|
// to the maximum and stays there, long after whatever caused it went away.
|
|
settled = 30 * time.Second
|
|
)
|
|
wait := first
|
|
|
|
for {
|
|
began := time.Now()
|
|
|
|
// The link runs under a context this loop can cancel, so being roused ends the current
|
|
// attempt rather than only shortening the wait after it.
|
|
//
|
|
// **That is the whole of it.** After a resume the socket looks perfectly healthy from
|
|
// inside this process — there is no error and no close, because nothing has tried to
|
|
// send anything. It is heartbeats that eventually discover it, twenty or thirty seconds
|
|
// later. A machine that knows it just woke does not have to wait to be told.
|
|
//
|
|
// A rouse that turns out to be spurious costs one reconnect, which is cheap and
|
|
// idempotent: the node redeclares its queue and anything unacknowledged is redelivered.
|
|
// The alternative costs half a minute of believing it is in a mesh it has left.
|
|
trying, done := context.WithCancel(ctx)
|
|
if roused != nil {
|
|
go func() {
|
|
select {
|
|
case <-trying.Done():
|
|
case <-roused:
|
|
say("woken or moved — dropping the link and opening it again")
|
|
done()
|
|
}
|
|
}()
|
|
}
|
|
err := run(trying)
|
|
done()
|
|
if ctx.Err() != nil {
|
|
return nil
|
|
}
|
|
if time.Since(began) > settled {
|
|
wait = first
|
|
}
|
|
|
|
switch {
|
|
case errors.Is(err, ErrWrongCertificate):
|
|
// Said in full every time rather than folded into a retry count. This does not mean
|
|
// the network is down; it means what answered is not the mesh this node joined, and
|
|
// no amount of waiting fixes it. The node keeps running what it was last told, which
|
|
// is the right thing to do while somebody works out what happened.
|
|
say("the broker is not the one this node joined: " + err.Error())
|
|
say("this will not fix itself. This node keeps running what it was last told.")
|
|
case err != nil:
|
|
say(fmt.Sprintf("disconnected: %v — trying again in %s", err, wait))
|
|
default:
|
|
say(fmt.Sprintf("the link closed — trying again in %s", wait))
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case <-roused:
|
|
// And it does not serve out a wait computed for a broker that was restarting, either.
|
|
//
|
|
// The backoff is not *reset* by this. Being roused says the machine changed, not that
|
|
// whatever was refusing the connection has stopped — a laptop woken repeatedly on a
|
|
// network with no route would otherwise retry at full speed for as long as somebody
|
|
// keeps opening the lid.
|
|
say("woken or moved — trying again now")
|
|
case <-time.After(wait):
|
|
}
|
|
if wait *= 2; wait > most {
|
|
wait = most
|
|
}
|
|
}
|
|
}
|
|
|
|
// Run holds the link open once, applying what arrives and reporting what happened.
|
|
//
|
|
// Outbound only, and nothing listens on this machine. Returns when the link ends, for any reason;
|
|
// Hold is what decides whether to open it again.
|
|
func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration,
|
|
outbox Outbox) error {
|
|
if say == nil {
|
|
say = func(string) {}
|
|
}
|
|
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)
|
|
}
|
|
|
|
// Applying is one at a time — two at once would race on the same filesystem — but SEEING is
|
|
// not: with a prefetch of one the host could never know that a newer declaration was already
|
|
// waiting, and so applied every one of a backlog in turn, at the better part of a minute each,
|
|
// becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged
|
|
// deliveries lets it drain to the newest; each declaration still survives a restart on the
|
|
// broker until it is acknowledged, which happens only after it is applied or set aside.
|
|
if err := channel.Qos(drainDepth, 0, false); err != nil {
|
|
return err
|
|
}
|
|
|
|
deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Said, because it is the event anybody watching actually wants. Without it a node logs
|
|
// every failure and nothing on success, so a log full of "trying again" and then silence
|
|
// reads as still broken when it means the opposite.
|
|
say("in the mesh, consuming " + queue)
|
|
|
|
// A word every so often, so the mesh can tell a node that is quiet from one that is gone.
|
|
// Cheap on purpose: it carries a name and nothing else, because anything more would be a
|
|
// report, and reports are rare where this is constant.
|
|
beat := time.NewTicker(AliveEvery)
|
|
defer beat.Stop()
|
|
publishAlive(ctx, channel, m, say, timeout)
|
|
|
|
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
|
|
|
// Published mandatory, so the broker hands back anything it cannot route rather than
|
|
// dropping it. Without this a report goes to an exchange with no matching binding, the
|
|
// publisher is told nothing, and the mesh believes this node never answered while the node
|
|
// believes it did — which is what happened when `report` was left unbound on the other side.
|
|
returned := channel.NotifyReturn(make(chan amqp.Return, 4))
|
|
go func() {
|
|
for r := range returned {
|
|
say(fmt.Sprintf("the broker could not route this node's %s: %s (%d %s)",
|
|
r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText))
|
|
}
|
|
}()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case <-beat.C:
|
|
publishAlive(ctx, channel, m, say, timeout)
|
|
case unasked := <-outbox:
|
|
// Said without having been asked: a reconcile found what an adopted node holds, or
|
|
// its firewall, changed since it last said.
|
|
published := publishReport(ctx, channel, m, unasked.Report, say, timeout)
|
|
if unasked.Done != nil {
|
|
unasked.Done(published)
|
|
}
|
|
case reason := <-closed:
|
|
return fmt.Errorf("the link closed: %v", reason)
|
|
case delivery, ok := <-deliveries:
|
|
if !ok {
|
|
return errors.New("the broker stopped delivering")
|
|
}
|
|
// Whatever else is already waiting supersedes this one. Each set-aside declaration
|
|
// is reported as such, then acknowledged unapplied.
|
|
delivery, superseded := newest(deliveries, delivery, drainWindow)
|
|
for _, old := range superseded {
|
|
say("set aside a declaration: a newer one arrived with it")
|
|
publishReport(ctx, channel, m, Report{Node: m.Node, Declared: declaredIn(old.Body),
|
|
Superseded: declaredIn(delivery.Body)}, say, timeout)
|
|
_ = old.Ack(false)
|
|
}
|
|
report := handle(ctx, m, apply, delivery)
|
|
switch {
|
|
case report.Refused != "":
|
|
say("refused a declaration: " + report.Refused)
|
|
case len(report.Failed) > 0:
|
|
say(fmt.Sprintf("applied %d and failed: %v", len(report.Applied), report.Failed))
|
|
default:
|
|
say(fmt.Sprintf("applied %d resource(s)", len(report.Applied)))
|
|
}
|
|
publishReport(ctx, channel, m, report, say, 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)
|
|
}
|
|
}
|
|
}
|
|
|
|
// drainDepth is how many declarations the host will hold unacknowledged while it looks for a newer
|
|
// one; drainWindow is how long it waits for another to follow the one it has. Both small: a push is
|
|
// rare and a backlog is the exception this exists for, not the shape of ordinary traffic.
|
|
const (
|
|
drainDepth = 16
|
|
drainWindow = 750 * time.Millisecond
|
|
)
|
|
|
|
// newest takes what is already waiting behind `first` and returns the last of them to apply, and
|
|
// the rest to set aside. It waits `window` for a straggler after each arrival and no longer: a
|
|
// declaration in flight from the mesh arrives within that; one that does not is the next push.
|
|
func newest(deliveries <-chan amqp.Delivery, first amqp.Delivery, window time.Duration) (amqp.Delivery, []amqp.Delivery) {
|
|
latest := first
|
|
var superseded []amqp.Delivery
|
|
for {
|
|
select {
|
|
case next, ok := <-deliveries:
|
|
if !ok {
|
|
return latest, superseded
|
|
}
|
|
superseded = append(superseded, latest)
|
|
latest = next
|
|
case <-time.After(window):
|
|
return latest, superseded
|
|
}
|
|
}
|
|
}
|
|
|
|
// declaredIn is the id a signed declaration carries, for a report about one that was not applied.
|
|
// Empty if the message is not one — a forged or garbled message is refused by handleBody when its
|
|
// turn comes; here it is only named.
|
|
func declaredIn(body []byte) string {
|
|
var signed Signed
|
|
if err := json.Unmarshal(body, &signed); err != nil {
|
|
return ""
|
|
}
|
|
var d struct {
|
|
Declared string `json:"declared"`
|
|
}
|
|
if err := json.Unmarshal(signed.Declaration, &d); err != nil {
|
|
return ""
|
|
}
|
|
return d.Declared
|
|
}
|
|
|
|
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, signed.Signature)
|
|
}
|
|
|
|
// publishReport tells the mesh what this node did, and says whether the broker took it.
|
|
// Publish sends one report on this node's own connection and returns: the one-shot path for a
|
|
// report a command makes rather than the running host — a rekey (novox/hq ADR 0105). The same
|
|
// account, the same pinned certificate and the same exchange as the running host's reports.
|
|
func Publish(ctx context.Context, m Membership, report Report, 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),
|
|
})
|
|
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()
|
|
var said string
|
|
if !publishReport(ctx, channel, m, report, func(s string) { said = s }, timeout) {
|
|
return errors.New(said)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report,
|
|
say Announce, timeout time.Duration) bool {
|
|
report.Node = m.Node
|
|
body, err := json.Marshal(report)
|
|
if err != nil {
|
|
say("cannot encode this node's own report: " + err.Error())
|
|
return false
|
|
}
|
|
publish, cancel := context.WithTimeout(ctx, timeout)
|
|
defer cancel()
|
|
|
|
// Said rather than swallowed. A report that fails to publish leaves the mesh believing this
|
|
// node never answered, while the node believes it did — and the two would go on disagreeing
|
|
// with nothing anywhere saying so. That shape of fault is the one this project keeps finding.
|
|
if err := channel.PublishWithContext(publish, Exchange, KeyReport, true, false,
|
|
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
|
say(fmt.Sprintf("applied, and could not tell the mesh: %v", err))
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
// publishAlive says this node is here, and nothing else.
|
|
func publishAlive(ctx context.Context, channel *amqp.Channel, m Membership, say Announce,
|
|
timeout time.Duration) {
|
|
body, err := json.Marshal(Alive{Node: m.Node})
|
|
if err != nil {
|
|
return
|
|
}
|
|
publish, cancel := context.WithTimeout(ctx, timeout)
|
|
defer cancel()
|
|
// Not mandatory, unlike a report. Losing one is nothing: the next is a minute away, and the
|
|
// mesh is reading a gap rather than counting arrivals. Insisting on delivery would turn a
|
|
// harmless miss into a logged failure every minute.
|
|
if err := channel.PublishWithContext(publish, Exchange, KeyAlive, false, false,
|
|
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
|
say("could not tell the mesh this node is here: " + err.Error())
|
|
}
|
|
}
|