Merge pull request 'A report that arrives while the store restarts is kept, not lost (issue 082)' (#42) from multiple-fixes into main

This commit was merged in pull request #42.
This commit is contained in:
2026-09-22 13:53:13 +02:00
5 changed files with 309 additions and 3 deletions
+52
View File
@@ -0,0 +1,52 @@
package inventory
import (
"errors"
"io"
"net"
"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; 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 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 pg *pgconn.PgError
if errors.As(err, &pg) {
switch pg.Code {
case "57P01", "57P02", "57P03":
return true
}
return len(pg.Code) == 5 && pg.Code[:2] == "08"
}
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)
}
+76
View File
@@ -0,0 +1,76 @@
package inventory
import (
"context"
"errors"
"fmt"
"io"
"os"
"strings"
"testing"
"time"
"github.com/jackc/pgx/v5/pgconn"
)
// 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
}{
{"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("%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)
}
}
+8 -1
View File
@@ -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")
}
+104
View File
@@ -0,0 +1,104 @@
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 }
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(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(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(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")
}
}
+69 -2
View File
@@ -52,8 +52,29 @@ 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. 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
// 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
// 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
@@ -266,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:
@@ -302,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)
@@ -311,6 +332,27 @@ 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) && 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).
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)
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)
@@ -325,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"}