From 3fd0c37d229f42af43130440978dffdeaab8cf5e Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 13:34:04 +0200 Subject: [PATCH] Review of 082: a store killed mid-conversation is an outage and a wrong password is not; one report holds the queue at most two minutes, then is let go loudly; shutdown does not wait out the pause --- internal/inventory/reachable.go | 45 +++++++++++++------ internal/inventory/reachable_test.go | 66 ++++++++++++++++++++++++---- internal/link/report_retry_test.go | 49 +++++++++++++++++++-- internal/link/serve.go | 54 ++++++++++++++++++++--- 4 files changed, 182 insertions(+), 32 deletions(-) diff --git a/internal/inventory/reachable.go b/internal/inventory/reachable.go index 00f79d9..5f53f93 100644 --- a/internal/inventory/reachable.go +++ b/internal/inventory/reachable.go @@ -2,7 +2,8 @@ package inventory import ( "errors" - "strings" + "io" + "net" "github.com/jackc/pgx/v5/pgconn" ) @@ -11,25 +12,41 @@ import ( // 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). +// when the foundation is adopted (ADR 0078), and for those seconds every write fails; 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. +// Unreachable is the connection failing and the server saying "not now". The connection failing +// is a network error, a connection that ended mid-conversation (a store killed rather than +// stopped gives a bare unexpected EOF), a timeout, or anything pgx marks safe to retry. The server +// saying "not now" is a connection exception (class 08) or it shutting down or starting up (57P01, +// 57P02, 57P03). Everything else is an answer, and asking again gets the same one: a wrong +// password, a database that does not exist, a statement cancelled, a constraint. Those are +// checked first, even inside a failed connection, because a connection that failed on a wrong +// password would otherwise read as a network fault and be asked again for ever. 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") + switch pg.Code { + case "57P01", "57P02", "57P03": + return true + } + return len(pg.Code) == 5 && pg.Code[:2] == "08" } - return false + if errors.Is(err, io.ErrUnexpectedEOF) || errors.Is(err, io.EOF) { + return true + } + if pgconn.Timeout(err) || pgconn.SafeToRetry(err) { + return true + } + var network net.Error + if errors.As(err, &network) { + return true + } + var connect *pgconn.ConnectError + return errors.As(err, &connect) } diff --git a/internal/inventory/reachable_test.go b/internal/inventory/reachable_test.go index f25c830..0454358 100644 --- a/internal/inventory/reachable_test.go +++ b/internal/inventory/reachable_test.go @@ -1,28 +1,76 @@ package inventory import ( + "context" "errors" "fmt" + "io" + "os" + "strings" "testing" + "time" "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) { +// The store restarting or dying is "not now"; an answer the store gave is not (issue 082). +func TestAStoreThatIsGoneForNowIsUnreachableAndAnAnswerIsNot(t *testing.T) { for _, c := range []struct { + what string 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}, + {"starting up", &pgconn.PgError{Code: "57P03"}, true}, + {"stopped by its administrator, wrapped", fmt.Errorf("recording: %w", &pgconn.PgError{Code: "57P01"}), true}, + {"crashed", &pgconn.PgError{Code: "57P02"}, true}, + {"a connection exception", &pgconn.PgError{Code: "08006"}, true}, + {"killed under a live connection", io.ErrUnexpectedEOF, true}, + {"killed, wrapped", fmt.Errorf("recording: %w", io.ErrUnexpectedEOF), true}, + {"a statement cancelled", &pgconn.PgError{Code: "57014"}, false}, + {"the database dropped", &pgconn.PgError{Code: "57P04"}, false}, + {"a wrong password", &pgconn.PgError{Code: "28P01"}, false}, + {"a constraint", &pgconn.PgError{Code: "23503"}, false}, + {"a plain error", errors.New("a report named no node"), false}, + {"nothing", nil, false}, } { if got := Unreachable(c.err); got != c.want { - t.Errorf("Unreachable(%v) = %v, want %v", c.err, got, c.want) + t.Errorf("%s: Unreachable(%v) = %v, want %v", c.what, c.err, got, c.want) } } } + +// A connection refused is the shape a store that is down gives — a real one, to a port on +// loopback nothing listens on. +func TestAStoreThatRefusesTheConnectionIsUnreachable(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _, err := pgconn.Connect(ctx, "postgres://nobody@127.0.0.1:1/nothing?sslmode=disable&connect_timeout=2") + if err == nil { + t.Skip("something listens on port 1") + } + if !Unreachable(err) { + t.Fatalf("a refused connection is not unreachable: %T %v", err, err) + } +} + +// A connection the store refused on a wrong password is an answer, even though it arrives as a +// failed connection — asked again for ever, it would hold every message behind it. +func TestAWrongPasswordIsAnAnswerNotAnOutage(t *testing.T) { + dsn := os.Getenv("MESH_TEST_POSTGRES") + if dsn == "" { + t.Skip("MESH_TEST_POSTGRES is not set") + } + wrong := strings.Replace(dsn, "postgres:check@", "postgres:not-the-password@", 1) + if wrong == dsn { + t.Skip("the test store's address does not carry the expected credential") + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _, err := pgconn.Connect(ctx, wrong) + if err == nil { + t.Fatal("a wrong password connected") + } + if Unreachable(err) { + t.Fatalf("a wrong password was taken for an outage and would be retried for ever: %v", err) + } +} diff --git a/internal/link/report_retry_test.go b/internal/link/report_retry_test.go index 7d5dd16..7e2f61e 100644 --- a/internal/link/report_retry_test.go +++ b/internal/link/report_retry_test.go @@ -7,6 +7,7 @@ import ( "io" "log" "testing" + "time" amqp "github.com/rabbitmq/amqp091-go" ) @@ -40,22 +41,64 @@ func TestAReportTheStoreCouldNotTakeIsKeptAndOneItRefusedIsNot(t *testing.T) { notNow := &saidTo{} s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet, again: 1} - s.handleReport(aReport(t, notNow)) + s.handleReport(context.Background(), 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)) + s.handleReport(context.Background(), 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)) + s.handleReport(context.Background(), aReport(t, recorded)) if !recorded.acked || recorded.nacked { t.Fatalf("a recorded report was not acknowledged: %+v", recorded) } } + +// A store that has not come back within the bound is not restarting: the report is let go, loudly, +// rather than holding every enrolment and report behind it for ever (novox/hq issue 082, review). +func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) { + quiet := log.New(io.Discard, "", 0) + s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))}, + log: quiet, again: 1, giveUp: time.Millisecond} + + first := &saidTo{} + s.handleReport(context.Background(), aReport(t, first)) + if !first.requeued { + t.Fatalf("the first failure was not handed back: %+v", first) + } + time.Sleep(5 * time.Millisecond) + later := &saidTo{} + s.handleReport(context.Background(), aReport(t, later)) + if !later.acked || later.nacked { + t.Fatalf("a report past the bound was not let go, so it would hold the queue for ever: %+v", later) + } + // Let go, and forgotten: the same report arriving fresh starts a new clock. + again := &saidTo{} + s.handleReport(context.Background(), aReport(t, again)) + if !again.requeued { + t.Fatalf("a report let go was not forgotten, so its next arrival is given no chance: %+v", again) + } +} + +// Shutting down does not wait out the pause. +func TestAReportBeingTriedAgainDoesNotHoldUpShutdown(t *testing.T) { + quiet := log.New(io.Discard, "", 0) + s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, + log: quiet, again: time.Hour} + ctx, cancel := context.WithCancel(context.Background()) + cancel() + done := make(chan struct{}) + go func() { s.handleReport(ctx, aReport(t, &saidTo{})); close(done) }() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("a cancelled context still waited out the pause") + } +} diff --git a/internal/link/serve.go b/internal/link/serve.go index f0b9039..a3c7f9c 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -53,8 +53,11 @@ type Server struct { 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 + // the broker; zero means TryAgainAfter. giveUp is how long one report is kept trying before + // it is let go; zero means GiveUpAfter. waiting is when each report still trying first failed. + again time.Duration + giveUp time.Duration + waiting map[string]time.Time } // ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and @@ -66,6 +69,12 @@ var ErrTryAgain = errors.New("not now, try again") // a tight loop; 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 report holds the queue. The consumer takes one message at a +// time, so a report being tried again holds every enrolment, build result and other report +// behind it; a store that has not come back in this long is not restarting, and holding the +// mesh's control queue for it would turn one lost report into a mesh that answers nothing. +const GiveUpAfter = 2 * time.Minute + // 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 @@ -278,7 +287,7 @@ func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) { case KeyEnrol: s.handleEnrol(ctx, delivery) case KeyReport: - s.handleReport(delivery) + s.handleReport(ctx, delivery) case KeyAlive: s.handleAlive(delivery) case KeyBuilt: @@ -314,7 +323,7 @@ func (s *Server) handleAlive(delivery amqp.Delivery) { // 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(delivery amqp.Delivery) { +func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { var report Report if err := json.Unmarshal(delivery.Body, &report); err != nil { s.log.Printf("a report could not be read: %v", err) @@ -323,7 +332,7 @@ 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) { + if errors.Is(err, ErrTryAgain) && s.keepTrying(report) { // 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). @@ -332,10 +341,18 @@ func (s *Server) handleReport(delivery amqp.Delivery) { 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) + select { + case <-ctx.Done(): + case <-time.After(again): + } _ = delivery.Nack(false, true) return } + if errors.Is(err, ErrTryAgain) { + s.log.Printf("LOST %s's report of declaration %s: the store has not come back in %s, "+ + "and the control queue cannot wait longer — the node is current but the mesh will "+ + "read it as unanswered until its next push: %v", report.Node, report.Declared, s.giveUpAfter(), err) + } // 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) @@ -350,9 +367,34 @@ func (s *Server) handleReport(delivery amqp.Delivery) { default: s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied)) } + s.stopTrying(report) _ = delivery.Ack(false) } +// keepTrying says whether a report the store could not take is still within the time it may hold +// the queue, starting that clock on its first failure. +func (s *Server) keepTrying(r Report) bool { + if s.waiting == nil { + s.waiting = map[string]time.Time{} + } + key := r.Node + " " + r.Declared + first, seen := s.waiting[key] + if !seen { + s.waiting[key] = time.Now() + return true + } + return time.Since(first) < s.giveUpAfter() +} + +func (s *Server) stopTrying(r Report) { delete(s.waiting, r.Node+" "+r.Declared) } + +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) { reply := EnrolReply{Refusal: "that token cannot be used"}