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
This commit is contained in:
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
+48
-6
@@ -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"}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user