From 047830a6b576ff26e9922b8a03df7dc5b93cfc4b Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 13:25:13 +0200 Subject: [PATCH] A report the store cannot take yet is handed back to the broker, not acknowledged and lost (novox/hq issue 082) --- internal/inventory/reachable.go | 35 ++++++++++++++++ internal/inventory/reachable_test.go | 28 +++++++++++++ internal/link/enrolment.go | 9 +++- internal/link/report_retry_test.go | 61 ++++++++++++++++++++++++++++ internal/link/serve.go | 25 ++++++++++++ 5 files changed, 157 insertions(+), 1 deletion(-) create mode 100644 internal/inventory/reachable.go create mode 100644 internal/inventory/reachable_test.go create mode 100644 internal/link/report_retry_test.go diff --git a/internal/inventory/reachable.go b/internal/inventory/reachable.go new file mode 100644 index 0000000..00f79d9 --- /dev/null +++ b/internal/inventory/reachable.go @@ -0,0 +1,35 @@ +package inventory + +import ( + "errors" + "strings" + + "github.com/jackc/pgx/v5/pgconn" +) + +// Unreachable says whether an error means the store could not be asked right now, rather than +// that it answered no. +// +// The difference decides whether something worth keeping is kept or lost. The store is recreated +// when the foundation is adopted (ADR 0078), and for those seconds every write fails with "the +// database system is starting up" or a refused connection; a caller that treats that like a +// refusal throws away what it was writing — a node's report of the apply that caused the restart +// was lost exactly so, and the mesh never heard from the node again (novox/hq issue 082). +// +// Unreachable is a connection that failed, and the server classes that mean "not now": 08 +// (connection exception) and 57 (operator intervention — starting up, shutting down, admin +// shutdown, cannot connect now). Everything else is an answer, and asking again gets the same one. +func Unreachable(err error) bool { + if err == nil { + return false + } + var connect *pgconn.ConnectError + if errors.As(err, &connect) { + return true + } + var pg *pgconn.PgError + if errors.As(err, &pg) { + return strings.HasPrefix(pg.Code, "08") || strings.HasPrefix(pg.Code, "57") + } + return false +} diff --git a/internal/inventory/reachable_test.go b/internal/inventory/reachable_test.go new file mode 100644 index 0000000..f25c830 --- /dev/null +++ b/internal/inventory/reachable_test.go @@ -0,0 +1,28 @@ +package inventory + +import ( + "errors" + "fmt" + "testing" + + "github.com/jackc/pgx/v5/pgconn" +) + +// The store restarting is "not now"; a constraint the store enforced is an answer (issue 082). +func TestAStoreThatIsStartingIsUnreachableAndARefusalIsNot(t *testing.T) { + for _, c := range []struct { + err error + want bool + }{ + {&pgconn.PgError{Code: "57P03", Message: "the database system is starting up"}, true}, + {fmt.Errorf("recording: %w", &pgconn.PgError{Code: "57P01"}), true}, + {&pgconn.PgError{Code: "08006"}, true}, + {&pgconn.PgError{Code: "23503", Message: "violates foreign key constraint"}, false}, + {errors.New("a report named no node"), false}, + {nil, false}, + } { + if got := Unreachable(c.err); got != c.want { + t.Errorf("Unreachable(%v) = %v, want %v", c.err, got, c.want) + } + } +} diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index a2a77cf..b7975bb 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -144,7 +144,14 @@ 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. -func (e Enrolment) Heard(ctx context.Context, report Report) error { +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). + defer func() { + if inventory.Unreachable(err) { + err = fmt.Errorf("%w: %w", ErrTryAgain, err) + } + }() if report.Node == "" { return errors.New("a report named no node") } diff --git a/internal/link/report_retry_test.go b/internal/link/report_retry_test.go new file mode 100644 index 0000000..7d5dd16 --- /dev/null +++ b/internal/link/report_retry_test.go @@ -0,0 +1,61 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "io" + "log" + "testing" + + 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 } + +type heardWith struct{ err error } + +func (h heardWith) Heard(context.Context, Report) error { return h.err } + +func aReport(t *testing.T, to *saidTo) amqp.Delivery { + t.Helper() + body, err := json.Marshal(Report{Node: "anchor", Applied: []string{"store"}}) + if err != nil { + t.Fatal(err) + } + return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body} +} + +// A report the store could not take right now goes back to the broker to be asked again; one the +// store answered no to is acknowledged, or it would come back for ever (novox/hq issue 082). +func TestAReportTheStoreCouldNotTakeIsKeptAndOneItRefusedIsNot(t *testing.T) { + quiet := log.New(io.Discard, "", 0) + + notNow := &saidTo{} + s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet, again: 1} + s.handleReport(aReport(t, notNow)) + if notNow.acked || !notNow.nacked || !notNow.requeued { + t.Fatalf("a report the store could not take yet was not handed back to be asked again: %+v", notNow) + } + + refused := &saidTo{} + s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet, again: 1} + s.handleReport(aReport(t, refused)) + if !refused.acked || refused.nacked { + t.Fatalf("a report the store answered no to was not acknowledged, so it would spin: %+v", refused) + } + + recorded := &saidTo{} + s = &Server{listener: heardWith{}, log: quiet} + s.handleReport(aReport(t, recorded)) + if !recorded.acked || recorded.nacked { + t.Fatalf("a recorded report was not acknowledged: %+v", recorded) + } +} diff --git a/internal/link/serve.go b/internal/link/serve.go index f0a8ec3..f0b9039 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -52,8 +52,20 @@ type Server struct { log *log.Logger upgrader Upgrader replayer Replayer + // again is how long a report the store could not take waits before it is handed back to + // the broker; zero means TryAgainAfter. + again 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") + +// TryAgainAfter is the pause before a report the store could not take goes back to the broker. +// The consumer takes one message at a time, so without it a restarting store would be asked in +// a tight loop; a store comes back in seconds, and a report a few seconds late is still current. +const TryAgainAfter = 2 * time.Second + // 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 @@ -311,6 +323,19 @@ func (s *Server) handleReport(delivery amqp.Delivery) { } if s.listener != nil { if err := s.listener.Heard(context.Background(), report); err != nil { + if errors.Is(err, ErrTryAgain) { + // Kept, not acknowledged. 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 (novox/hq issue 082). + again := s.again + if again == 0 { + again = TryAgainAfter + } + s.log.Printf("could not record %s's report yet, and will again in %s: %v", report.Node, again, err) + time.Sleep(again) + _ = delivery.Nack(false, true) + return + } // 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)