diff --git a/core/artwork/gate.go b/core/artwork/gate.go index ff77f733d..534bd8ce3 100644 --- a/core/artwork/gate.go +++ b/core/artwork/gate.go @@ -17,6 +17,11 @@ import ( const ( breakerThreshold = 5 breakerProbeAfter = time.Minute + // breakerRecoveries is how many consecutive answers an open breaker needs before it trusts the + // provider again. One is not enough: a provider that is rate-limiting or blocking us still + // answers the occasional request, and closing on the first of those puts the agent straight + // back to full rate, which is what earns the next block. + breakerRecoveries = 3 ) var errBreakerOpen = errors.New("artwork: external circuit breaker open") @@ -44,7 +49,8 @@ type extGate struct { // gate runs a named external step through that agent's rate limiter and circuit breaker. func (w *Worker) gate(name string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) { g := w.gateFor(name) - if !g.breaker.allow() { + allowed, gen := g.breaker.allow() + if !allowed { log.Debug(w.runCtx, "Artwork: Skipping agent, circuit breaker open", "agent", name) return nil, "", errBreakerOpen } @@ -55,7 +61,7 @@ func (w *Worker) gate(name string, f func() (io.ReadCloser, string, error)) (io. } callStart := time.Now() r, path, err := f() - g.breaker.record(name, err) + g.breaker.record(name, gen, err) log.Trace(w.runCtx, "Artwork: External agent call", "agent", name, "hit", r != nil, "limiterWait", callStart.Sub(waitStart), "elapsed", time.Since(callStart), err) return r, path, err @@ -84,41 +90,65 @@ type breaker struct { mu sync.Mutex failures int openedAt time.Time + // recoveries counts consecutive good answers while open; a single failure discards them. + recoveries int + // generation identifies the current open episode, so an answer from a call admitted before + // the breaker opened cannot be mistaken for evidence that it has recovered. + generation int } func newBreaker() *breaker { return &breaker{} } -func (b *breaker) allow() bool { +// allow reports whether a call may proceed, and the open episode it was admitted under: zero +// when the breaker was closed, the current generation when admitted as a half-open probe. +func (b *breaker) allow() (bool, int) { b.mu.Lock() defer b.mu.Unlock() if b.failures < breakerThreshold { - return true + return true, 0 } if time.Since(b.openedAt) >= breakerProbeAfter { b.openedAt = time.Now() // start a fresh probe window so only one caller passes - return true + return true, b.generation } - return false + return false, 0 } -func (b *breaker) record(name string, err error) { +func (b *breaker) record(name string, gen int, err error) { // A cancelled run says nothing about the provider, so it neither counts nor clears. if errors.Is(err, context.Canceled) { return } b.mu.Lock() defer b.mu.Unlock() - if !isTransientExternal(err) { - if b.failures >= breakerThreshold { - log.Info("Artwork: Circuit breaker closed for agent", "agent", name) + if isTransientExternal(err) { + b.recoveries = 0 + b.failures++ + if b.failures == breakerThreshold { + b.openedAt = time.Now() + b.generation++ + log.Warn("Artwork: Circuit breaker opened for agent", "agent", name, + "consecutiveFailures", b.failures, "probeAfter", breakerProbeAfter, err) } + return + } + if b.failures < breakerThreshold { b.failures = 0 return } - b.failures++ - if b.failures == breakerThreshold { - b.openedAt = time.Now() - log.Warn("Artwork: Circuit breaker opened for agent", "agent", name, - "consecutiveFailures", b.failures, "probeAfter", breakerProbeAfter, err) + // Only a probe from this open episode is evidence of recovery. The worker drains concurrently, + // so answers keep arriving from calls admitted before the breaker opened; counting those would + // close it with no probe interval elapsed, which is the burst this exists to prevent. + if gen == 0 || gen != b.generation { + return } + // A not-found counts because the provider did answer, but on its own it is thin evidence that + // a provider which just blocked us is well. + b.recoveries++ + if b.recoveries < breakerRecoveries { + return + } + log.Info("Artwork: Circuit breaker closed for agent", "agent", name, + "consecutiveAnswers", b.recoveries) + b.failures, b.recoveries = 0, 0 } diff --git a/core/artwork/gate_test.go b/core/artwork/gate_test.go new file mode 100644 index 000000000..abe723508 --- /dev/null +++ b/core/artwork/gate_test.go @@ -0,0 +1,43 @@ +package artwork + +import ( + "errors" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// allowed drops the generation token when a caller only cares about admission. +func allowed(b *breaker) bool { ok, _ := b.allow(); return ok } + +var _ = Describe("breaker", func() { + // The worker drains concurrently, so when the breaker opens there are already calls past + // allow(), queued in the rate limiter or waiting on a response. Their answers arrive + // afterwards. Counting those as recovery closes the breaker with no probe interval elapsed, + // which is the burst the ramp exists to prevent. No clock is involved: the race is an + // ordering, so it is reproduced by making the calls in the order concurrency produces. + It("ignores answers from calls admitted before it opened", func() { + b := newBreaker() + + // A batch clears allow() while the breaker is still closed. + for range breakerThreshold + breakerRecoveries { + ok, gen := b.allow() + Expect(ok).To(BeTrue()) + Expect(gen).To(BeZero(), "admitted with the breaker closed, so not a probe") + } + + // The fast failures in that batch open it. + for range breakerThreshold { + b.record("agentA", 0, errors.New("blocked")) + } + Expect(allowed(b)).To(BeFalse(), "breaker is open") + + // The slower answers from the same batch land now. + for range breakerRecoveries { + b.record("agentA", 0, nil) + } + + Expect(allowed(b)).To(BeFalse(), + "answers from calls admitted before the breaker opened must not close it") + }) +}) diff --git a/core/artwork/worker_timing_test.go b/core/artwork/worker_timing_test.go index f7315b20a..63f12d03c 100644 --- a/core/artwork/worker_timing_test.go +++ b/core/artwork/worker_timing_test.go @@ -20,24 +20,63 @@ func TestArtworkBreakerHalfOpen(t *testing.T) { b := newBreaker() for range breakerThreshold { - b.record("agentA", errors.New("boom")) + b.record("agentA", 0, errors.New("boom")) } - g.Expect(b.allow()).To(BeFalse(), "breaker opens after consecutive errors") + g.Expect(allowed(b)).To(BeFalse(), "breaker opens after consecutive errors") time.Sleep(breakerProbeAfter - time.Nanosecond) - g.Expect(b.allow()).To(BeFalse(), "still open before the probe interval") + g.Expect(allowed(b)).To(BeFalse(), "still open before the probe interval") time.Sleep(time.Nanosecond) - g.Expect(b.allow()).To(BeTrue(), "half-open: one probe is granted") - g.Expect(b.allow()).To(BeFalse(), "only a single probe per interval") + ok, gen := b.allow() + g.Expect(ok).To(BeTrue(), "half-open: one probe is granted") + g.Expect(gen).ToNot(BeZero(), "a probe carries the open episode it belongs to") + g.Expect(allowed(b)).To(BeFalse(), "only a single probe per interval") - b.record("agentA", errors.New("boom")) // probe fails -> stay open + b.record("agentA", gen, errors.New("boom")) // probe fails -> stay open time.Sleep(breakerProbeAfter) - g.Expect(b.allow()).To(BeTrue(), "another probe after the next interval") + ok, gen = b.allow() + g.Expect(ok).To(BeTrue(), "another probe after the next interval") - b.record("agentA", nil) // probe succeeds -> close - g.Expect(b.allow()).To(BeTrue(), "closed breaker admits freely") - g.Expect(b.allow()).To(BeTrue()) + // One good answer must not reopen the floodgates: closing here is what let a burst out at + // full rate and got the provider to escalate from throttling to blocking. + b.record("agentA", gen, nil) + g.Expect(allowed(b)).To(BeFalse(), "a single good answer does not close the breaker") + + for range breakerRecoveries - 1 { + time.Sleep(breakerProbeAfter) + ok, gen = b.allow() + g.Expect(ok).To(BeTrue()) + b.record("agentA", gen, nil) + } + g.Expect(allowed(b)).To(BeTrue(), "closed breaker admits freely") + g.Expect(allowed(b)).To(BeTrue()) + }) +} + +// The failure seen in production: while an agent was blocked, the occasional answer it did serve +// reset the breaker, releasing a burst that immediately re-tripped it. Open and closed pairs were +// seconds apart, over and over. +func TestArtworkBreakerDoesNotCloseOnAnIsolatedAnswer(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + g := NewWithT(t) + b := newBreaker() + open := func() { + for range breakerThreshold { + b.record("agentA", 0, errors.New("blocked")) + } + } + open() + + // A not-found is an answer, so it counts toward recovery, but never on its own. + for range breakerRecoveries * 2 { + time.Sleep(breakerProbeAfter) + ok, gen := b.allow() + g.Expect(ok).To(BeTrue(), "one probe per interval") + b.record("agentA", gen, agents.ErrNotFound) + b.record("agentA", 0, errors.New("blocked")) // the very next call is refused again + g.Expect(allowed(b)).To(BeFalse(), "an answer between failures must not close the breaker") + } }) }