Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
114 changes: 97 additions & 17 deletions go/internal/ingest/board_webhook.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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.
Expand Down Expand Up @@ -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,
Expand All @@ -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,
}
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand All @@ -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)
}
Expand All @@ -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.
Expand Down
100 changes: 100 additions & 0 deletions go/internal/ingest/pull_request_sink_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package ingest
import (
"context"
"errors"
"runtime"
"slices"
"sync"
"testing"
Expand Down Expand Up @@ -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()
}
}
}
Loading