A report that arrives while the store restarts is kept, not lost (issue 082) #42
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user