diff --git a/internal/link/hearing.go b/internal/link/hearing.go new file mode 100644 index 0000000..303d7e7 --- /dev/null +++ b/internal/link/hearing.go @@ -0,0 +1,84 @@ +package link + +import ( + "context" + "time" +) + +// What a host hears, as the host's own words for it. +// +// The outbound half went behind `Bus` (bus.go) and a node's two statements stopped naming a +// transport. This is the other half — dialling, and the declarations that arrive — and it is where +// the transport reached furthest: the run loop selected on a channel of the client library's own +// delivery type, so every part of holding a node in its mesh knew which bus it was on. +// +// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005). This is its own +// interface over its own libraries, and it agrees with the controller only because a conformance +// fixture holds both to one envelope. + +// Link is this node's live connection to its mesh: what it hears, and what it says. +// +// One interface rather than two, because **dialling is where the transport is chosen** and choosing +// it twice is how one half of a node ends up on a different bus from the other. +type Link interface { + // Bus is what this node says: what it applied, and that it is here. + Bus + + // Declarations is what the mesh tells this node to be. + Declarations() <-chan Declaration + + // Lost says the link ended, and why. + // + // **Read rather than discovered.** A node that finds out by noticing silence is a node that + // believed it was in the mesh for as long as the silence lasted, which is the one state ADR + // 0004 says must never look like being connected. + Lost() <-chan error + + // Close lets go of whatever was dialled. + Close() +} + +// Declaration is one thing the mesh told this node to be. +// +// **Handled, once — after the report is published.** A node that dies between applying and +// reporting leaves the declaration with the mesh and applies it again on return, which is safe +// because applying is reconciliation: it converges rather than repeating. +// +// There is one way of being done rather than two. A declaration set aside because a newer arrived +// with it is settled exactly as an applied one is, on both buses, and the difference between them +// is a fact the *report* carries — a second method here would be a distinction the transport does +// not make. +type Declaration interface { + // Body is the signed declaration as it arrived, bytes unchanged: a node verifies what it + // received rather than what it re-encoded. + Body() []byte + + // Handled settles it. Called after the report for it has been published, either way. + Handled() error +} + +// Open opens this node's link to its mesh. +// +// Named Open rather than Dial because Dial is this package's raw TLS dial, which the enrolment path +// uses to see a certificate before it trusts anything. +// +// **Both transports ship and this is the one place that chooses** (novox/hq ADR 0116: nothing moves +// a node's bus before step 5). Until then every membership names the bus the mesh runs on today, +// and the rollout is this switch and the credential behind it — not a change anywhere in the loop +// that reads from what comes back. +func Open(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { + switch m.Transport { + case OnNATS: + return dialNats(ctx, m, timeout) + default: + return dialCurrent(ctx, m, timeout) + } +} + +// The buses a node can be on. Empty is the one the mesh runs on today, which is every node until +// the rollout — so a membership recorded before any of this existed reads as correct rather than as +// unset. +const ( + OnCurrent = "" + OnNATS = "nats" +) diff --git a/internal/link/hearing_current.go b/internal/link/hearing_current.go new file mode 100644 index 0000000..772e17e --- /dev/null +++ b/internal/link/hearing_current.go @@ -0,0 +1,164 @@ +package link + +import ( + "context" + "errors" + "fmt" + "net/url" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The host's link on the bus the mesh runs on today. +// +// Moved out of the run loop rather than changed: the dial, the queue, the prefetch window and the +// return handler are what they were, because the mesh is running on this and a bus nothing speaks +// yet is no reason to alter the one every node is on. + +// currentLink is this node's connection as a channel. +type currentLink struct { + conn *amqp.Connection + channel *amqp.Channel + arrived chan Declaration + lost chan error +} + +func dialCurrent(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { + config, err := PinnedConfig(m.Fingerprint) + if err != nil { + return nil, 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 nil, err + } + return nil, fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err) + } + + channel, err := conn.Channel() + if err != nil { + conn.Close() + return nil, err + } + + queue := QueueFor(m.Node) + if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil { + conn.Close() + return nil, 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 { + conn.Close() + return nil, err + } + + deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil) + if err != nil { + conn.Close() + return nil, err + } + + l := ¤tLink{ + conn: conn, channel: channel, + arrived: make(chan Declaration, drainDepth), + lost: make(chan 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 { + select { + case l.lost <- fmt.Errorf("the broker could not route this node's %s: %s (%d %s)", + r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText): + default: + } + } + }() + + closed := conn.NotifyClose(make(chan *amqp.Error, 1)) + go func() { + select { + case reason := <-closed: + select { + case l.lost <- fmt.Errorf("the link closed: %v", reason): + default: + } + case <-ctx.Done(): + } + }() + + // One goroutine turning the library's deliveries into the mesh's words, so the run loop selects + // on one kind of thing whichever bus it is on. + go func() { + defer close(l.arrived) + for { + select { + case <-ctx.Done(): + return + case delivery, ok := <-deliveries: + if !ok { + select { + case l.lost <- errors.New("the broker stopped delivering"): + default: + } + return + } + select { + case l.arrived <- currentDeclaration{delivery}: + case <-ctx.Done(): + return + } + } + } + }() + + return l, nil +} + +func (l *currentLink) Declarations() <-chan Declaration { return l.arrived } +func (l *currentLink) Lost() <-chan error { return l.lost } + +func (l *currentLink) Close() { + if l.channel != nil { + _ = l.channel.Close() + } + if l.conn != nil { + _ = l.conn.Close() + } +} + +// Report and Alive are the outbound half, over the channel this link holds. +func (l *currentLink) Report(ctx context.Context, node string, body []byte) error { + return OverCurrent{Channel: l.channel}.Report(ctx, node, body) +} + +func (l *currentLink) Alive(ctx context.Context, node string, body []byte) error { + return OverCurrent{Channel: l.channel}.Alive(ctx, node, body) +} + +// currentDeclaration is one delivery from the bus the mesh has. +type currentDeclaration struct{ delivery amqp.Delivery } + +func (d currentDeclaration) Body() []byte { return d.delivery.Body } +func (d currentDeclaration) Handled() error { return d.delivery.Ack(false) } diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go new file mode 100644 index 0000000..e0f90e5 --- /dev/null +++ b/internal/link/hearing_nats.go @@ -0,0 +1,182 @@ +package link + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "github.com/nats-io/nats.go" +) + +// The host's link on the bus being built. +// +// Two things a host does here that it cannot do on the other bus, and one it must not try. +// +// **It declares nothing.** On the bus the mesh has, a host declares its own queue on connecting, +// because a queue that is not there means a node that hears nothing. Here the object it reads +// through is a durable consumer, and a host's account reaches no part of the JetStream API — by +// design, because the controller is the only writer of consumer definitions (design 25 §3). So the +// host **binds** to a consumer the controller made when this node enrolled, and a missing one is +// said as what it is rather than quietly created with whatever configuration this client happens to +// default to. +// +// **It gets order for free, and keeps the drain anyway.** The declaration subject is last-per-subject +// (design 29 §4), so a node that was away receives exactly the current declaration rather than a +// queue of superseded ones, and the stream's sequence orders them definitively — the wire-level +// answer to novox/hq issue 107. What the drain in run.go still answers is the live case: three +// pushes to a *connected* node are three deliveries whatever the stream later retains. + +// DeclareSubject is where this node's declaration lands. Its own, and no other node's: a host's +// account subscribes exactly this and the subject is the authority on which node a declaration is +// for. +func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" } + +// natsURL is a bus address as the client wants it. A membership records host and port, because that +// is what genesis sealed into it and what the other transport takes; the scheme is this transport's +// own business. +func natsURL(address string) string { + if strings.Contains(address, "://") { + return address + } + return "nats://" + address +} + +// natsLink is this node's connection as a JetStream subscription. +type natsLink struct { + conn *nats.Conn + js nats.JetStreamContext + sub *nats.Subscription + node string + arrived chan Declaration + lost chan error +} + +func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { + // Pinned exactly as the other transport is, and for once the Go client makes that easy: it + // takes a *tls.Config, so the same PinnedConfig with the same VerifyPeerCertificate does the + // work. **The constraint recorded against the tool runtime does not apply here** — that client + // takes PEM strings with no verify hook, which is why the bus's certificate must carry a name + // matching the address *modules* dial it by. A host checks the fingerprint and nothing else. + config, err := PinnedConfig(m.Fingerprint) + if err != nil { + return nil, err + } + opts := []nats.Option{ + nats.Secure(config), + nats.UserInfo(m.Node, m.Password), + nats.Name("mesh-host/" + m.Node), + nats.Timeout(timeout), + // A node that has silently lost its route notices, rather than holding a connection the + // server forgot about and believing it is still in the mesh. + nats.PingInterval(10 * time.Second), + nats.MaxPingsOutstanding(2), + // Reconnection is the caller's: Hold already decides when to try again and how long to + // wait, and a client quietly reconnecting underneath it would make that reasoning a + // duplicate of the library's. + nats.NoReconnect(), + } + + conn, err := nats.Connect(natsURL(m.Broker), opts...) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return nil, err + } + return nil, fmt.Errorf("cannot reach the bus at %s: %w", m.Broker, err) + } + js, err := conn.JetStream() + if err != nil { + conn.Close() + return nil, fmt.Errorf("the bus at %s has no JetStream: %w", m.Broker, err) + } + + l := &natsLink{ + conn: conn, js: js, node: m.Node, + arrived: make(chan Declaration, drainDepth), + lost: make(chan error, 1), + } + + // Bound to the consumer the controller made for this node, named after the node because that is + // what the node's own ack grant allows (`$JS.ACK.NODES..>`). + feed := make(chan *nats.Msg, drainDepth) + // The subject as well as the binding: the client checks what is asked for against the + // consumer's own filter, and an empty subject is refused rather than taken to mean "whatever + // that consumer delivers". + sub, err := js.ChanSubscribe(DeclareSubject(m.Node), feed, nats.Bind("NODES", m.Node)) + if err != nil { + conn.Close() + return nil, fmt.Errorf( + "this node cannot read its declarations: %w. The mesh creates that when a node enrols, "+ + "and a host may not create one itself — so this is the mesh's to answer, not this "+ + "machine's", err) + } + l.sub = sub + + conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) { + select { + case l.lost <- fmt.Errorf("the link dropped: %w", err): + default: + } + }) + conn.SetClosedHandler(func(*nats.Conn) { + select { + case l.lost <- errors.New("the link closed"): + default: + } + }) + + go func() { + defer close(l.arrived) + for { + select { + case <-ctx.Done(): + return + case msg, ok := <-feed: + if !ok { + select { + case l.lost <- errors.New("the bus stopped delivering"): + default: + } + return + } + select { + case l.arrived <- natsDeclaration{msg}: + case <-ctx.Done(): + return + } + } + } + }() + + return l, nil +} + +func (l *natsLink) Declarations() <-chan Declaration { return l.arrived } +func (l *natsLink) Lost() <-chan error { return l.lost } + +func (l *natsLink) Close() { + if l.sub != nil { + _ = l.sub.Unsubscribe() + } + if l.conn != nil { + l.conn.Close() + } +} + +func (l *natsLink) Report(ctx context.Context, node string, body []byte) error { + return OverNATS{Conn: l.conn, JS: l.js}.Report(ctx, node, body) +} + +func (l *natsLink) Alive(ctx context.Context, node string, body []byte) error { + return OverNATS{Conn: l.conn, JS: l.js}.Alive(ctx, node, body) +} + +// natsDeclaration is one declaration off the NODES stream. +type natsDeclaration struct{ msg *nats.Msg } + +func (d natsDeclaration) Body() []byte { return d.msg.Data } + +// Handled acknowledges it. The ack goes to this node's own ack subject, which is the one thing +// besides its reports a node's account may publish. +func (d natsDeclaration) Handled() error { return d.msg.Ack() } diff --git a/internal/link/hearing_nats_test.go b/internal/link/hearing_nats_test.go new file mode 100644 index 0000000..9a33e38 --- /dev/null +++ b/internal/link/hearing_nats_test.go @@ -0,0 +1,228 @@ +package link + +import ( + "context" + "crypto/ed25519" + "encoding/json" + "fmt" + "os" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// The host's link against a real server, because every claim here is about one. +// +// Whether binding to a consumer the host did not create works, whether a declaration on the node's +// own subject arrives, whether acknowledging it removes it from the consumer's pending — none of +// that can be reasoned out, and the first two are the ones that would leave a node silently hearing +// nothing: +// +// docker run -d --rm --name t -p 14223:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14223 go test ./internal/link/ -run TestNats + +func aBus(t *testing.T) (*nats.Conn, nats.JetStreamContext) { + t.Helper() + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + conn, err := nats.Connect(url) + if err != nil { + t.Fatal(err) + } + t.Cleanup(conn.Close) + js, err := conn.JetStream() + if err != nil { + t.Fatal(err) + } + // **Ensured and purged, not deleted and recreated.** Delete-then-add looked like a reset and is + // not one: a test that did that inherited the previous test's messages, and the symptom was a + // declaration counted as delivered twice — which reads as a redelivery bug in the code under + // test rather than as a dirty stream. Purge is defined to empty a stream; recreating one is a + // race with the server's own teardown. + for _, want := range []*nats.StreamConfig{ + {Name: "NODES", Subjects: []string{"mesh.node.*.declare"}, MaxMsgsPerSubject: 1}, + {Name: "CONTROL", Subjects: []string{"mesh.control.*.report", "mesh.control.enrol"}, + Retention: nats.WorkQueuePolicy}, + } { + if _, err := js.StreamInfo(want.Name); err != nil { + if _, err := js.AddStream(want); err != nil { + t.Fatal(err) + } + } + if err := js.PurgeStream(want.Name); err != nil { + t.Fatal(err) + } + } + return conn, js +} + +// theMeshMakes is the consumer the controller creates when a node enrols. Made here by the test +// because the host may not: its account reaches no part of the JetStream API, which is the whole +// reason this binds rather than subscribes. +// +// Removed afterwards, and each test names its own node: two tests sharing a consumer name share its +// delivery count and its pending list, and the first thing that goes wrong reads as a fault in the +// host rather than in the test beside it. +func theMeshMakes(t *testing.T, js nats.JetStreamContext, node string) { + t.Helper() + t.Cleanup(func() { _ = js.DeleteConsumer("NODES", node) }) + if _, err := js.AddConsumer("NODES", &nats.ConsumerConfig{ + Durable: node, + FilterSubject: DeclareSubject(node), + AckPolicy: nats.AckExplicitPolicy, + AckWait: 300 * time.Second, + DeliverSubject: "_DELIVER." + node, + }); err != nil { + t.Fatal(err) + } +} + +func signedBy(t *testing.T, key ed25519.PrivateKey, declaration []byte) []byte { + t.Helper() + body, err := json.Marshal(Signed{ + Declaration: declaration, Signature: ed25519.Sign(key, declaration), + }) + if err != nil { + t.Fatal(err) + } + return body +} + +// A declaration on this node's own subject reaches the host, is applied, and acknowledging it +// empties the consumer — which is what tells the mesh the node has it. +func TestNatsADeclarationReachesTheHostAndIsSettled(t *testing.T) { + conn, js := aBus(t) + const node = "settling" + theMeshMakes(t, js, node) + + public, private, _ := ed25519.GenerateKey(nil) + m := Membership{Node: node, Signer: public} + + // Dialled directly rather than through Open: the test server has no TLS, and what is being + // checked is the subscription and the settling, not the pin — which PinnedConfig owns and its + // own tests cover. + l := &natsLink{conn: conn, js: js, node: node, + arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1)} + feed := make(chan *nats.Msg, drainDepth) + sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node)) + if err != nil { + t.Fatalf("the host could not bind to the consumer the mesh made for it: %v", err) + } + defer func() { _ = sub.Unsubscribe() }() + go func() { + for msg := range feed { + l.arrived <- natsDeclaration{msg} + } + }() + + if _, err := js.Publish(DeclareSubject(node), + signedBy(t, private, []byte(`{"declared":"d1"}`))); err != nil { + t.Fatal(err) + } + + select { + case d := <-l.Declarations(): + report := handleBody(context.Background(), m, d.Body(), + func(context.Context, []byte, []byte) Report { + return Report{Applied: []string{"store"}} + }) + if report.Refused != "" { + t.Fatalf("a declaration the mesh signed was refused: %s", report.Refused) + } + if err := d.Handled(); err != nil { + t.Fatalf("the node could not acknowledge its own declaration: %v", err) + } + case <-time.After(8 * time.Second): + t.Fatal("no declaration reached the host") + } + + // **Nothing pending is the property**; a delivery count is not. Delivery is at-least-once by + // design, so pinning "delivered exactly once" would be asserting something the mesh does not + // rely on. What matters is that the acknowledgement landed, so the mesh can tell the node has + // it — and that no redelivery was needed to get there, which is what would say the node was + // too slow to answer for its own ack wait. + deadline := time.Now().Add(5 * time.Second) + var last string + for time.Now().Before(deadline) { + info, err := js.ConsumerInfo("NODES", node) + switch { + case err != nil: + last = err.Error() + case info.NumAckPending == 0 && info.NumRedelivered == 0: + return + default: + last = fmt.Sprintf("pending %d, redelivered %d", info.NumAckPending, info.NumRedelivered) + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("the declaration was not settled, so the mesh cannot tell the node has it: %s", last) +} + +// **A node that was away gets exactly the current declaration and nothing older.** Three pushed +// while nothing is listening leave one on the stream, and it is the newest — the wire-level answer +// to novox/hq issue 107, and the half of the drain that stops being the host's problem. +func TestNatsANodeThatWasAwayGetsOnlyTheNewest(t *testing.T) { + _, js := aBus(t) + const node = "returning" + _, private, _ := ed25519.GenerateKey(nil) + + for _, id := range []string{"d1", "d2", "d3"} { + if _, err := js.Publish(DeclareSubject(node), + signedBy(t, private, []byte(`{"declared":"`+id+`"}`))); err != nil { + t.Fatal(err) + } + } + info, err := js.StreamInfo("NODES") + if err != nil { + t.Fatal(err) + } + if info.State.Msgs != 1 { + t.Fatalf("%d declarations survived for one node; a node that was away would apply a backlog "+ + "of things nobody wants any more", info.State.Msgs) + } + + theMeshMakes(t, js, node) + feed := make(chan *nats.Msg, drainDepth) + sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node)) + if err != nil { + t.Fatal(err) + } + defer func() { _ = sub.Unsubscribe() }() + + select { + case msg := <-feed: + if declaredIn(msg.Data) != "d3" { + t.Fatalf("the node was given %q rather than the newest", declaredIn(msg.Data)) + } + case <-time.After(8 * time.Second): + t.Fatal("the node that was away was given nothing") + } +} + +// A report goes through the stream and a heartbeat does not: the one that must survive the +// controller's store restarting is kept, and the one that must not is not. +func TestNatsAReportIsKeptAndAHeartbeatIsNot(t *testing.T) { + conn, js := aBus(t) + bus := OverNATS{Conn: conn, JS: js} + ctx := context.Background() + + body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1"}) + if err := bus.Report(ctx, "anchor", body); err != nil { + t.Fatal(err) + } + beat, _ := json.Marshal(Alive{Node: "anchor"}) + if err := bus.Alive(ctx, "anchor", beat); err != nil { + t.Fatal(err) + } + + info, err := js.StreamInfo("CONTROL") + if err != nil { + t.Fatal(err) + } + if info.State.Msgs != 1 { + t.Fatalf("%d messages were kept; a report must be and a heartbeat must not", info.State.Msgs) + } +} diff --git a/internal/link/newest_test.go b/internal/link/newest_test.go index 7ad6cbb..5eab4a0 100644 --- a/internal/link/newest_test.go +++ b/internal/link/newest_test.go @@ -3,33 +3,46 @@ package link import ( "testing" "time" - - amqp "github.com/rabbitmq/amqp091-go" ) +// said is one declaration as a test hands it over, with no transport under it — which is what the +// seam bought: the drain's reasoning was reachable only through a real broker before. +type said struct { + body []byte + handled bool +} + +func (s *said) Body() []byte { return s.body } +func (s *said) Handled() error { s.handled = true; return nil } + +func arriving(bodies ...string) chan Declaration { + ch := make(chan Declaration, 8) + for _, b := range bodies { + ch <- &said{body: []byte(b)} + } + return ch +} + // A machine asked to be five things becomes the last one: what is already waiting supersedes what // arrived first, and everything set aside is named so it can be reported. func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) { - deliveries := make(chan amqp.Delivery, 8) - for _, id := range []string{"two", "three", "four"} { - deliveries <- amqp.Delivery{Body: []byte(id)} + waiting := arriving("two", "three", "four") + apply, superseded := newest(waiting, &said{body: []byte("one")}, 50*time.Millisecond) + if string(apply.Body()) != "four" { + t.Fatalf("applied %q, not the newest", apply.Body()) } - apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("one")}, 50*time.Millisecond) - if string(apply.Body) != "four" { - t.Fatalf("applied %q, not the newest", apply.Body) - } - if len(superseded) != 3 || string(superseded[0].Body) != "one" || string(superseded[2].Body) != "three" { + if len(superseded) != 3 || string(superseded[0].Body()) != "one" || + string(superseded[2].Body()) != "three" { t.Fatalf("set aside %d: %v", len(superseded), superseded) } } // One declaration with nothing behind it is applied as it always was, after the window. func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) { - deliveries := make(chan amqp.Delivery, 1) began := time.Now() - apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("only")}, 30*time.Millisecond) - if string(apply.Body) != "only" || len(superseded) != 0 { - t.Fatalf("got %q with %d set aside", apply.Body, len(superseded)) + apply, superseded := newest(arriving(), &said{body: []byte("only")}, 30*time.Millisecond) + if string(apply.Body()) != "only" || len(superseded) != 0 { + t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded)) } if time.Since(began) < 30*time.Millisecond { t.Fatal("did not wait the window for a straggler") @@ -38,13 +51,13 @@ func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) { // A straggler within the window is taken; one after it is the next push. func TestAStragglerWithinTheWindowIsTaken(t *testing.T) { - deliveries := make(chan amqp.Delivery, 2) + waiting := arriving() go func() { time.Sleep(20 * time.Millisecond) - deliveries <- amqp.Delivery{Body: []byte("late")} + waiting <- &said{body: []byte("late")} }() - apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("first")}, 100*time.Millisecond) - if string(apply.Body) != "late" || len(superseded) != 1 { - t.Fatalf("got %q with %d set aside", apply.Body, len(superseded)) + apply, superseded := newest(waiting, &said{body: []byte("first")}, 100*time.Millisecond) + if string(apply.Body()) != "late" || len(superseded) != 1 { + t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded)) } } diff --git a/internal/link/run.go b/internal/link/run.go index a5687cb..0288a4f 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -6,10 +6,7 @@ import ( "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. @@ -33,6 +30,10 @@ type Membership struct { Fingerprint string Password string Signer ed25519.PublicKey + // Transport is which bus this node speaks (hearing.go). Empty is the one the mesh runs on + // today, which is every node until the rollout — so a membership recorded before any of this + // existed reads as correct rather than as unset. + Transport string } // Applier is what the host does with a declaration that has been proved to come from the mesh. @@ -190,108 +191,62 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout if say == nil { say = func(string) {} } - config, err := PinnedConfig(m.Fingerprint) + link, err := Open(ctx, m, timeout) if err != nil { return err } + defer link.Close() - 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) + // 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, hearing what this node should be") // 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, OverCurrent{Channel: channel}, m, say, timeout) + publishAlive(ctx, link, 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)) - } - }() + declarations := link.Declarations() for { select { case <-ctx.Done(): return nil case <-beat.C: - publishAlive(ctx, OverCurrent{Channel: channel}, m, say, timeout) + publishAlive(ctx, link, 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, OverCurrent{Channel: channel}, m, unasked.Report, say, timeout) + published := publishReport(ctx, link, 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: + case reason := <-link.Lost(): + return reason + case declaration, ok := <-declarations: if !ok { - return errors.New("the broker stopped delivering") + // The link's own reason, when it has managed to say one: "stopped delivering" on + // its own says nothing about why, and why is the whole of what an operator wants. + select { + case reason := <-link.Lost(): + return reason + default: + return errors.New("the mesh stopped sending this node declarations") + } } - // 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) + // Whatever else is already waiting supersedes this one. Each set-aside declaration is + // reported as such, then settled unapplied. + declaration, superseded := newest(declarations, declaration, drainWindow) for _, old := range superseded { say("set aside a declaration: a newer one arrived with it") - publishReport(ctx, OverCurrent{Channel: channel}, m, Report{Node: m.Node, Declared: declaredIn(old.Body), - Superseded: declaredIn(delivery.Body)}, say, timeout) - _ = old.Ack(false) + publishReport(ctx, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()), + Superseded: declaredIn(declaration.Body())}, say, timeout) + _ = old.Handled() } - report := handle(ctx, m, apply, delivery) + report := handleBody(ctx, m, declaration.Body(), apply) switch { case report.Refused != "": say("refused a declaration: " + report.Refused) @@ -300,12 +255,12 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout default: say(fmt.Sprintf("applied %d resource(s)", len(report.Applied))) } - publishReport(ctx, OverCurrent{Channel: 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, + publishReport(ctx, link, m, report, say, timeout) + // Settled after the report is published. A node that dies between applying and + // reporting leaves the declaration with the mesh and applies it again on return, // which is safe because applying is reconciliation — it converges rather than // repeating. - _ = delivery.Ack(false) + _ = declaration.Handled() } } } @@ -333,12 +288,12 @@ const ( // three deliveries, whatever the stream later retains. So this is narrowed at the rollout, not // deleted — and saying which half goes is worth more than a note that it "can probably be // removed", which is how a load-bearing window gets deleted by somebody in a hurry. -func newest(deliveries <-chan amqp.Delivery, first amqp.Delivery, window time.Duration) (amqp.Delivery, []amqp.Delivery) { +func newest(arriving <-chan Declaration, first Declaration, window time.Duration) (Declaration, []Declaration) { latest := first - var superseded []amqp.Delivery + var superseded []Declaration for { select { - case next, ok := <-deliveries: + case next, ok := <-arriving: if !ok { return latest, superseded } @@ -367,10 +322,6 @@ func declaredIn(body []byte) string { 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 { @@ -393,30 +344,13 @@ func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) R // 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) + link, err := Open(ctx, m, timeout) 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() + defer link.Close() var said string - if !publishReport(ctx, OverCurrent{Channel: channel}, m, report, func(s string) { said = s }, timeout) { + if !publishReport(ctx, link, m, report, func(s string) { said = s }, timeout) { return errors.New(said) } return nil