diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index fb736b1..00eedf6 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -782,6 +782,25 @@ func (i *Inventory) RecordSent(ctx context.Context, node, digest string) error { return err } +// Outstanding is the digest of the declaration a machine was last sent, by its name, and empty +// for one that has never been sent anything. +// +// **By name rather than by id**, because the caller is the serving loop and what a node puts in a +// report is its name. Asked of one machine rather than read from Waiting's sweep, because it is +// asked per message: a report names the declaration it is about, and a report about one the mesh +// has already moved past is not acted on (design 25 §3). +func (i *Inventory) Outstanding(ctx context.Context, name string) (string, error) { + var sent string + err := i.store.Pool().QueryRow(ctx, + `select coalesce(sent, '') from node where name = $1`, name).Scan(&sent) + if errors.Is(err, pgx.ErrNoRows) { + // Not an error worth carrying up: a report from a machine the mesh has no record of has + // nothing to be stale against, and whatever is wrong with it is the listener's to say. + return "", nil + } + return sent, err +} + // Waiting is every machine whose declaration has changed since it was last sent one. // // The caller works out what each machine should be now, because only it can — resolution is the diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index 54b8967..657ec21 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -234,6 +234,17 @@ var ErrNoBrokerManagement = errors.New("no broker management configured") // A node states; the owning context writes (novox/hq ADR 0006). What a node says it applied is // its own account of its own machine, kept as a copy for recovery — so this writes it down and // decides nothing from it. +// Outstanding is the declaration the mesh last sent a node, so a report about an older one is not +// acted on (design 25 §3, window.go). +// +// **Here rather than on Listener.** A report is recorded by whatever keeps records, and a great +// many things that record reports have no idea what was sent — every test in this package among +// them. So the serving loop asks for this when the listener happens to be able to answer, and +// where it cannot, a report has nothing to be stale against and is simply acted on. +func (e Enrolment) Outstanding(ctx context.Context, node string) (string, error) { + return e.Inventory.Outstanding(ctx, node) +} + func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { // A store that could not be asked right now is said as such, so the report is kept for // another attempt rather than acknowledged and lost (novox/hq issue 082). diff --git a/internal/link/receive.go b/internal/link/receive.go new file mode 100644 index 0000000..eda491c --- /dev/null +++ b/internal/link/receive.go @@ -0,0 +1,125 @@ +package link + +import ( + "context" + "time" +) + +// The consume side of the bus, in the mesh's own words. +// +// The outbound half went behind `Bus` (bus.go) and the transport stopped reaching its callers. +// This is the other half, and it is the larger one: everything a node or a module says arrives +// here, and until now every handler took the transport's own delivery type — so the serving loop +// could not be moved to another bus without moving enrolment, reports, builds, upgrades and +// catch-up with it in one breath. +// +// **Two implementations, both shipping** (novox/hq ADR 0116: nothing moves a node's bus before +// step 5). Both shipping is what makes them comparable, and it is what lets the store-window +// guarantee (ADR 0083) be stated once — in window.go, pure — rather than twice, once per +// transport, where the two would eventually disagree about the thing that matters most. + +// The kinds of message the controller acts on. +// +// A transport maps its own addressing onto these — a routing key on the bus the mesh has, a +// subject on the one being built — and nothing past this point knows which it was. They are never +// on the wire: the wire is the transport's business, and a kind that travelled would be a third +// name for the same thing. +const ( + KindEnrolment = "enrolment" + KindReport = "report" + KindHeartbeat = "heartbeat" + KindBuilt = "built" + KindModuleMoved = "module-moved" + KindCatchUp = "catch-up" +) + +// Control is one thing a node or a module said, as the controller must act on it. +// +// **Settling is stated as what the mesh means, not as the transport's verbs.** The two buses +// spell them differently — an ack and a reject against a delivery tag, an ack and a term against +// a stream sequence — and the guarantee is the same either way: `Took` is done with, `Drop` is +// understood and not worth another attempt, and `Hold` is the store window, where the message is +// kept and comes back. +// +// A handler that returns without calling any of the three leaves the message unsettled on +// purpose. That is the right answer while shutting down: a cancelled context is not an answer +// about a message, and the bus should hand it to whatever consumes next (novox/hq issue 083). +type Control interface { + // Kind is which of the constants above this is. + Kind() string + + // Body is the message itself — the payload alone, never the envelope. + Body() []byte + + // Redelivered says the bus has handed this message over before. An enrolment cares and + // nothing else does: one already spent is not finished a second time. + Redelivered() bool + + // HeldFor is how long this message has been waiting to be taken. Zero on a first delivery. + // + // **Read from the message rather than remembered by the controller.** On the bus being built + // it is the age of the publish, which a controller that restarted mid-window still reads + // correctly — the whole reason the holding moves into the server. On the bus the mesh has it + // is how long this process has held it, which is the most that transport can say. + HeldFor() time.Duration + + // Answer replies to whoever is waiting on this message; only an enrolment expects one. + // + // Each transport knows where its own answer goes, and they do not agree about it: one carries + // a reply queue in the delivery, and on the other the field that would have carried it has + // been claimed by the consumer's own ack subject, so the address travels in the payload + // (design 25 §2, verified). That difference is exactly what this seam exists to keep out of + // the handler. + Answer(ctx context.Context, body []byte) error + + // Took settles the message: acted on, or understood and needing no action. + Took() error + + // About names what this message is about — a node's report, one module's move, one build's + // outcome — and is said before the store is asked. + // + // A transport that holds messages **in memory** uses it to set aside anything older it is + // holding about the same thing: the older is the past, and letting it come back after the + // newer was acted on would undo the newer. + // + // **This is the one thing holding-in-memory can do that holding-in-the-server cannot**, and + // naming it here rather than hiding it is deliberate. On the bus being built the message + // belongs to the server and comes back whatever happened meanwhile, so this is ignored and the + // digest a report carries answers the same question instead (window.go, design 25 §3). + About(what string) + + // Hold keeps the message and asks for it again after the delay — the store window. + Hold(after time.Duration) error + + // Drop settles the message without acting on it: refused, stale, or given up on. It is not + // delivered again. + Drop() error +} + +// Inbound is where control messages come from. +type Inbound interface { + // Also asks for one more kind to be delivered. + // + // **Nothing is subscribed unless something is listening for it.** A durable queue or a + // durable stream consumer that nobody reads fills quietly, and the first symptom is a bus out + // of disk rather than anything about modules. + Also(kind string) error + + // Receive delivers every message to act until the context ends, and says why it stopped. + Receive(ctx context.Context, act func(context.Context, Control)) error + + // Close lets go of whatever the implementation holds. + Close() +} + +// Outstanding answers which declaration the mesh last sent a node — the digest, not the +// declaration. +// +// **Asked before the store is waited on** (design 25 §3): a report about a declaration the mesh +// has already moved past is not worth holding a slot in the window that a current message needs. +// It is a separate interface from Listener rather than a method on it, because a controller that +// only publishes needs neither and something that records reports need not also be able to say +// what was sent. +type Outstanding interface { + Outstanding(ctx context.Context, node string) (string, error) +} diff --git a/internal/link/receive_current.go b/internal/link/receive_current.go new file mode 100644 index 0000000..efeeae5 --- /dev/null +++ b/internal/link/receive_current.go @@ -0,0 +1,317 @@ +package link + +import ( + "context" + "errors" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The consume side on the bus the mesh runs on today. +// +// Everything here was the serving loop's until the seam went in: the queues, the binds, the +// prefetch, and the list of messages the store could not take yet. It moved rather than changed — +// the behaviour this transport has is the behaviour it had, because the mesh is running on it and +// a bus nothing speaks yet is no reason to alter the one every node is on (ADR 0116). + +// Prefetch is how many messages the bus hands the controller before it has settled them. +// +// More than one because a message the store could not take is held, unsettled, while the loop +// goes on answering others — an enrolment above all, which a host is waiting on (novox/hq issue +// 083). Bounded, because what is held is also what the bus has not kept on its own disk as +// pending. +const Prefetch = 64 + +// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to +// answer — an enrolment above all — while others wait for the store. +const PrefetchHeadroom = 8 + +// TryAgainAfter is how often the held are looked at. A store comes back in seconds, and a report a +// few seconds late is still current. +const TryAgainAfter = 2 * time.Second + +// currentInbound consumes what nodes say over the bus the mesh has. +type currentInbound struct { + conn *amqp.Connection + channel *amqp.Channel + // upgrades and catchups are bound only when something is listening (Also). + upgrades bool + catchups bool + // held is every message the store could not take, by the bus's own delivery tag. Kept here + // rather than in the serving loop because holding a delivery unacknowledged is this + // transport's way of keeping it, and the other's is to hand it back to the server. + held map[uint64]*holding + // again is how often the held are looked at; zero means TryAgainAfter. Set by tests. + again time.Duration +} + +// holding is one message kept for the store, and when to try it again. +type holding struct { + message *currentControl + due time.Time + about string +} + +// Current is the consume side of the bus the mesh runs on today. +func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound { + return ¤tInbound{conn: conn, channel: channel, held: map[uint64]*holding{}} +} + +// Also binds the queue one more kind arrives on. +// +// The kinds nodes publish all share one queue and are bound at Connect, because a node may +// publish any of them and binding one while forgetting another is a message the bus accepts, finds +// no queue for, and drops — the publisher sees success and the consumer sees nothing. The two that +// are events get their own queue each, and only when something is listening. +func (c *currentInbound) Also(kind string) error { + switch kind { + case KindModuleMoved: + if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil { + return err + } + c.upgrades = true + case KindCatchUp: + if err := c.bindEvent(CatchUpQueue, KeyCatchingUp); err != nil { + return err + } + c.catchups = true + default: + return fmt.Errorf("nothing binds a queue for %s on this bus", kind) + } + return nil +} + +func (c *currentInbound) bindEvent(queue, key string) error { + if _, err := c.channel.QueueDeclare(queue, true, false, false, false, nil); err != nil { + return fmt.Errorf("cannot declare the %s queue: %w", queue, err) + } + if err := c.channel.QueueBind(queue, key, EventsExchange, false, nil); err != nil { + return fmt.Errorf("cannot bind %s to %s/%s: %w", queue, EventsExchange, key, err) + } + return nil +} + +func (c *currentInbound) Close() {} + +// Receive consumes until the context ends. +// +// One consumer per queue, deliberately: with two on one queue the bus would round-robin between +// them and each would receive half of what it expects — a fault this project has already had, +// between a module's daemon and its capability server. +func (c *currentInbound) Receive(ctx context.Context, act func(context.Context, Control)) error { + // A bounded prefetch rather than one. The loop still takes messages one at a time; what the + // prefetch buys is that a message the store could not take can be held while the loop goes on + // to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything + // held goes back to the bus if the controller stops, because nothing held is acknowledged. + if err := c.channel.Qos(Prefetch, 0, false); err != nil { + return err + } + + deliveries, err := c.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane", + false, false, false, false, nil) + if err != nil { + return err + } + + // Its own queue and its own consumer for each event, for the reason above: two consumers on + // one queue split its messages between them, and an upgrade or a catch-up request going to + // whichever half was not listening is a gap that looks like a working mesh. + var upgrades, catchups <-chan amqp.Delivery + if c.upgrades { + upgrades, err = c.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades", + false, false, false, false, nil) + if err != nil { + return err + } + } + if c.catchups { + catchups, err = c.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup", + false, false, false, false, nil) + if err != nil { + return err + } + } + + closed := c.conn.NotifyClose(make(chan *amqp.Error, 1)) + + again := c.again + if again == 0 { + again = TryAgainAfter + } + ticker := time.NewTicker(again) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return nil + case <-ticker.C: + if ctx.Err() != nil { + return nil + } + c.retryHeld(ctx, act) + case delivery, ok := <-catchups: + if !ok { + if catchups != nil { + return errors.New("the bus stopped delivering catch-up requests") + } + continue + } + act(ctx, c.wrap(KindCatchUp, delivery)) + case delivery, ok := <-upgrades: + // A nil channel blocks for ever, so this case simply never fires when nothing is + // listening for upgrades. Closed is different, and means the bus stopped. + if !ok { + if upgrades != nil { + return errors.New("the bus stopped delivering upgrades") + } + continue + } + act(ctx, c.wrap(KindModuleMoved, delivery)) + case reason := <-closed: + // Said rather than returned quietly. A controller whose bus connection dropped is a + // mesh where nothing can be told anything, and the reason is the first thing anybody + // will want. + return fmt.Errorf("the bus connection closed: %v", reason) + case delivery, ok := <-deliveries: + if !ok { + return errors.New("the bus stopped delivering") + } + kind, known := kindOfKey[delivery.RoutingKey] + if !known { + // Rejected without requeue: a message nothing understands will not be understood + // on the next attempt either, and requeuing it would spin. + _ = delivery.Reject(false) + continue + } + act(ctx, c.wrap(kind, delivery)) + } + } +} + +// kindOfKey is how this transport's addressing becomes what the mesh calls a message. +var kindOfKey = map[string]string{ + KeyEnrol: KindEnrolment, + KeyReport: KindReport, + KeyAlive: KindHeartbeat, + KeyBuilt: KindBuilt, + KeyModuleUpgraded: KindModuleMoved, + KeyCatchingUp: KindCatchUp, +} + +func (c *currentInbound) wrap(kind string, delivery amqp.Delivery) *currentControl { + return ¤tControl{kind: kind, delivery: delivery, on: c} +} + +// retryHeld hands every message whose delay has passed back to the loop. Each handler holds it +// again, settles it, or lets it go past the bound. +func (c *currentInbound) retryHeld(ctx context.Context, act func(context.Context, Control)) { + now := time.Now() + due := make([]*currentControl, 0, len(c.held)) + for _, h := range c.held { + if !h.due.After(now) { + due = append(due, h.message) + } + } + for _, m := range due { + if ctx.Err() != nil { + return + } + act(ctx, m) + } +} + +// currentControl is one delivery from the bus the mesh has, as the controller reads it. +type currentControl struct { + kind string + delivery amqp.Delivery + on *currentInbound + // about is what this message is about, as the handler named it; empty until it does. + about string + // first is when this message was first held for the store; zero while it has not been. + first time.Time +} + +func (m *currentControl) Kind() string { return m.kind } +func (m *currentControl) Body() []byte { return m.delivery.Body } +func (m *currentControl) Redelivered() bool { return m.delivery.Redelivered } + +func (m *currentControl) HeldFor() time.Duration { + if m.first.IsZero() { + return 0 + } + return time.Since(m.first) +} + +// Answer publishes to the reply queue the request named. +func (m *currentControl) Answer(ctx context.Context, body []byte) error { + if m.delivery.ReplyTo == "" { + return errors.New("that request named no reply queue, so nothing can be told the answer") + } + return m.on.channel.PublishWithContext(ctx, "", m.delivery.ReplyTo, false, false, + amqp.Publishing{ + ContentType: "application/json", + CorrelationId: m.delivery.CorrelationId, + Body: body, + }) +} + +func (m *currentControl) Took() error { + m.forget() + return m.delivery.Ack(false) +} + +// Drop rejects without requeue: on this bus that is what "understood, and not worth another +// attempt" is spelled as, and it is what feeds a dead-letter queue where one is configured. +func (m *currentControl) Drop() error { + m.forget() + return m.delivery.Reject(false) +} + +// About names what this message is about, and lets go of whatever is held about the same thing: +// the held one is the past, and acting on it after this one would undo this one. Acknowledged +// rather than left to come back, because a held message nothing will act on is a place in the +// prefetch nothing gets back. +func (m *currentControl) About(what string) { + m.about = what + if what == "" { + return + } + for tag, h := range m.on.held { + if h.about != what || tag == m.delivery.DeliveryTag { + continue + } + delete(m.on.held, tag) + _ = h.message.delivery.Ack(false) + } +} + +// Hold keeps the message unacknowledged and sets it aside to be handed back after the delay. +// +// Held no further than the prefetch leaves room: past that the bus would hand the loop nothing +// new — enrolments included — until something held was let go. A message that cannot be held says +// so, and the handler settles it its own way. +func (m *currentControl) Hold(after time.Duration) error { + if m.on.held == nil { + m.on.held = map[uint64]*holding{} + } + if _, already := m.on.held[m.delivery.DeliveryTag]; !already { + if len(m.on.held) >= Prefetch-PrefetchHeadroom { + return fmt.Errorf("%d messages are already held for the store, and holding more "+ + "would stop the queue", len(m.on.held)) + } + m.first = time.Now() + } + m.on.held[m.delivery.DeliveryTag] = &holding{ + message: m, due: time.Now().Add(after), about: m.about, + } + return nil +} + +func (m *currentControl) forget() { + if m.on != nil { + delete(m.on.held, m.delivery.DeliveryTag) + } +} diff --git a/internal/link/receive_current_test.go b/internal/link/receive_current_test.go new file mode 100644 index 0000000..462e3f3 --- /dev/null +++ b/internal/link/receive_current_test.go @@ -0,0 +1,71 @@ +package link + +import ( + "context" + "encoding/json" + "io" + "log" + "testing" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The harness for the consume side on the bus the mesh runs on today. +// +// Messages arrive through the seam, so what these tests exercise is the controller's decision +// about a message and this transport's way of keeping one — which is what the seam separated. A +// fake acknowledger stands in for the bus, because what is asserted is how a message was settled +// and that needs no server. + +// settled is how the bus was told to settle one message. +type settled struct{ acked, nacked, requeued, rejected bool } + +func (a *settled) Ack(uint64, bool) error { a.acked = true; return nil } +func (a *settled) Nack(_ uint64, _ bool, requeue bool) error { + a.nacked, a.requeued = true, requeue + return nil +} +func (a *settled) Reject(uint64, bool) error { a.rejected = true; return nil } + +// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will +// hand it to whatever consumes next if the controller stops. +func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected } + +var tag uint64 + +func quiet() *log.Logger { return log.New(io.Discard, "", 0) } + +// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder, +// an upgrader or a replayer to be set on it. +func serving() (*Server, *currentInbound) { + in := ¤tInbound{held: map[uint64]*holding{}} + return &Server{inbound: in, bus: OverCurrent{}, log: quiet()}, in +} + +// sends is one message arriving over this transport, as the controller reads it. +func (c *currentInbound) sends(t *testing.T, to *settled, kind string, v any) Control { + t.Helper() + body, err := json.Marshal(v) + if err != nil { + t.Fatal(err) + } + tag++ + return ¤tControl{kind: kind, on: c, delivery: amqp.Delivery{ + Acknowledger: to, Body: body, DeliveryTag: tag, + }} +} + +// dueNow brings every held message forward, so a test need not wait out the backoff a real store +// restart would be given (RedeliverAfter). +func (c *currentInbound) dueNow() { + for _, h := range c.held { + h.due = time.Now().Add(-time.Second) + } +} + +// retries hands every held message back to the controller, the way the ticker does. +func (c *currentInbound) retries(ctx context.Context, s *Server) { + c.dueNow() + c.retryHeld(ctx, s.act) +} diff --git a/internal/link/report_retry_test.go b/internal/link/report_retry_test.go index 8a175d7..7549437 100644 --- a/internal/link/report_retry_test.go +++ b/internal/link/report_retry_test.go @@ -2,26 +2,11 @@ package link import ( "context" - "encoding/json" "errors" - "io" - "log" "testing" "time" - - amqp "github.com/rabbitmq/amqp091-go" ) -type saidTo struct{ acked, nacked, requeued bool } - -func (a *saidTo) Ack(uint64, bool) error { a.acked = true; return nil } -func (a *saidTo) Nack(_ uint64, _ bool, requeue bool) error { - a.nacked, a.requeued = true, requeue - return nil -} -func (a *saidTo) Reject(uint64, bool) error { return nil } -func (a *saidTo) unsettled() bool { return !a.acked && !a.nacked } - type heardWith struct{ err error } func (h heardWith) Heard(context.Context, Report) error { return h.err } @@ -31,39 +16,31 @@ type switchable struct{ err error } func (h *switchable) Heard(context.Context, Report) error { return h.err } -var tag uint64 - -func aReport(t *testing.T, to *saidTo, node, declared string) amqp.Delivery { - t.Helper() - body, err := json.Marshal(Report{Node: node, Declared: declared, Applied: []string{"store"}}) - if err != nil { - t.Fatal(err) - } - tag++ - return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body, DeliveryTag: tag} +func aReport(node, declared string) Report { + return Report{Node: node, Declared: declared, Applied: []string{"store"}} } -func quiet() *log.Logger { return log.New(io.Discard, "", 0) } - // A report the store could not take right now is held, unsettled, and recorded when the store is // back; one the store answered no to is acknowledged; one recorded is acknowledged (issue 082, 083). func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) { store := &switchable{err: errors.Join(ErrTryAgain, errors.New("starting up"))} - s := &Server{listener: store, log: quiet()} - held := &saidTo{} - s.handleReport(context.Background(), aReport(t, held, "anchor", "d1")) - if !held.unsettled() || len(s.parked) != 1 { - t.Fatalf("a report the store could not take was not held: %+v, %d held", held, len(s.parked)) + s, in := serving() + s.listener = store + held := &settled{} + s.act(context.Background(), in.sends(t, held, KindReport, aReport("anchor", "d1"))) + if !held.unsettled() || len(in.held) != 1 { + t.Fatalf("a report the store could not take was not held: %+v, %d held", held, len(in.held)) } store.err = nil - s.retryHeld(context.Background()) - if !held.acked || len(s.parked) != 0 { - t.Fatalf("a held report was not recorded once the store was back: %+v, %d held", held, len(s.parked)) + in.retries(context.Background(), s) + if !held.acked || len(in.held) != 0 { + t.Fatalf("a held report was not recorded once the store was back: %+v, %d held", held, len(in.held)) } - refused := &saidTo{} - s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet()} - s.handleReport(context.Background(), aReport(t, refused, "anchor", "d1")) + refused := &settled{} + s, in = serving() + s.listener = heardWith{err: errors.New("a report named no node")} + s.act(context.Background(), in.sends(t, refused, KindReport, aReport("anchor", "d1"))) if !refused.acked || refused.nacked { t.Fatalf("a report the store answered no to was not acknowledged: %+v", refused) } @@ -72,33 +49,35 @@ func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) { // A newer report from the same node supersedes one of its reports still held: recorded after the // newer, the older would overwrite what the node is doing now. func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) { - s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet()} - older, newer, other := &saidTo{}, &saidTo{}, &saidTo{} - s.handleReport(context.Background(), aReport(t, older, "anchor", "d1")) - s.handleReport(context.Background(), aReport(t, other, "laptop", "d7")) - s.handleReport(context.Background(), aReport(t, newer, "anchor", "d2")) + s, in := serving() + s.listener = heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))} + older, newer, other := &settled{}, &settled{}, &settled{} + s.act(context.Background(), in.sends(t, older, KindReport, aReport("anchor", "d1"))) + s.act(context.Background(), in.sends(t, other, KindReport, aReport("laptop", "d7"))) + s.act(context.Background(), in.sends(t, newer, KindReport, aReport("anchor", "d2"))) if !older.acked { t.Fatalf("the older report was not set aside by the newer: %+v", older) } - if !newer.unsettled() || !other.unsettled() || len(s.parked) != 2 { + if !newer.unsettled() || !other.unsettled() || len(in.held) != 2 { t.Fatalf("the newer report and another node's were not both held: newer %+v other %+v, %d held", - newer, other, len(s.parked)) + newer, other, len(in.held)) } } // A store that has not come back within the bound is not restarting: the report is let go, loudly, // rather than held for ever. func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) { - s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))}, - log: quiet(), giveUp: time.Millisecond} - held := &saidTo{} - s.handleReport(context.Background(), aReport(t, held, "anchor", "d1")) + s, in := serving() + s.listener = heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))} + s.giveUp = time.Millisecond + held := &settled{} + s.act(context.Background(), in.sends(t, held, KindReport, aReport("anchor", "d1"))) if !held.unsettled() { t.Fatalf("the first failure was not held: %+v", held) } time.Sleep(5 * time.Millisecond) - s.retryHeld(context.Background()) - if !held.acked || len(s.parked) != 0 { - t.Fatalf("a report past the bound was not let go: %+v, %d held", held, len(s.parked)) + in.retries(context.Background(), s) + if !held.acked || len(in.held) != 0 { + t.Fatalf("a report past the bound was not let go: %+v, %d held", held, len(in.held)) } } diff --git a/internal/link/serve.go b/internal/link/serve.go index b61f633..f1c5750 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -7,31 +7,30 @@ import ( "encoding/json" "errors" "fmt" - "github.com/novox/mesh-controller/internal/envfile" - "github.com/novox/mesh-controller/internal/inventory" "log" "os" "time" amqp "github.com/rabbitmq/amqp091-go" + + "github.com/novox/mesh-controller/internal/envfile" ) -// AMQPVar is the control plane's own connection to the broker. +// AMQPVar is the controller's own connection to the bus the mesh runs on today. const AMQPVar = "MESH_BROKER_AMQP" -// Enroller is what the control plane does with an enrolment request. +// Enroller is what the controller does with an enrolment request. // -// An interface so the serving loop can be tested against a real broker without a database, and -// so the two concerns — moving messages, and deciding — stay apart. +// An interface so the serving loop can be tested against a real bus without a database, and so the +// two concerns — moving messages, and deciding — stay apart. type Enroller interface { // Enrol spends the token, records the key, and reports the node's name. The error is // returned to the node as a refusal; it must be the same for every reason a token can fail. Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, error) } -// Server consumes what nodes say. -// Listener is what the control plane does with a report. Separate from Enroller so the two can -// be given independently, and so a server that only sends declarations needs neither. +// Listener is what the controller does with a report. Separate from Enroller so the two can be +// given independently, and so a server that only sends declarations needs neither. type Listener interface { Heard(ctx context.Context, report Report) error } @@ -46,106 +45,77 @@ type Recorder interface { Built(ctx context.Context, result BuildResult) error } -type Server struct { - conn *amqp.Connection - channel *amqp.Channel - enroller Enroller - listener Listener - recorder Recorder - log *log.Logger - upgrader Upgrader - replayer Replayer - // Messages the store could not take right now, held unacknowledged and tried again on a - // ticker, by subject (novox/hq issues 082, 083). again is the ticker's interval, zero meaning - // TryAgainAfter; giveUp is how long one is kept, zero meaning GiveUpAfter. - again time.Duration - giveUp time.Duration - parked map[string]*held -} - -// held is one message the store could not take, kept to be tried again. -type held struct { - delivery amqp.Delivery - retry func(context.Context, amqp.Delivery) - what string - first time.Time -} - -// ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and -// asking again, as when the store is restarting (novox/hq issue 082). -var ErrTryAgain = errors.New("not now, try again") - -// TryAgainAfter is how often messages the store could not take are tried again. A store comes -// back in seconds, and a report a few seconds late is still current. -const TryAgainAfter = 2 * time.Second - -// GiveUpAfter bounds how long one message is kept trying. A store that has not come back in this -// long is not restarting, and the message is let go with a line saying it was lost. -const GiveUpAfter = 2 * time.Minute - -// Prefetch is how many messages the broker hands the control plane before it has settled them. -// More than one because a message the store could not take is held, unsettled, while the loop goes -// on answering others — an enrolment above all, which a host is waiting on (novox/hq issue 083). -// Bounded, because what is held is also what the broker has not kept on its own disk as pending. -const Prefetch = 64 - -// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to -// answer — an enrolment above all — while others wait for the store. -const PrefetchHeadroom = 8 - -// Records tells the server where to keep build results. -// -// Set after Connect rather than passed to it, because a control plane that only publishes — the -// `build` command, which waits for its own answer — needs a connection and no recorder, and -// making it supply one would have it construct something it never uses. -func (s *Server) Records(r Recorder) { s.recorder = r } - -// Upgrader is what the control plane does when the catalogue says a module moved. +// Upgrader is what the controller does when the catalogue says a module moved. // // An interface for the same reason Enroller is one: deciding what an upgrade means for the // machines running it is a different concern from noticing that one was announced, and only the // first needs a database. type Upgrader interface { // Upgraded is told which module moved and between which commits. An error is logged and the - // message is not requeued: an upgrade the control plane could not act on is not one it will - // act on by being handed the same message again, and a poison message on a durable queue - // would stop every upgrade behind it — except the store unreachable for the moment, which is - // asked again for a bounded time (novox/hq issue 083). + // message is not handed back: an upgrade the controller could not act on is not one it will + // act on by being given the same message again, and a poison message on a durable queue would + // stop every upgrade behind it — except the store unreachable for the moment, which is asked + // again for a bounded time (novox/hq issue 083). Upgraded(ctx context.Context, u Upgraded) error } -// Follows says what to do about upgrades, and binds the queue they arrive on. +// Server acts on what nodes and modules say. // -// **Not bound unless something is listening.** A durable queue bound to every upgrade with no -// consumer fills up quietly, and the first symptom is a broker out of disk rather than anything -// about modules. +// **It holds no transport.** What arrives comes through Inbound and what it publishes goes through +// Bus, so this file is the controller's *decisions* about messages and nothing about wires. The +// connection below is the bus the mesh runs on today, kept because the command line publishes over +// the same one until the rollout (ADR 0116). +type Server struct { + inbound Inbound + bus Bus + conn *amqp.Connection + channel *amqp.Channel + + enroller Enroller + listener Listener + recorder Recorder + upgrader Upgrader + replayer Replayer + + log *log.Logger + // giveUp is how long one message is held for the store; zero means GiveUpAfter. + giveUp time.Duration +} + +// ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and +// asking again, as when the store is restarting (novox/hq issue 082). +var ErrTryAgain = errors.New("not now, try again") + +// GiveUpAfter bounds how long one message is kept trying. A store that has not come back in this +// long is not restarting, and the message is let go with a line saying it was lost. +const GiveUpAfter = 2 * time.Minute + +// Records tells the server where to keep build results. +// +// Set after Connect rather than passed to it, because a controller that only publishes — the +// `build` command, which waits for its own answer — needs a connection and no recorder, and making +// it supply one would have it construct something it never uses. +func (s *Server) Records(r Recorder) { s.recorder = r } + +// Follows says what to do about upgrades, and asks for them to be delivered. func (s *Server) Follows(u Upgrader) error { - if _, err := s.channel.QueueDeclare(UpgradeQueue, true, false, false, false, nil); err != nil { - return fmt.Errorf("cannot declare the %s queue: %w", UpgradeQueue, err) - } - if err := s.channel.QueueBind(UpgradeQueue, KeyModuleUpgraded, EventsExchange, false, nil); err != nil { - return fmt.Errorf("cannot bind %s to %s/%s: %w", UpgradeQueue, EventsExchange, KeyModuleUpgraded, err) + if err := s.inbound.Also(KindModuleMoved); err != nil { + return err } s.upgrader = u return nil } -// Answers binds the queue a catalogue's catch-up request arrives on. -// -// **Not bound unless something is listening**, for the same reason upgrades are not: a durable -// queue with no consumer fills quietly and the first symptom is a broker out of disk. +// Answers says what to do about a catalogue's catch-up request, and asks for them to be delivered. func (s *Server) Answers(r Replayer) error { - if _, err := s.channel.QueueDeclare(CatchUpQueue, true, false, false, false, nil); err != nil { - return fmt.Errorf("cannot declare the %s queue: %w", CatchUpQueue, err) - } - if err := s.channel.QueueBind(CatchUpQueue, KeyCatchingUp, EventsExchange, false, nil); err != nil { - return fmt.Errorf("cannot bind %s to %s/%s: %w", CatchUpQueue, EventsExchange, KeyCatchingUp, err) + if err := s.inbound.Also(KindCatchUp); err != nil { + return err } s.replayer = r return nil } -// Connect opens the control plane's own connection to the broker. +// Connect opens the controller's own connection to the bus. // // On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq // 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may @@ -164,7 +134,7 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) { conn, err := amqp.Dial(url) if err != nil { - // Not quoted back: the URL carries the control plane's own broker password. + // Not quoted back: the URL carries the controller's own bus password. return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err) } channel, err := conn.Channel() @@ -173,17 +143,17 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) { return nil, err } - // Declared here rather than assumed. The control plane is the only thing that may create - // them — a node's account can write to this exchange and read its own queue, and configure - // nothing else, so a node arriving before the control plane has ever run finds nothing and - // says so, rather than quietly creating a topology nobody designed. + // Declared here rather than assumed. The controller is the only thing that may create them — a + // node's account can write to this exchange and read its own queue, and configure nothing + // else, so a node arriving before the controller has ever run finds nothing and says so, + // rather than quietly creating a topology nobody designed. if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil { conn.Close() return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err) } - // The events exchange too. The control plane is not the only publisher on it — modules - // announce onto it with their own accounts — but it is the only thing permitted to create it, - // for the same reason it is the only thing permitted to create the direct one. + // The events exchange too. The controller is not the only publisher on it — modules announce + // onto it with their own accounts — but it is the only thing permitted to create it, for the + // same reason it is the only thing permitted to create the direct one. if err := channel.ExchangeDeclare(EventsExchange, "topic", true, false, false, false, nil); err != nil { conn.Close() return nil, fmt.Errorf("cannot declare the %s exchange: %w", EventsExchange, err) @@ -192,7 +162,7 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) { conn.Close() return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err) } - // Every key a node may publish. Binding one and forgetting another is a message the broker + // Every key a node may publish. Binding one and forgetting another is a message the bus // accepts, finds no queue for, and drops — the publisher sees success and the consumer sees // nothing. That is exactly what happened to reports: `report` was left unbound while `enrol` // worked, so nodes announced what they had applied into a void for an afternoon. @@ -203,14 +173,24 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) { } } - return &Server{conn: conn, channel: channel, enroller: enroller, listener: listener, - log: log.New(os.Stdout, "", log.LstdFlags)}, nil + return &Server{ + inbound: Current(conn, channel), + bus: OverCurrent{Channel: channel}, + conn: conn, + channel: channel, + enroller: enroller, + listener: listener, + log: log.New(os.Stdout, "", log.LstdFlags), + }, nil } -// Channel is the control plane's channel, for sending declarations. +// Channel is the controller's channel, for the command line's own publishing. func (s *Server) Channel() *amqp.Channel { return s.channel } func (s *Server) Close() { + if s.inbound != nil { + s.inbound.Close() + } if s.channel != nil { _ = s.channel.Close() } @@ -219,131 +199,120 @@ func (s *Server) Close() { } } -// Serve consumes until the context ends. -// -// One consumer, deliberately: with two, the broker would round-robin between them and each would -// receive half of what it expects — a fault this project has already had, between a module's -// daemon and its capability server. +// Serve acts on what arrives until the context ends. func (s *Server) Serve(ctx context.Context) error { - // A bounded prefetch rather than one. The loop still takes messages one at a time; what the - // prefetch buys is that a message the store could not take can be held while the loop goes on - // to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything held - // goes back to the broker if the control plane stops, because nothing held is acknowledged. - if err := s.channel.Qos(Prefetch, 0, false); err != nil { - return err - } - - deliveries, err := s.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane", - false, false, false, false, nil) - if err != nil { - return err - } - - // The upgrade queue, when something is listening for them. A second queue rather than a - // second consumer on the first: two consumers on one queue split its messages between them, - // which is the fault the comment above exists about. Two queues share nothing. - var upgrades <-chan amqp.Delivery + s.log.Printf("consuming what nodes say: %s, %s, %s, %s", + KindEnrolment, KindReport, KindHeartbeat, KindBuilt) if s.upgrader != nil { - upgrades, err = s.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades", - false, false, false, false, nil) - if err != nil { - return err - } + s.log.Printf("following %s", KindModuleMoved) } - - // Its own queue and its own consumer, for the reason above: two consumers on one queue split - // its messages, and a catch-up request going to whichever half was not listening is a gap that - // looks like a working mesh. - var catchups <-chan amqp.Delivery if s.replayer != nil { - catchups, err = s.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup", - false, false, false, false, nil) - if err != nil { - return err - } - } - - closed := s.conn.NotifyClose(make(chan *amqp.Error, 1)) - s.log.Printf("consuming %s, bound to %s/{%s,%s,%s,%s}", - ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt) - if s.upgrader != nil { - s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded) - } - - again := s.again - if again == 0 { - again = TryAgainAfter - } - ticker := time.NewTicker(again) - defer ticker.Stop() - - for { - select { - case <-ctx.Done(): - return nil - case <-ticker.C: - if ctx.Err() != nil { - return nil - } - s.retryHeld(ctx) - case delivery, ok := <-catchups: - if !ok { - if catchups != nil { - return errors.New("the broker stopped delivering catch-up requests") - } - continue - } - s.catchingUp(ctx, delivery) - case delivery, ok := <-upgrades: - // A nil channel blocks for ever, so this case simply never fires when nothing is - // listening for upgrades. Closed is different, and means the broker stopped. - if !ok { - if upgrades != nil { - return errors.New("the broker stopped delivering upgrades") - } - continue - } - s.upgraded(ctx, delivery) - case reason := <-closed: - // Said rather than returned quietly. A control plane whose broker connection dropped - // is a mesh where nothing can be told anything, and the reason is the first thing - // anybody will want. - return fmt.Errorf("the broker connection closed: %v", reason) - case delivery, ok := <-deliveries: - if !ok { - return errors.New("the broker stopped delivering") - } - s.handle(ctx, delivery) - } + s.log.Printf("answering %s", KindCatchUp) } + return s.inbound.Receive(ctx, s.act) } -func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) { - switch delivery.RoutingKey { - case KeyEnrol: - s.handleEnrol(ctx, delivery) - case KeyReport: - s.handleReport(ctx, delivery) - case KeyAlive: - s.handleAlive(delivery) - case KeyBuilt: - s.handleBuilt(ctx, delivery) +// act is one message, whichever bus it came over. +func (s *Server) act(ctx context.Context, m Control) { + switch m.Kind() { + case KindEnrolment: + s.enrolling(ctx, m) + case KindReport: + s.reported(ctx, m) + case KindHeartbeat: + s.heartbeat(m) + case KindBuilt: + s.wasBuilt(ctx, m) + case KindModuleMoved: + s.moved(ctx, m) + case KindCatchUp: + s.catchingUp(ctx, m) default: - // Rejected without requeue: a message nothing understands will not be understood on the - // next attempt either, and requeuing it would spin. - s.log.Printf("refusing a message with routing key %q", delivery.RoutingKey) - _ = delivery.Reject(false) + // Dropped: a message nothing understands will not be understood on the next attempt + // either, and asking for it again would spin. + s.log.Printf("refusing a message the mesh has no name for: %q", m.Kind()) + _ = m.Drop() } } -// handleAlive records that a node was heard from, and nothing else. +// decide asks the window what to do about one message, and holds it when that is the answer — +// because holding is the one verdict that means the same thing everywhere. // -// Deliberately silent: a node saying it is there every minute would fill the log with the -// ordinary case, and a log where the ordinary case is loud is a log nobody reads. -func (s *Server) handleAlive(delivery amqp.Delivery) { +// The other three come back, because "settled without acting" is a build result dropped and an +// upgrade taken, and only the handler knows which its message is. Hold means the message is the +// bus's problem now and **the caller must not settle it**; that is also what a cancelled context +// gets, because shutting down is not an answer about a message (novox/hq issue 083, on review). +// +// declaredIn is the declaration this message is about, empty when it is about none; outstanding is +// what the mesh last sent that node, empty when it has sent none or could not be asked. +func (s *Server) decide(ctx context.Context, m Control, what, declaredIn, outstanding string, + err error) Verdict { + + if ctx.Err() != nil { + return Hold + } + window := StoreWindow{GiveUpAfter: s.giveUpAfter()} + switch v := window.Decide(err, declaredIn, outstanding, m.HeldFor()); v { + case Hold: + // Said once, on the first hold. Every redelivery saying it again would fill the log with + // one store restart. + first := m.HeldFor() == 0 + if held := m.Hold(RedeliverAfter(m.HeldFor())); held != nil { + s.log.Printf("LOST %s: it could not be held for the store (%v): %v", what, held, err) + return GiveUp + } + if first { + s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err) + } + return Hold + case Stale: + s.log.Printf("set aside %s: the mesh has moved past that declaration", what) + return Stale + case GiveUp: + s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err) + return GiveUp + default: + return v + } +} + +func (s *Server) giveUpAfter() time.Duration { + if s.giveUp == 0 { + return GiveUpAfter + } + return s.giveUp +} + +// outstanding is the digest of the declaration the mesh last sent a node. +// +// Nothing is guessed when it cannot be answered: a listener that keeps no record of what was sent +// gives a report nothing to be stale against, and a store that cannot be read will hold the report +// anyway, so there is nothing for the check to decide. +func (s *Server) outstanding(ctx context.Context, node string) string { + if node == "" { + return "" + } + asks, ok := s.listener.(Outstanding) + if !ok { + return "" + } + digest, err := asks.Outstanding(ctx, node) + if err != nil { + return "" + } + return digest +} + +// heartbeat records that a node was heard from, and nothing else. +// +// Deliberately silent: a node saying it is there every minute would fill the log with the ordinary +// case, and a log where the ordinary case is loud is a log nobody reads. Not held for the store +// either — the next heartbeat is a minute away, and a heartbeat kept for two minutes to be written +// late says nothing the one after it will not say better. +func (s *Server) heartbeat(m Control) { var alive Alive - if err := json.Unmarshal(delivery.Body, &alive); err != nil || alive.Node == "" { - _ = delivery.Reject(false) + if err := json.Unmarshal(m.Body(), &alive); err != nil || alive.Node == "" { + _ = m.Drop() return } if s.listener != nil { @@ -351,34 +320,49 @@ func (s *Server) handleAlive(delivery amqp.Delivery) { s.log.Printf("could not record that %s is here: %v", alive.Node, err) } } - _ = delivery.Ack(false) + _ = m.Took() } -// handleReport records what a node says it did. +// reported records what a node says it did. // // A node states; nothing here writes anything the node claimed about itself beyond that it was // heard from. What it applied is its own account of its own machine, and the mesh keeps the last // one as a copy for recovery rather than as a source (novox/hq 09-the-node-lifecycle). -func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { +func (s *Server) reported(ctx context.Context, m Control) { var report Report - if err := json.Unmarshal(delivery.Body, &report); err != nil { + if err := json.Unmarshal(m.Body(), &report); err != nil { s.log.Printf("a report could not be read: %v", err) - _ = delivery.Reject(false) + _ = m.Drop() return } - // A node's newer report supersedes one of its older reports still held: the older is its - // past, and recorded after the newer it would overwrite what the node is doing now. - subject := "report " + report.Node - s.supersede(subject, delivery) + m.About("report " + report.Node) + what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) + if s.listener != nil { - if err := s.listener.Heard(context.Background(), report); err != nil { - // Held, not acknowledged, while the store cannot take it: the node reports an apply - // once, and a report lost here is a node the mesh never hears from again — the store + // **Asked before the store, not after** (design 25 §3). A report about a declaration the + // mesh has moved past would otherwise wait out a restarting store to be written and then + // overwrite what the node is doing now — and on the bus being built, where the holding is + // the server's, it comes back after the newer was applied whatever the controller does. + outstanding := s.outstanding(ctx, report.Node) + declaredIn := staleAgainst(report) + if Superseded(declaredIn, outstanding) { + s.log.Printf("set aside %s: the mesh has moved past that declaration", what) + _ = m.Took() + return + } + + err := s.listener.Heard(context.Background(), report) + switch s.decide(ctx, m, what, declaredIn, outstanding, err) { + case Hold: + // Held, not settled, while the store cannot take it: the node reports an apply once, + // and a report lost here is a node the mesh never hears from again — the store // restarting under the adoption that node just applied lost exactly that (issue 082). - what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) - if s.tryLater(ctx, delivery, subject, what, err, s.handle) { - return - } + return + case Stale, GiveUp: + _ = m.Took() + return + } + if err != nil { // Said rather than swallowed. A report the mesh heard and failed to write down is a // node whose recovery copy is silently older than it looks. s.log.Printf("could not record %s's report: %v", report.Node, err) @@ -393,124 +377,49 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { default: s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied)) } - s.settled(subject, delivery) - _ = delivery.Ack(false) + _ = m.Took() } -// tryLater holds a message the store could not take right now, to be tried again on the ticker, -// and says whether it did (novox/hq issues 082, 083). +// staleAgainst is the declaration a report may be judged stale against, and it is empty for a +// report that is not only an account of an apply. // -// "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error the -// inventory reads as an outage. Anything else is an answer, and is left to the caller to settle. -// Held means unacknowledged and set aside: the loop goes on to the next message, so an enrolment a -// host is waiting on is answered while a report waits for the store. One message is held at most -// giveUpAfter; past it, it is let go with a line saying it was lost, and the caller settles it. -func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, subject, what string, err error, - retry func(context.Context, amqp.Delivery)) bool { - // Shutting down: nothing is settled. Unsettled, the broker hands the message to whatever - // consumes next — a cancelled context is not an answer about the message (issue 083, review). - if ctx.Err() != nil { - return true +// **A report carries two different things, and only one of them is about a declaration.** What the +// node applied is; what the machine *is* — the tunnel it took over, the ports its own bundle holds, +// what an adopted node found and is keeping, a node moving its overlay key — is not. Those reach +// the mesh on a report because a report is the message a node sends, and nowhere else: a rekey +// dropped as stale is a node whose overlay key never moves, and no retry is coming, because the +// node said it once. +// +// So staleness is asked only of a report that is purely an apply's account. The rest is acted on +// whenever it arrives, which is the behaviour the mesh has had all along. +func staleAgainst(report Report) string { + if report.Rekey != nil || report.Tunnel != nil || len(report.Held) > 0 || + report.Firewall != "" || len(report.Reachable) > 0 || len(report.Carried) > 0 { + return "" } - if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) { - return false - } - if s.parked == nil { - s.parked = map[string]*held{} - } - h, ok := s.parked[subject] - if !ok || h.delivery.DeliveryTag != delivery.DeliveryTag { - // Held no further than the prefetch leaves room: past it, the broker would hand the loop - // nothing new — enrolments included — until something held was let go. - if !ok && len(s.parked) >= Prefetch-PrefetchHeadroom { - s.log.Printf("LOST %s: %d messages are already held for the store, and holding more "+ - "would stop the queue: %v", what, len(s.parked), err) - return false - } - h = &held{delivery: delivery, retry: retry, what: what, first: time.Now()} - s.parked[subject] = h - s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err) - return true - } - if time.Since(h.first) >= s.giveUpAfter() { - delete(s.parked, subject) - s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err) - return false - } - return true + return report.Declared } -// supersede drops a message held for a subject when a newer one for it arrives: the older is -// acknowledged, because acting on it after the newer would undo the newer. -func (s *Server) supersede(subject string, newer amqp.Delivery) { - h, ok := s.parked[subject] - if !ok || h.delivery.DeliveryTag == newer.DeliveryTag { - return - } - delete(s.parked, subject) - s.log.Printf("set aside %s: a newer one arrived", h.what) - _ = h.delivery.Ack(false) -} - -// settled forgets a message once it has been handled either way. -func (s *Server) settled(subject string, delivery amqp.Delivery) { - if h, ok := s.parked[subject]; ok && h.delivery.DeliveryTag == delivery.DeliveryTag { - delete(s.parked, subject) - } -} - -// digest names a message by its content. -func digest(body []byte) string { - sum := sha256.Sum256(body) - return hex.EncodeToString(sum[:]) -} - -// retryHeld tries every held message again. Each handler holds it again, settles it, or lets it -// go past the bound. -func (s *Server) retryHeld(ctx context.Context) { - for _, h := range s.snapshot() { - if ctx.Err() != nil { - return - } - h.retry(ctx, h.delivery) - } -} - -func (s *Server) snapshot() []*held { - out := make([]*held, 0, len(s.parked)) - for _, h := range s.parked { - out = append(out, h) - } - return out -} - -func (s *Server) giveUpAfter() time.Duration { - if s.giveUp == 0 { - return GiveUpAfter - } - return s.giveUp -} - -func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { +func (s *Server) enrolling(ctx context.Context, m Control) { reply := EnrolReply{Refusal: "that token cannot be used"} var request EnrolRequest - if err := json.Unmarshal(delivery.Body, &request); err != nil { + if err := json.Unmarshal(m.Body(), &request); err != nil { s.log.Printf("an enrolment request could not be read: %v", err) } else { - request.Redelivered = delivery.Redelivered + request.Redelivered = m.Redelivered() accepted, err := s.enroller.Enrol(ctx, request) switch { case errors.Is(err, ErrTryAgain): // Not a refusal: nothing was spent, and the same request asked again will be - // answered. Replied at once rather than held, so the node — which is waiting on - // this answer — decides when to ask, and the queue behind it moves (issue 083). + // answered. Replied at once rather than held, so the node — which is waiting on this + // answer — decides when to ask, and the queue behind it moves (issue 083). reply = EnrolReply{TryAgain: true, Refusal: "the mesh cannot answer right now; ask again"} s.log.Printf("asked %q to enrol again shortly: %v", request.Node, err) case err != nil && request.Redelivered: - // Said as what it most likely is: the broker handed this request over again after - // the control plane stopped mid-answer, and an enrolment already spent is not - // finished a second time. The node may need a new token. + // Said as what it most likely is: the bus handed this request over again after the + // controller stopped mid-answer, and an enrolment already spent is not finished a + // second time. The node may need a new token. s.log.Printf("refusing a redelivered enrolment for %q — it may have finished before "+ "the control plane stopped, and if the node did not get its answer it needs a new "+ "token: %v", request.Node, err) @@ -524,169 +433,160 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { } } - s.reply(ctx, delivery, reply) - - // Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the - // request on the broker rather than having consumed it silently. Asked again by the same - // presenter, an enrolment finishes: the token is held for its key and spent last (issue 083). - _ = delivery.Ack(false) -} - -func (s *Server) reply(ctx context.Context, delivery amqp.Delivery, reply EnrolReply) { - if delivery.ReplyTo == "" { - s.log.Print("an enrolment request named no reply queue, so nothing can be told the answer") - return - } - body, err := json.Marshal(reply) - if err != nil { + if body, err := json.Marshal(reply); err != nil { s.log.Printf("cannot encode a reply: %v", err) - return + } else { + answer, cancel := context.WithTimeout(ctx, 10*time.Second) + if err := m.Answer(answer, body); err != nil { + s.log.Printf("cannot answer an enrolment: %v", err) + } + cancel() } - timeout, cancel := context.WithTimeout(ctx, 10*time.Second) - defer cancel() - if err := s.channel.PublishWithContext(timeout, "", delivery.ReplyTo, false, false, - amqp.Publishing{ - ContentType: "application/json", - CorrelationId: delivery.CorrelationId, - Body: body, - }); err != nil { - s.log.Printf("cannot reply to %s: %v", delivery.ReplyTo, err) - } + // Settled after the answer is sent, so a controller that dies mid-answer leaves the request on + // the bus rather than having consumed it silently. Asked again by the same presenter, an + // enrolment finishes: the token is held for its key and spent last (issue 083). + _ = m.Took() } -// handleBuilt keeps what a builder said, whichever way it went. +// wasBuilt keeps what a builder said, whichever way it went. // // This is for results nobody was waiting for. A build asked for with `build` is answered directly // to the asker; one triggered any other way is published here, and without this it would be // reported into the void — which is the same as not reporting it. -func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { +func (s *Server) wasBuilt(ctx context.Context, m Control) { var result BuildResult - if err := json.Unmarshal(delivery.Body, &result); err != nil { + if err := json.Unmarshal(m.Body(), &result); err != nil { s.log.Printf("a build result could not be read: %v", err) - _ = delivery.Reject(false) + _ = m.Drop() return } if s.recorder == nil { - // Nothing to keep it in. Rejected rather than dropped silently, so the broker's own - // counters show something arriving that nothing handles. + // Nothing to keep it in. Dropped rather than swallowed, so the bus's own counters show + // something arriving that nothing handles. s.log.Printf("a build result arrived and this control plane keeps none") - _ = delivery.Reject(false) + _ = m.Drop() return } // Each build result its own subject: none supersedes another, and recording one twice is // harmless — the build is kept by its id. - subject := "build " + digest(delivery.Body) - s.supersede(subject, delivery) - if err := s.recorder.Built(ctx, result); err != nil { + m.About("build " + digest(m.Body())) + err := s.recorder.Built(ctx, result) + switch s.decide(ctx, m, fmt.Sprintf("a build result from %s", result.On), "", "", err) { + case Hold: // Held while the store cannot take it: a build result lost here is never announced (083). - if s.tryLater(ctx, delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) { - return - } + return + case Stale, GiveUp: + _ = m.Drop() + return + } + if err != nil { s.log.Printf("cannot keep a build result from %s: %v", result.On, err) - s.settled(subject, delivery) - _ = delivery.Reject(false) + _ = m.Drop() return } - s.settled(subject, delivery) switch { case result.Failed != "": s.log.Printf("%s could not build %s", result.On, result.Repository) default: s.log.Printf("%s built %s from %s", result.On, result.Repository, result.Commit) } - _ = delivery.Ack(false) + _ = m.Took() } // catchingUp answers a catalogue that has just started and may have missed builds. // -// Acknowledged after the work. A replay that fails for a reason other than the store is not one -// that succeeds by being handed the same request again, so that is acknowledged and said; but a -// store that could not be read right now is held and asked again, bounded, rather than the -// request lost until the catalogue next restarts (issue 083). One request stands for all: a newer -// one supersedes one still held. -func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { - const subject = "catch-up" - s.supersede(subject, delivery) - holding := false - defer func() { - if !holding { - s.settled(subject, delivery) - _ = delivery.Ack(false) - } - }() +// Settled after the work. A replay that fails for a reason other than the store is not one that +// succeeds by being handed the same request again, so that is settled and said; but a store that +// could not be read right now is held and asked again, bounded, rather than the request lost until +// the catalogue next restarts (issue 083). One request stands for all: a newer one supersedes one +// still held. +func (s *Server) catchingUp(ctx context.Context, m Control) { + m.About("catch-up") if s.replayer == nil { s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay") + _ = m.Took() return } announcements, err := s.replayer.Announceable(ctx) + switch s.decide(ctx, m, "a catalogue's request to catch up", "", "", err) { + case Hold: + return + case Stale, GiveUp: + _ = m.Took() + return + } if err != nil { - if s.tryLater(ctx, delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) { - holding = true - return - } s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err) + _ = m.Took() return } sent := 0 for _, a := range announcements { a.Replay = true - if err := EmitEvent(ctx, OverCurrent{Channel: s.channel}, KeyModuleBuilt, "control-plane", "", a); err != nil { + if err := EmitEvent(ctx, s.bus, KeyModuleBuilt, "control-plane", "", a); err != nil { // Said and abandoned rather than retried: the catalogue asks again every time it // starts, and half a graph delivered twice is no better than half delivered once. s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v", a.Module, short(a.Commit), err) + _ = m.Took() return } sent++ } s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent) + _ = m.Took() } -// upgraded hands one announcement to whatever is following them. +// moved hands one announcement to whatever is following them. // -// **Acknowledged whatever happens, but one thing.** A failure here is usually the control plane -// being unable to act on an upgrade — a machine that cannot be resolved, a broker that will not -// take a declaration — and none of those get better by being handed the same message again. The -// one exception is the upgrader saying the store could not be read for the moment (ErrTryAgain): -// that is held and asked again, bounded (novox/hq issue 083). Only the upgrader's word counts -// here, not an error that merely looks like an outage — a push that timed out on the second -// machine is not asked again, or the first would be pushed every few seconds for two minutes. -// A newer move of the same module supersedes one still held. -func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { +// **Settled whatever happens, but one thing.** A failure here is usually the controller being +// unable to act on an upgrade — a machine that cannot be resolved, a bus that will not take a +// declaration — and none of those get better by being handed the same message again. The one +// exception is the upgrader saying the store could not be read for the moment (ErrTryAgain): that +// is held and asked again, bounded (novox/hq issue 083). Only the upgrader's word counts here, not +// an error that merely looks like an outage — a push that timed out on the second machine is not +// asked again, or the first would be pushed every few seconds for two minutes. A newer move of the +// same module supersedes one still held. +func (s *Server) moved(ctx context.Context, m Control) { var u Upgraded - _ = json.Unmarshal(delivery.Body, &u) - subject := "upgrade " + u.Module - s.supersede(subject, delivery) - holding := false - defer func() { - if !holding { - s.settled(subject, delivery) - _ = delivery.Ack(false) - } - }() - if err := json.Unmarshal(delivery.Body, &u); err != nil { + if err := json.Unmarshal(m.Body(), &u); err != nil { s.log.Printf("an upgrade announcement could not be read: %v", err) + _ = m.Took() return } + m.About("upgrade " + u.Module) if u.Module == "" { s.log.Printf("an upgrade announcement named no module; ignored") + _ = m.Took() return } - if err := s.upgrader.Upgraded(ctx, u); err != nil { - // Shutting down is not an answer about the announcement: left for the broker. - if ctx.Err() != nil { - holding = true - return - } - if errors.Is(err, ErrTryAgain) && - s.tryLater(ctx, delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) { - holding = true - return - } + err := s.upgrader.Upgraded(ctx, u) + // Only "not now" is worth holding. Anything else is an answer, and the window would read a + // timed-out push as an outage and ask for it again. + holdable := err + if !errors.Is(err, ErrTryAgain) { + holdable = nil + } + what := fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)) + switch s.decide(ctx, m, what, "", "", holdable) { + case Hold: + return + case Stale, GiveUp: + _ = m.Took() + return + } + if err != nil { s.log.Printf("%s moved to %s and the mesh could not act on it: %v", u.Module, short(u.Commit), err) } + _ = m.Took() +} + +// digest names a message by its content. +func digest(body []byte) string { + sum := sha256.Sum256(body) + return hex.EncodeToString(sum[:]) } // short is a commit as people read it. diff --git a/internal/link/stale_report_test.go b/internal/link/stale_report_test.go new file mode 100644 index 0000000..2e0b949 --- /dev/null +++ b/internal/link/stale_report_test.go @@ -0,0 +1,135 @@ +package link + +import ( + "context" + "errors" + "testing" +) + +// Supersession as a check rather than a memory (design 25 §3). +// +// Holding a message in memory let the controller drop an older report when a newer one for the +// same node arrived. On the bus being built the message belongs to the server and comes back +// whatever happened meanwhile — so the older report is redelivered *after* the newer was applied, +// and acting on it would undo the newer. +// +// The answer was already in the message: a report carries the digest of the declaration it is +// about, so "is this the past?" is a question the message answers. + +// sentAndHeard records reports and knows what was last sent, which is the pair the check needs. +type sentAndHeard struct { + sent string + heard []Report + err error +} + +func (s *sentAndHeard) Heard(_ context.Context, r Report) error { + if s.err != nil { + return s.err + } + s.heard = append(s.heard, r) + return nil +} + +func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil } + +// A report about a declaration the mesh has moved past is settled and not acted on. Settled rather +// than dropped, because there is nothing wrong with the message — it is simply the past, and +// redelivering it for ever is worse than letting it go. +func TestAReportAboutASupersededDeclarationIsNotActedOn(t *testing.T) { + store := &sentAndHeard{sent: "d2"} + s, in := serving() + s.listener = store + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindReport, aReport("anchor", "d1"))) + + if len(store.heard) != 0 { + t.Fatalf("a report about a superseded declaration was acted on: %+v", store.heard) + } + if !to.acked { + t.Fatalf("a superseded report was not settled, so it comes back for ever: %+v", *to) + } +} + +// The report about the declaration that *is* outstanding is acted on, and so is one from a node +// the mesh has no digest for — an older host that says nothing about which declaration it applied +// has nothing to be judged against, and refusing it would silence every node built before reports +// carried the digest. +func TestAReportAboutTheOutstandingDeclarationIsActedOn(t *testing.T) { + for _, c := range []struct{ what, sent, declared string }{ + {"the one outstanding", "d2", "d2"}, + {"a report that says nothing about which", "d2", ""}, + {"a node nothing was ever sent", "", "d1"}, + } { + store := &sentAndHeard{sent: c.sent} + s, in := serving() + s.listener = store + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindReport, aReport("anchor", c.declared))) + if len(store.heard) != 1 || !to.acked { + t.Errorf("%s: was not acted on and acknowledged: heard %+v, settled %+v", + c.what, store.heard, *to) + } + } +} + +// **Staleness is decided before the store is waited on**, not after: a redelivery that lost its +// race is not worth holding a slot in the window that a current message needs. +func TestASupersededReportIsNotHeldForTheStore(t *testing.T) { + store := &sentAndHeard{sent: "d2", err: errors.Join(ErrTryAgain, errors.New("starting up"))} + s, in := serving() + s.listener = store + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindReport, aReport("anchor", "d1"))) + if !to.acked || len(in.held) != 0 { + t.Fatalf("a superseded report waited for the store: %+v, %d held", *to, len(in.held)) + } +} + +// **Half of a report is not about a declaration, and that half is never stale.** +// +// What the machine *is* — the tunnel it took over, the ports its own bundle holds, what an adopted +// node found and is keeping, a node moving its overlay key — reaches the mesh on a report and +// nowhere else. A rekey set aside as stale is a node whose overlay key never moves, and no retry is +// coming, because the node said it once. So a report carrying any of these is acted on whenever it +// arrives, however far the mesh has moved on. +func TestAReportCarryingWhatOnlyTheNodeKnowsIsActedOnHoweverOldItIs(t *testing.T) { + for _, c := range []struct { + what string + report Report + }{ + {"a rekey", Report{Node: "anchor", Declared: "d1", + Rekey: &Rekey{Previous: "k1", OverlayKey: "k2"}}}, + {"the tunnel it carried", Report{Node: "anchor", Declared: "d1", + Tunnel: &CarriedTunnel{Interface: "wg0", State: "taken"}}}, + {"what an adopted node holds", Report{Node: "anchor", Declared: "d1", + Held: []Held{{ID: "conf", Module: "web", Kind: "file"}}}}, + {"the firewall it found", Report{Node: "anchor", Declared: "d1", Firewall: "ufw"}}, + {"what is reachable on it", Report{Node: "anchor", Declared: "d1", + Reachable: []Reach{{Protocol: "tcp", Port: 443}}}}, + {"the ports its own bundle holds", Report{Node: "anchor", Declared: "d1", + Carried: []int{5432}}}, + } { + store := &sentAndHeard{sent: "d9"} + s, in := serving() + s.listener = store + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindReport, c.report)) + if len(store.heard) != 1 { + t.Errorf("%s was set aside as stale, and the mesh will never hear it again: %+v", + c.what, *to) + } + } +} + +// A heartbeat is not held for the store: the next one is a minute away, and one kept for two +// minutes to be written late says nothing the one after it will not say better. +func TestAHeartbeatIsNotHeldForTheStore(t *testing.T) { + s, in := serving() + s.listener = heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))} + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindHeartbeat, Alive{Node: "anchor"})) + if !to.acked || len(in.held) != 0 { + t.Fatalf("a heartbeat was held for the store: %+v, %d held", *to, len(in.held)) + } +} diff --git a/internal/link/store_window_test.go b/internal/link/store_window_test.go index f0fcd2d..80b6a4b 100644 --- a/internal/link/store_window_test.go +++ b/internal/link/store_window_test.go @@ -2,15 +2,11 @@ package link import ( "context" - "encoding/json" "errors" "fmt" - "io" - "log" "testing" "github.com/jackc/pgx/v5/pgconn" - amqp "github.com/rabbitmq/amqp091-go" ) // What a store restarting under an adoption answers with (novox/hq issues 082, 083). @@ -28,29 +24,6 @@ type replaysWith struct{ err error } func (r replaysWith) Announceable(context.Context) ([]Announcement, error) { return nil, r.err } -type settledAs struct{ acked, nacked, requeued, rejected bool } - -func (a *settledAs) Ack(uint64, bool) error { a.acked = true; return nil } -func (a *settledAs) Nack(_ uint64, _ bool, requeue bool) error { - a.nacked, a.requeued = true, requeue - return nil -} -func (a *settledAs) Reject(uint64, bool) error { a.rejected = true; return nil } - -func a(t *testing.T, to *settledAs, key string, v any) amqp.Delivery { - t.Helper() - body, err := json.Marshal(v) - if err != nil { - t.Fatal(err) - } - tag++ - return amqp.Delivery{Acknowledger: to, RoutingKey: key, Body: body, DeliveryTag: tag} -} - -func quietServer() *Server { return &Server{log: log.New(io.Discard, "", 0)} } - -func (a *settledAs) held() bool { return !a.acked && !a.nacked && !a.rejected } - // A build result the store could not take right now is handed back; one it refused is rejected, // as before; one it kept is acknowledged. func TestABuildResultWaitsOutARestartingStore(t *testing.T) { @@ -58,16 +31,16 @@ func TestABuildResultWaitsOutARestartingStore(t *testing.T) { for _, c := range []struct { what string err error - want func(*settledAs) bool + want func(*settled) bool }{ - {"restarting", restarting, func(s *settledAs) bool { return s.held() }}, - {"refused", errors.New("no such module"), func(s *settledAs) bool { return s.rejected && !s.nacked }}, - {"kept", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, + {"restarting", restarting, func(s *settled) bool { return s.unsettled() }}, + {"refused", errors.New("no such module"), func(s *settled) bool { return s.rejected && !s.nacked }}, + {"kept", nil, func(s *settled) bool { return s.acked && !s.nacked }}, } { - s := quietServer() + s, in := serving() s.recorder = recordsWith{err: c.err} - to := &settledAs{} - s.handleBuilt(context.Background(), a(t, to, KeyBuilt, built)) + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindBuilt, built)) if !c.want(to) { t.Errorf("%s: a build result was settled as %+v", c.what, *to) } @@ -81,17 +54,17 @@ func TestAnUpgradeWaitsOutARestartingStoreAndNothingElse(t *testing.T) { for _, c := range []struct { what string err error - want func(*settledAs) bool + want func(*settled) bool }{ - {"the store away, said by the upgrader", errors.Join(ErrTryAgain, restarting), func(s *settledAs) bool { return s.held() }}, - {"a push that timed out", context.DeadlineExceeded, func(s *settledAs) bool { return s.acked && !s.nacked }}, - {"cannot act", errors.New("anchor cannot be resolved"), func(s *settledAs) bool { return s.acked && !s.nacked }}, - {"acted", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, + {"the store away, said by the upgrader", errors.Join(ErrTryAgain, restarting), func(s *settled) bool { return s.unsettled() }}, + {"a push that timed out", context.DeadlineExceeded, func(s *settled) bool { return s.acked && !s.nacked }}, + {"cannot act", errors.New("anchor cannot be resolved"), func(s *settled) bool { return s.acked && !s.nacked }}, + {"acted", nil, func(s *settled) bool { return s.acked && !s.nacked }}, } { - s := quietServer() + s, in := serving() s.upgrader = upgradesWith{err: c.err} - to := &settledAs{} - s.upgraded(context.Background(), a(t, to, "upgraded", moved)) + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindModuleMoved, moved)) if !c.want(to) { t.Errorf("%s: an upgrade was settled as %+v", c.what, *to) } @@ -104,16 +77,16 @@ func TestACatchUpWaitsOutARestartingStore(t *testing.T) { for _, c := range []struct { what string err error - want func(*settledAs) bool + want func(*settled) bool }{ - {"restarting", restarting, func(s *settledAs) bool { return s.held() }}, - {"unreadable", errors.New("a build row is malformed"), func(s *settledAs) bool { return s.acked && !s.nacked }}, - {"nothing to replay", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, + {"restarting", restarting, func(s *settled) bool { return s.unsettled() }}, + {"unreadable", errors.New("a build row is malformed"), func(s *settled) bool { return s.acked && !s.nacked }}, + {"nothing to replay", nil, func(s *settled) bool { return s.acked && !s.nacked }}, } { - s := quietServer() + s, in := serving() s.replayer = replaysWith{err: c.err} - to := &settledAs{} - s.catchingUp(context.Background(), a(t, to, "catch-up", map[string]string{})) + to := &settled{} + s.act(context.Background(), in.sends(t, to, KindCatchUp, map[string]string{})) if !c.want(to) { t.Errorf("%s: a catch-up request was settled as %+v", c.what, *to) } @@ -121,15 +94,15 @@ func TestACatchUpWaitsOutARestartingStore(t *testing.T) { } // Shutting down is not an answer about a message: one handled with a cancelled context is left -// unsettled, for the broker to hand to whatever consumes next (issue 083, review). -func TestAMessageHandledDuringShutdownIsLeftForTheBroker(t *testing.T) { +// unsettled, for the bus to hand to whatever consumes next (issue 083, review). +func TestAMessageHandledDuringShutdownIsLeftForTheBus(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() - s := quietServer() + s, in := serving() s.recorder = recordsWith{err: context.Canceled} - to := &settledAs{} - s.handleBuilt(ctx, a(t, to, KeyBuilt, BuildResult{On: "anchor", Repository: "/r", Commit: "abc"})) - if !to.held() { + to := &settled{} + s.act(ctx, in.sends(t, to, KindBuilt, BuildResult{On: "anchor", Repository: "/r", Commit: "abc"})) + if !to.unsettled() { t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to) } } @@ -137,45 +110,45 @@ func TestAMessageHandledDuringShutdownIsLeftForTheBroker(t *testing.T) { // Two identical build results: the newer sets the older aside rather than leaving it unsettled // for ever, holding a place in the prefetch. func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) { - s := quietServer() + s, in := serving() s.recorder = recordsWith{err: restarting} built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"} - first, second := &settledAs{}, &settledAs{} - s.handleBuilt(context.Background(), a(t, first, KeyBuilt, built)) - s.handleBuilt(context.Background(), a(t, second, KeyBuilt, built)) - if !first.acked || !second.held() || len(s.parked) != 1 { + first, second := &settled{}, &settled{} + s.act(context.Background(), in.sends(t, first, KindBuilt, built)) + s.act(context.Background(), in.sends(t, second, KindBuilt, built)) + if !first.acked || !second.unsettled() || len(in.held) != 1 { t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held", - *first, *second, len(s.parked)) + *first, *second, len(in.held)) } } // What is held stops short of the prefetch, so the loop always has room to answer an enrolment. func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) { - s := quietServer() + s, in := serving() s.recorder = recordsWith{err: restarting} - var last *settledAs + var last *settled for i := 0; i < Prefetch; i++ { - last = &settledAs{} - s.handleBuilt(context.Background(), a(t, last, KeyBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)})) + last = &settled{} + s.act(context.Background(), in.sends(t, last, KindBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)})) } - if len(s.parked) != Prefetch-PrefetchHeadroom { - t.Fatalf("%d messages were held; the ceiling is %d", len(s.parked), Prefetch-PrefetchHeadroom) + if len(in.held) != Prefetch-PrefetchHeadroom { + t.Fatalf("%d messages were held; the ceiling is %d", len(in.held), Prefetch-PrefetchHeadroom) } - if last.held() { + if last.unsettled() { t.Fatalf("a message past the ceiling was held: %+v", *last) } } -// An upgrade handled during shutdown is left for the broker too — the upgrader's error is the +// An upgrade handled during shutdown is left for the bus too — the upgrader's error is the // cancelled context, which is no answer about the announcement. -func TestAnUpgradeHandledDuringShutdownIsLeftForTheBroker(t *testing.T) { +func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() - s := quietServer() + s, in := serving() s.upgrader = upgradesWith{err: context.Canceled} - to := &settledAs{} - s.upgraded(ctx, a(t, to, "upgraded", Upgraded{Module: "gitea", Commit: "abcdef0123"})) - if !to.held() { - t.Fatalf("an upgrade handled during shutdown was settled, and so lost: %+v", *to) + to := &settled{} + s.act(ctx, in.sends(t, to, KindModuleMoved, Upgraded{Module: "gitea", Commit: "abcdef0123"})) + if !to.unsettled() { + t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to) } } diff --git a/internal/link/window.go b/internal/link/window.go index 2b89a15..91db759 100644 --- a/internal/link/window.go +++ b/internal/link/window.go @@ -52,7 +52,7 @@ func (w StoreWindow) Decide(err error, declaredIn, outstanding string, heldFor t // **Staleness is checked before the store, not after.** A redelivery that lost its race is // not worth waiting on a store for, and asking the store first would mean a message about a // superseded declaration holding a slot in the window that a current one needs. - if declaredIn != "" && outstanding != "" && declaredIn != outstanding { + if Superseded(declaredIn, outstanding) { return Stale } if err == nil { @@ -69,6 +69,20 @@ func (w StoreWindow) Decide(err error, declaredIn, outstanding string, heldFor t return Hold } +// Superseded says a message is about a declaration the mesh has already moved past. +// +// Stated on its own because it is asked in two places for one reason: here, so the window's whole +// decision is in one pure function, and by the serving loop *before* it asks the store, because +// that is the point — a message about the past must not wait on a store, or it holds a slot in the +// window that a current message needs. +// +// Unanswerable is not stale. A message that names no declaration, and a node the mesh has never +// sent one, both give nothing to compare: the mesh acts on the message rather than guessing, which +// is also what keeps a host built before reports carried the digest from going silent. +func Superseded(declaredIn, outstanding string) bool { + return declaredIn != "" && outstanding != "" && declaredIn != outstanding +} + // RedeliverAfter is how long the server should hold a naked message before trying again. // // Backed off, and bounded. A store restarting is back in seconds; a store that is gone is not