diff --git a/go/internal/ingest/board_webhook.go b/go/internal/ingest/board_webhook.go index cc8bd5626..c92aa2dfd 100644 --- a/go/internal/ingest/board_webhook.go +++ b/go/internal/ingest/board_webhook.go @@ -17,6 +17,7 @@ import ( "log/slog" "strings" "sync/atomic" + "time" "github.com/RigelBuild/compass/go/internal/forge" compassv1internal "github.com/RigelBuild/compass/go/internal/gen/compass/v1" @@ -30,10 +31,15 @@ var boardWebhookDrops = expvar.NewInt("compass_board_webhook_drops") // defaultBoardQueueSize is the bounded drain-queue depth. Sized to absorb a // normal edit-storm burst between drains; a sustained overflow (the drain paused -// on ErrBudgetExhausted while events keep arriving) drops with the metric+Warn -// and is healed by the T3 reconciler. +// on ErrBudgetExhausted while events keep arriving) drops with the metric+Warn. +// The sweep heals a dropped issue; a dropped CI-only PR change waits for the PR's +// next event. const defaultBoardQueueSize = 1024 +// defaultBoardRetryWait re-drains held PR keys when the budget error gave no +// reset hint. +const defaultBoardRetryWait = time.Minute + // issueHydrator is the conditional point-read seam (satisfied structurally by // *forge.GitHub via GetIssueConditional, notify_reader.go:129). Defined locally // so this package imports no concrete forge client on the hydrate path. @@ -82,6 +88,12 @@ type BoardWebhookArm struct { targets TargetChecker log *slog.Logger dropped atomic.Int64 + + // pending holds PR keys a budget pause left un-hydrated. The sweep cannot + // heal a CI-only change, so the drain retries them. Drain goroutine only. + pending map[boardCoord]struct{} + maxPending int + after func(time.Duration) <-chan time.Time } // NewBoardWebhookArm returns an arm that hydrates each accepted event through h, @@ -103,6 +115,10 @@ func NewBoardWebhookArm(h issueHydrator, ing *Ingester, targets TargetChecker, c numbers: cfg.PullNumbers, targets: targets, log: log, + + pending: map[boardCoord]struct{}{}, + maxPending: size, + after: time.After, } } @@ -130,18 +146,38 @@ func (a *BoardWebhookArm) Enqueue(_ context.Context, ev forge.ForgeEvent) { // Run drains the queue until ctx is cancelled — then it returns nil (clean // shutdown, driver.go:95-99 idiom). Each drain COALESCES per coordinate before -// hydrating, so an N-event burst on one issue costs one GET. +// hydrating, so an N-event burst on one issue costs one GET. After a budget +// pause it re-drains the held PR keys once the gate's reset hint passes. func (a *BoardWebhookArm) Run(ctx context.Context) error { + var retry <-chan time.Time for { select { case <-ctx.Done(): return nil case ev := <-a.queue: - a.drainBatch(ctx, ev) + if wait, paused := a.drainBatch(ctx, ev); paused && retry == nil { + retry = a.retryTimer(wait) + } + case <-retry: + retry = nil + if wait, paused := a.retryPending(ctx); paused { + retry = a.retryTimer(wait) + } } } } +// retryTimer arms the re-drain for held PR keys; nil when none are held. +func (a *BoardWebhookArm) retryTimer(wait time.Duration) <-chan time.Time { + if len(a.pending) == 0 { + return nil + } + if wait <= 0 { + wait = defaultBoardRetryWait + } + return a.after(wait) +} + // boardRelevant reports whether an event changes what the board shows. Issue // events count on OPENED, STATE or UPDATE; an issue COMMENT would burn one // hydrate per comment for nothing shown. Every PR event counts, since reviews @@ -173,11 +209,12 @@ func (a *BoardWebhookArm) boardRelevant(ev forge.ForgeEvent) bool { // drainBatch coalesces the event that woke the drain and every other queued // event into distinct keys: normalized repo, kind and number, or head SHA for a // CHECKS event with no number. An edit storm and a mixed-case duplicate both -// collapse to one key. It then hydrates + sinks each key once, in arrival order. -func (a *BoardWebhookArm) drainBatch(ctx context.Context, first forge.ForgeEvent) { +// collapse to one key. A key already held for retry is left to the timer, so a +// held PR whose budget is still out never pauses a fresh batch. It then hydrates +// + sinks each key once, in arrival order, and reports a pause. +func (a *BoardWebhookArm) drainBatch(ctx context.Context, first forge.ForgeEvent) (time.Duration, bool) { seen := map[boardCoord]struct{}{} var order []boardCoord - add := func(ev forge.ForgeEvent) { c := boardCoord{repo: normalizeBoardRepo(ev.Repo), kind: ev.Kind, number: ev.Number} if ev.Number == 0 { @@ -186,6 +223,9 @@ func (a *BoardWebhookArm) drainBatch(ctx context.Context, first forge.ForgeEvent if _, ok := seen[c]; ok { return } + if _, held := a.pending[c]; held { + return + } seen[c] = struct{}{} order = append(order, c) } @@ -196,29 +236,69 @@ func (a *BoardWebhookArm) drainBatch(ctx context.Context, first forge.ForgeEvent case ev := <-a.queue: add(ev) default: - goto process + return a.process(ctx, order, seen) } } +} -process: - for _, c := range order { +// retryPending re-drains the PR keys held by an earlier budget pause. +func (a *BoardWebhookArm) retryPending(ctx context.Context) (time.Duration, bool) { + order, seen := a.takePending() + return a.process(ctx, order, seen) +} + +// takePending empties the held PR keys into a fresh batch. +func (a *BoardWebhookArm) takePending() ([]boardCoord, map[boardCoord]struct{}) { + seen := make(map[boardCoord]struct{}, len(a.pending)) + order := make([]boardCoord, 0, len(a.pending)) + for c := range a.pending { + seen[c] = struct{}{} + order = append(order, c) + } + clear(a.pending) + return order, seen +} + +// hold keeps the batch's un-hydrated PR keys for a retry; issue keys are left to +// the sweep, which sees their updated_at move. Past maxPending a key is dropped. +func (a *BoardWebhookArm) hold(rest []boardCoord) { + for _, c := range rest { + if c.kind != compassv1internal.ForgeArtifactKind_FORGE_ARTIFACT_KIND_PULL_REQUEST { + continue + } + if _, ok := a.pending[c]; !ok && len(a.pending) >= a.maxPending { + a.dropped.Add(1) + boardWebhookDrops.Add(1) + continue + } + a.pending[c] = struct{}{} + } +} + +// process hydrates each key once. A budget error holds the remaining PR keys +// and returns the gate's reset hint. +func (a *BoardWebhookArm) process(ctx context.Context, order []boardCoord, seen map[boardCoord]struct{}) (time.Duration, bool) { + for i, c := range order { if err := a.resolveAndSink(ctx, c, seen); err != nil { if errors.Is(err, forge.ErrBudgetExhausted) { - // Budget exhausted pauses the drain: abandon the rest of this - // batch (the reconciler heals the un-hydrated coordinates) and - // resume once the client gate reopens. - a.log.WarnContext(ctx, "board webhook: budget exhausted, pausing drain (reconciler heals)", - "repo", c.repo, "number", c.number, "head_sha", c.headSHA) - return + a.hold(order[i:]) + var wait time.Duration + if rle, ok := errors.AsType[*forge.RateLimitError](err); ok { + wait = rle.RetryAfter + } + a.log.WarnContext(ctx, "board webhook: budget exhausted, pausing drain", + "repo", c.repo, "number", c.number, "head_sha", c.headSHA, "held", len(a.pending)) + return wait, true } // Per-event errors log-and-continue (driver.go:96-98 idiom). a.log.WarnContext(ctx, "board webhook: hydrate/sink failed (isolated)", "repo", c.repo, "number", c.number, "head_sha", c.headSHA, "err", err) } if ctx.Err() != nil { - return + return 0, false } } + return 0, false } // resolveAndSink maps a head-SHA CHECKS key to its PR number, then hydrates. diff --git a/go/internal/ingest/pull_request_sink_test.go b/go/internal/ingest/pull_request_sink_test.go index a71bf6278..9904ee061 100644 --- a/go/internal/ingest/pull_request_sink_test.go +++ b/go/internal/ingest/pull_request_sink_test.go @@ -3,6 +3,7 @@ package ingest import ( "context" "errors" + "runtime" "slices" "sync" "testing" @@ -446,3 +447,102 @@ func TestReconcileBackfillRetriesAfterRowFailure(t *testing.T) { t.Fatal("repo not marked after the clean retry") } } + +// TestArmBudgetPauseHoldsPRKeysForRetry: a budget pause keeps the batch's +// un-hydrated PR keys, including a resolved CHECKS key, and the retry hydrates +// them once; issue keys are left to the sweep. +func TestArmBudgetPauseHoldsPRKeysForRetry(t *testing.T) { + pulls := &fakePulls{errFor: map[uint64]error{9: &forge.RateLimitError{RetryAfter: 30 * time.Second}}} + arm, sink := newPRArm(t, pulls, fakeNumbers{"abc": 12}) + arm.Enqueue(context.Background(), prEvent(9, changeUpdate)) + arm.Enqueue(context.Background(), checksEvent()) + arm.Enqueue(context.Background(), issueEvent("owner/repo", 4, changeUpdate)) + wait, paused := arm.drainBatch(context.Background(), <-arm.queue) + if !paused || wait != 30*time.Second { + t.Fatalf("pause = %v wait %v, want paused with the 30s hint", paused, wait) + } + if len(arm.pending) != 2 { + t.Fatalf("held %v, want the two PR keys", arm.pending) + } + + pulls.errFor = nil + if _, paused := arm.retryPending(context.Background()); paused { + t.Fatal("retry paused with the gate open") + } + got := pulls.readNumbers() + slices.Sort(got) + if !slices.Equal(got, []uint64{9, 9, 12}) || len(sink.got) != 2 || len(arm.pending) != 0 { + t.Fatalf("reads = %v sank %d held %d, want [9 9 12], 2 and 0", got, len(sink.got), len(arm.pending)) + } +} + +// TestArmHeldKeyNeverPausesFreshBatch: while a PR key is held, a new event for +// it is left to the timer and a fresh batch hydrates without touching it. +func TestArmHeldKeyNeverPausesFreshBatch(t *testing.T) { + pulls := &fakePulls{errFor: map[uint64]error{9: forge.ErrBudgetExhausted}} + arm, sink := newPRArm(t, pulls, nil) + arm.Enqueue(context.Background(), prEvent(9, changeUpdate)) + drainAll(context.Background(), arm) + arm.Enqueue(context.Background(), prEvent(9, changeComment)) + arm.Enqueue(context.Background(), prEvent(10, changeUpdate)) + drainAll(context.Background(), arm) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{9, 10}) || len(sink.got) != 1 { + t.Fatalf("reads = %v sank %d, want [9 10] and 1", got, len(sink.got)) + } + if _, held := arm.pending[boardCoord{repo: "owner/repo", kind: boardKindPR, number: 9}]; !held || len(arm.pending) != 1 { + t.Fatalf("held = %v, want only PR 9", arm.pending) + } +} + +// TestArmRunRetriesAfterResetHint: Run waits the gate's reset hint, then +// re-drains the held keys with no new event. +func TestArmRunRetriesAfterResetHint(t *testing.T) { + pulls := &fakePulls{errFor: map[uint64]error{9: &forge.RateLimitError{RetryAfter: 45 * time.Second}}} + arm, sink := newPRArm(t, pulls, nil) + waits := make(chan time.Duration, 1) + fire := make(chan time.Time) + arm.after = func(d time.Duration) <-chan time.Time { + waits <- d + return fire + } + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- arm.Run(ctx) }() + + defer cancel() + arm.Enqueue(ctx, prEvent(9, changeUpdate)) + select { + case d := <-waits: + if d != 45*time.Second { + t.Fatalf("retry wait = %v, want the 45s hint", d) + } + case <-time.After(5 * time.Second): + t.Fatal("Run armed no retry after the budget pause") + } + pulls.mu.Lock() + pulls.errFor = nil + pulls.mu.Unlock() + fire <- time.Time{} + waitForReads(t, pulls, 2) + cancel() + if err := <-done; err != nil { + t.Fatalf("Run = %v", err) + } + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{9, 9}) || len(sink.got) != 1 { + t.Fatalf("reads = %v sank %d, want [9 9] and 1", got, len(sink.got)) + } +} + +// waitForReads waits until pulls has seen n reads; the retry runs on Run's goroutine. +func waitForReads(t *testing.T, pulls *fakePulls, n int) { + t.Helper() + deadline := time.After(5 * time.Second) + for len(pulls.readNumbers()) < n { + select { + case <-deadline: + t.Fatalf("reads = %v, want %d", pulls.readNumbers(), n) + default: + runtime.Gosched() + } + } +}