The bus on NATS: both transports behind seams, and the rollout switch #87

Merged
jschoubben merged 40 commits from feat/nats-genesis into main 2026-09-27 17:36:41 +00:00
10 changed files with 1094 additions and 550 deletions
Showing only changes of commit 06cf3c04e5 - Show all commits
+19
View File
@@ -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
+11
View File
@@ -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).
+125
View File
@@ -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)
}
+317
View File
@@ -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 &currentInbound{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 &currentControl{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)
}
}
+71
View File
@@ -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 := &currentInbound{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 &currentControl{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)
}
+31 -52
View File
@@ -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))
}
}
+322 -422
View File
@@ -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.
+135
View File
@@ -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))
}
}
+48 -75
View File
@@ -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)
}
}
+15 -1
View File
@@ -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