diff --git a/go/internal/ingest/board_reconcile.go b/go/internal/ingest/board_reconcile.go index 66b212215..c19559e5b 100644 --- a/go/internal/ingest/board_reconcile.go +++ b/go/internal/ingest/board_reconcile.go @@ -37,8 +37,18 @@ type BoardStore interface { // StoreRepoWatermark persists the watermark + ETag AFTER the repo's rows // sank (advance-after-sink, the idempotency invariant). StoreRepoWatermark(ctx context.Context, repo string, mark time.Time, etag string) error + // PullRequestUpdatedAt returns the stored PR row's forge updated_at; ok is + // false when the PR was never stored. + PullRequestUpdatedAt(ctx context.Context, repo string, number uint64) (time.Time, bool, error) + // PRsBackfilledAt returns when the repo's PR backfill ran; ok is false if never. + PRsBackfilledAt(ctx context.Context, repo string) (time.Time, bool, error) + // MarkPRsBackfilled records that the repo's PR backfill ran at at. + MarkPRsBackfilled(ctx context.Context, repo string, at time.Time) error } +// prBackfillWindow bounds which closed PRs a cold start or backfill hydrates. +const prBackfillWindow = 30 * 24 * time.Hour + // updatedLister is the conditional updated-order list surface, satisfied by // forge.GitHub.ListUpdatedIssues at T5. LOCAL + structural: this package never // imports the concrete provider. @@ -55,6 +65,8 @@ type BoardReconcileConfig struct { Pace time.Duration // Log is the sweep logger; nil uses slog.Default(). Log *slog.Logger + // Pulls hydrates listed PR rows and runs the backfill; nil skips PR rows. + Pulls *PullRequestHydrator } // BoardReconciler drives the backstop sweep over the enabled-repo set, listing @@ -65,6 +77,7 @@ type BoardReconciler struct { lister updatedLister ingester *Ingester store BoardStore + pulls *PullRequestHydrator backstop time.Duration pace time.Duration log *slog.Logger @@ -89,6 +102,7 @@ func NewBoardReconciler(l updatedLister, ing *Ingester, st BoardStore, cfg Board lister: l, ingester: ing, store: st, + pulls: cfg.Pulls, backstop: backstop, pace: pace, log: log, @@ -149,23 +163,11 @@ func (rc *BoardReconciler) sweep(ctx context.Context) { } } -// reconcileRepo conditionally lists one repo's updated-order issues since its +// reconcileRepo conditionally lists one repo's updated-order rows since its // stored watermark, sinks them, and advances the watermark AFTER the sink. A 304 // costs no sink and leaves the watermark untouched (the stored ETag remains the -// truth). A zero/absent watermark lists everything (cold-start/backfill). -// -// Rows are sunk one at a time so a poison row is isolated: it is skipped and -// counted, the rest sink, and the watermark advances only past the rows that -// actually sank — the re-walk window stays bounded rather than growing every -// sweep behind a pinned watermark. -// -// Tradeoff (deliberate, do not "fix"): the watermark advances to the max -// timestamp over the HEALTHY rows, so a row that fails only TRANSIENTLY while -// co-batched with a newer healthy row is left below the advanced watermark and -// is dropped until its next forge update re-lists it. Capping the advance below -// the min failed-row timestamp instead would re-introduce the poison-pin -// livelock this isolation exists to prevent — the bounded re-walk window is the -// correct resolution of that tradeoff. +// truth). A zero/absent watermark lists everything (cold-start/backfill). The +// PR backfill runs after the walk, whether or not it was a 304. func (rc *BoardReconciler) reconcileRepo(ctx context.Context, repo string) error { since, etag, err := rc.store.LoadRepoWatermark(ctx, repo) if err != nil { @@ -175,46 +177,152 @@ func (rc *BoardReconciler) reconcileRepo(ctx context.Context, repo string) error if err != nil { return err } - if res.NotModified { - return nil // content unchanged — the stored watermark + ETag stay the truth + pullsFailed := 0 + if !res.NotModified { + if pullsFailed, err = rc.sinkRows(ctx, repo, since, res); err != nil { + return err + } } + return rc.backfillPRs(ctx, repo, since, pullsFailed) +} +// sinkRows sinks one listed window, isolating each row so a poison row is +// skipped and counted, and advances the watermark only past the rows that sank. +// A PR hydrate that runs out of budget is not poison: the watermark stays at +// since with no ETag, so the next sweep re-lists the whole window. It returns +// the count of PR rows that failed to hydrate. +// +// Tradeoff (deliberate, do not "fix"): the watermark advances to the max +// timestamp over the HEALTHY rows, so a row that fails only TRANSIENTLY while +// co-batched with a newer healthy row is left below the advanced watermark and +// is dropped until its next forge update re-lists it. Capping the advance below +// the min failed-row timestamp instead would re-introduce the poison-pin +// livelock this isolation exists to prevent. +func (rc *BoardReconciler) sinkRows(ctx context.Context, repo string, since time.Time, res forge.ConditionalResult[forge.UpdatedRows]) (int, error) { var maxMark time.Time - poison := 0 + poison, pullsFailed := 0, 0 for _, row := range res.V.Issues { if ctx.Err() != nil { - return ctx.Err() + return 0, ctx.Err() } if serr := rc.ingester.IngestIssues(ctx, repo, []forge.Issue{row}); serr != nil { - // Isolate the poison row: skip + count it, keep sinking the rest, - // and do NOT let its timestamp advance the watermark. poison++ rc.log.WarnContext(ctx, "board reconcile: row sink failed (isolated)", "repo", repo, "number", row.Number, "error", serr) continue } - if row.UpdatedAt.After(maxMark) { - maxMark = row.UpdatedAt + maxMark = later(maxMark, row.UpdatedAt) + } + if rc.pulls != nil { + for _, row := range res.V.Pulls { + if ctx.Err() != nil { + return 0, ctx.Err() + } + if !since.IsZero() || inBackfillWindow(row, time.Now()) { + if herr := rc.hydrateIfNewer(ctx, repo, row); herr != nil { + if errors.Is(herr, forge.ErrBudgetExhausted) { + return 0, rc.abortWindow(ctx, repo, since, herr) + } + poison++ + pullsFailed++ + rc.log.WarnContext(ctx, "board reconcile: pull request hydrate failed (isolated)", + "repo", repo, "number", row.Number, "error", herr) + continue + } + } + maxMark = later(maxMark, row.UpdatedAt) } } // Nothing sank (empty list, or every row poison): leave the watermark where - // it is so a healthy row is re-listed next sweep. The window is bounded — it - // never accumulates rows behind a pinned watermark. + // it is so a healthy row is re-listed next sweep. if maxMark.IsZero() { - return nil + return pullsFailed, nil } - - // Advance-after-sink. On a clean sweep carry the fresh list ETag so the next - // sweep can 304; when a poison row was skipped, drop the ETag so the next - // sweep re-lists unconditionally (a fresh ETag would 304-suppress the - // still-unsunk row's retry) — bounded to the rows at/after maxMark. + // On a clean sweep carry the fresh list ETag so the next sweep can 304; when + // a poison row was skipped, drop it so the next sweep re-lists. storeETag := res.ETag if poison > 0 { storeETag = "" } - if serr := rc.store.StoreRepoWatermark(ctx, repo, maxMark, storeETag); serr != nil { - return serr + return pullsFailed, rc.store.StoreRepoWatermark(ctx, repo, maxMark, storeETag) +} + +// abortWindow keeps the watermark at since and clears the ETag, so the next +// sweep re-lists every row of the window instead of getting a 304, then +// returns cause. +func (rc *BoardReconciler) abortWindow(ctx context.Context, repo string, since time.Time, cause error) error { + if err := rc.store.StoreRepoWatermark(ctx, repo, since, ""); err != nil { + return errors.Join(cause, err) + } + return cause +} + +// hydrateIfNewer hydrates a listed PR row unless the stored row is as new. +func (rc *BoardReconciler) hydrateIfNewer(ctx context.Context, repo string, row forge.UpdatedPull) error { + stored, ok, err := rc.store.PullRequestUpdatedAt(ctx, repo, row.Number) + if err != nil { + return err + } + if ok && !row.UpdatedAt.After(stored) { + return nil + } + return rc.pulls.Hydrate(ctx, repo, row.Number) +} + +// backfillPRs runs once per repo: it hydrates every open PR and every PR row +// updated in the window before the watermark, then marks the repo. On a cold +// start the main walk covered that window, so its PR failures hold the mark. +func (rc *BoardReconciler) backfillPRs(ctx context.Context, repo string, since time.Time, walkFailed int) error { + if rc.pulls == nil { + return nil + } + if _, done, err := rc.store.PRsBackfilledAt(ctx, repo); err != nil || done { + return err + } + open, err := rc.pulls.listOpen(ctx, repo) + if err != nil { + return err + } + rows := open + if !since.IsZero() { + win, err := rc.lister.ListUpdatedIssues(ctx, repo, since.Add(-prBackfillWindow), "") + if err != nil { + return err + } + rows = append(rows, win.V.Pulls...) + } + failed := walkFailed + for _, row := range rows { + if ctx.Err() != nil { + return ctx.Err() + } + if err := rc.hydrateIfNewer(ctx, repo, row); err != nil { + if errors.Is(err, forge.ErrBudgetExhausted) { + return err + } + failed++ + rc.log.WarnContext(ctx, "board reconcile: backfill hydrate failed (isolated)", + "repo", repo, "number", row.Number, "error", err) + } + } + // A failed row would never be re-listed, so retry the pass next sweep; the + // updated_at gate makes the rows that sank cheap to skip. + if failed > 0 { + return nil + } + return rc.store.MarkPRsBackfilled(ctx, repo, time.Now()) +} + +// inBackfillWindow reports whether a cold start hydrates the row: open, or +// updated within the backfill window. +func inBackfillWindow(row forge.UpdatedPull, now time.Time) bool { + return row.State == "open" || row.UpdatedAt.After(now.Add(-prBackfillWindow)) +} + +func later(a, b time.Time) time.Time { + if b.After(a) { + return b } - return nil + return a } diff --git a/go/internal/ingest/board_reconcile_test.go b/go/internal/ingest/board_reconcile_test.go index be9f6500f..d0bdafb94 100644 --- a/go/internal/ingest/board_reconcile_test.go +++ b/go/internal/ingest/board_reconcile_test.go @@ -63,10 +63,27 @@ type fakeBoardStore struct { loadErr error storeErr error storeCalls []storedMark + prUpdated map[uint64]time.Time + backfilled map[string]time.Time +} + +func (s *fakeBoardStore) PullRequestUpdatedAt(_ context.Context, _ string, number uint64) (time.Time, bool, error) { + at, ok := s.prUpdated[number] + return at, ok, nil +} + +func (s *fakeBoardStore) PRsBackfilledAt(_ context.Context, repo string) (time.Time, bool, error) { + at, ok := s.backfilled[repo] + return at, ok, nil +} + +func (s *fakeBoardStore) MarkPRsBackfilled(_ context.Context, repo string, at time.Time) error { + s.backfilled[repo] = at + return nil } func newBoardStore(repos ...string) *fakeBoardStore { - return &fakeBoardStore{repos: repos, marks: map[string]storedMark{}} + return &fakeBoardStore{repos: repos, marks: map[string]storedMark{}, prUpdated: map[uint64]time.Time{}, backfilled: map[string]time.Time{}} } func (s *fakeBoardStore) ListEnabledRepos(_ context.Context) ([]string, error) { diff --git a/go/internal/ingest/board_webhook.go b/go/internal/ingest/board_webhook.go index aa8d34716..cc8bd5626 100644 --- a/go/internal/ingest/board_webhook.go +++ b/go/internal/ingest/board_webhook.go @@ -55,13 +55,20 @@ type BoardArmConfig struct { QueueSize int // Log is the arm logger; nil uses slog.Default(). Log *slog.Logger + // Pulls hydrates PR events onto the board; nil drops every PR event. + Pulls *PullRequestHydrator + // PullNumbers maps a CHECKS event's head SHA to its PR; nil drops those events. + PullNumbers PullNumberResolver } -// boardCoord is the (repo, number) coalescing key: the drain collapses every -// queued event for one coordinate to a single hydrate GET. +// boardCoord is the coalescing key: the drain collapses every queued event for +// one artifact to a single hydrate. A CHECKS event keys on its head SHA until +// the drain resolves it to a PR number. type boardCoord struct { - repo string - number uint64 + repo string + kind compassv1internal.ForgeArtifactKind + number uint64 + headSHA string } // BoardWebhookArm consumes board-relevant forge events from the webhook ingress @@ -70,6 +77,8 @@ type BoardWebhookArm struct { queue chan forge.ForgeEvent hydrator issueHydrator ing *Ingester + pulls *PullRequestHydrator + numbers PullNumberResolver targets TargetChecker log *slog.Logger dropped atomic.Int64 @@ -90,6 +99,8 @@ func NewBoardWebhookArm(h issueHydrator, ing *Ingester, targets TargetChecker, c queue: make(chan forge.ForgeEvent, size), hydrator: h, ing: ing, + pulls: cfg.Pulls, + numbers: cfg.PullNumbers, targets: targets, log: log, } @@ -100,12 +111,11 @@ func NewBoardWebhookArm(h issueHydrator, ing *Ingester, targets TargetChecker, c func (a *BoardWebhookArm) Dropped() int64 { return a.dropped.Load() } // Enqueue satisfies server.ForgeEventSink's contract (github_webhook.go:44-51): -// it MUST NOT block. It filters to board-relevant issue events (Change ∈ -// {OPENED, STATE, UPDATE} ∧ Kind == ISSUE — PR events and COMMENT-change events -// dropped, design.md:231-237) then channel try-sends; a full queue DROPS the -// event with the drop metric + a Warn (the T3 reconciler heals it). +// it MUST NOT block. It filters to board-relevant events, then channel +// try-sends; a full queue DROPS the event with the drop metric + a Warn (the +// reconciler heals it). func (a *BoardWebhookArm) Enqueue(_ context.Context, ev forge.ForgeEvent) { - if !boardRelevant(ev) { + if !a.boardRelevant(ev) { return } select { @@ -118,25 +128,6 @@ func (a *BoardWebhookArm) Enqueue(_ context.Context, ev forge.ForgeEvent) { } } -// boardRelevant reports whether an event is a board-relevant issue change: -// Kind == ISSUE and Change ∈ {OPENED, STATE, UPDATE}. COMMENT-change issue -// events (issue_comment also parses to Kind ISSUE, githubapp_webhook.go:196-210) -// and every PR-kind event are excluded — the board projects issues only, and -// admitting comments would burn one hydrate GET per comment (design.md:231-237). -func boardRelevant(ev forge.ForgeEvent) bool { - if ev.Kind != compassv1internal.ForgeArtifactKind_FORGE_ARTIFACT_KIND_ISSUE { - return false - } - switch ev.Change { - case compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_OPENED, - compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_STATE, - compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_UPDATE: - return true - default: - return false - } -} - // 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. @@ -151,18 +142,47 @@ func (a *BoardWebhookArm) Run(ctx context.Context) error { } } -// drainBatch coalesces first into a distinct-coordinate set (design.md:271-274): -// it seeds with the event that woke the drain, non-blockingly drains every other -// currently-queued event, and keys each on its NORMALIZED (repo, number) — so an -// edit storm on one issue, and a mixed-case duplicate of one repo, both collapse -// to a single coordinate. It then hydrates + sinks each distinct coordinate once, -// in arrival order. +// 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 +// and checks show on the PR, but only when the arm can hydrate PRs. +func (a *BoardWebhookArm) boardRelevant(ev forge.ForgeEvent) bool { + switch ev.Kind { + case compassv1internal.ForgeArtifactKind_FORGE_ARTIFACT_KIND_ISSUE: + switch ev.Change { + case compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_OPENED, + compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_STATE, + compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_UPDATE: + return true + default: + return false + } + case compassv1internal.ForgeArtifactKind_FORGE_ARTIFACT_KIND_PULL_REQUEST: + if a.pulls == nil { + return false + } + if ev.Number == 0 { + return ev.HeadSHA != "" && a.numbers != nil + } + return true + default: + return false + } +} + +// 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) { seen := map[boardCoord]struct{}{} var order []boardCoord add := func(ev forge.ForgeEvent) { - c := boardCoord{repo: normalizeBoardRepo(ev.Repo), number: ev.Number} + c := boardCoord{repo: normalizeBoardRepo(ev.Repo), kind: ev.Kind, number: ev.Number} + if ev.Number == 0 { + c.headSHA = ev.HeadSHA + } if _, ok := seen[c]; ok { return } @@ -182,18 +202,18 @@ func (a *BoardWebhookArm) drainBatch(ctx context.Context, first forge.ForgeEvent process: for _, c := range order { - if err := a.hydrateAndSink(ctx, c); err != nil { + 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) + "repo", c.repo, "number", c.number, "head_sha", c.headSHA) return } // 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, "err", err) + "repo", c.repo, "number", c.number, "head_sha", c.headSHA, "err", err) } if ctx.Err() != nil { return @@ -201,10 +221,34 @@ process: } } -// hydrateAndSink gates the coordinate's repo, hydrates the issue via an -// unconditional conditional GET (the event proves change; no stored per-issue -// ETag in v1, design.md:276-277), and sinks the fresh issue through the shared -// Ingester (ingest.go:82-99 — the one owner-strip/translate/stamp pipeline). A +// resolveAndSink maps a head-SHA CHECKS key to its PR number, then hydrates. +// A disabled repo spends no lookup. A commit with no PR is skipped, as is a PR +// already in this batch; seen gains the resolved key so a later duplicate is +// skipped too. +func (a *BoardWebhookArm) resolveAndSink(ctx context.Context, c boardCoord, seen map[boardCoord]struct{}) error { + if c.headSHA != "" { + enabled, err := a.targets.IsEnabledRepo(ctx, c.repo) + if err != nil || !enabled { + return err + } + n, err := a.numbers.PullNumberForSHA(ctx, c.repo, c.headSHA) + if errors.Is(err, forge.ErrNoPullRequestForSHA) { + return nil + } + if err != nil { + return err + } + c = boardCoord{repo: c.repo, kind: c.kind, number: n} + if _, dup := seen[c]; dup { + return nil + } + seen[c] = struct{}{} + } + return a.hydrateAndSink(ctx, c) +} + +// hydrateAndSink gates the coordinate's repo, then hydrates the artifact (the +// event proves change, so an issue read sends no stored ETag) and sinks it. A // non-enabled repo is dropped silently. func (a *BoardWebhookArm) hydrateAndSink(ctx context.Context, c boardCoord) error { enabled, err := a.targets.IsEnabledRepo(ctx, c.repo) @@ -214,6 +258,9 @@ func (a *BoardWebhookArm) hydrateAndSink(ctx context.Context, c boardCoord) erro if !enabled { return nil } + if c.kind == compassv1internal.ForgeArtifactKind_FORGE_ARTIFACT_KIND_PULL_REQUEST { + return a.pulls.Hydrate(ctx, c.repo, c.number) + } res, err := a.hydrator.GetIssueConditional(ctx, c.repo, c.number, "") if err != nil { return err diff --git a/go/internal/ingest/board_webhook_test.go b/go/internal/ingest/board_webhook_test.go index 5268a8c14..ef75005f7 100644 --- a/go/internal/ingest/board_webhook_test.go +++ b/go/internal/ingest/board_webhook_test.go @@ -160,9 +160,9 @@ func TestArmDropsNonEnabledRepo(t *testing.T) { } } -// TestArmDropsPRKind: a pull_request-kind event is filtered at Enqueue — never -// queued, never hydrated. -func TestArmDropsPRKind(t *testing.T) { +// TestArmDropsPRKindWithoutPullHydrator: with no PR hydrator configured, a +// pull_request-kind event is filtered at Enqueue — never queued, never hydrated. +func TestArmDropsPRKindWithoutPullHydrator(t *testing.T) { h := &fakeHydrator{result: forge.Issue{State: "open"}} arm, sink := newArmHarness(t, h, newFakeTargets("owner/repo"), 16) diff --git a/go/internal/ingest/notify_pull_cache.go b/go/internal/ingest/notify_pull_cache.go index ce38b81bf..0f95336de 100644 --- a/go/internal/ingest/notify_pull_cache.go +++ b/go/internal/ingest/notify_pull_cache.go @@ -8,6 +8,7 @@ package ingest import ( "context" "errors" + "strings" "sync" "time" @@ -114,7 +115,8 @@ func newCachedPullNumberResolver(base PullNumberResolver, ttl time.Duration, now // forge.ErrNoPullRequestForSHA so the router's step 0 cannot tell a cache hit // from a fresh resolve. func (c *CachedPullNumberResolver) PullNumberForSHA(ctx context.Context, repo, headSHA string) (uint64, error) { - key := pullNumberKey{repo: repo, sha: headSHA} + // The board lane passes a lowercased repo and the notify lane does not. + key := pullNumberKey{repo: strings.ToLower(repo), sha: headSHA} now := c.now() c.mu.Lock() diff --git a/go/internal/ingest/notify_pull_number_test.go b/go/internal/ingest/notify_pull_number_test.go index f49a9de1c..d2af8e776 100644 --- a/go/internal/ingest/notify_pull_number_test.go +++ b/go/internal/ingest/notify_pull_number_test.go @@ -472,3 +472,18 @@ func TestCachedPullNumberResolverLastObservedWins(t *testing.T) { t.Errorf("after a later observation = (%d, %v), want (88, nil)", num, err) } } + +// TestCachedPullNumberResolverIgnoresRepoCase: the board and notify lanes pass +// the same repo in different casing and share one lookup. +func TestCachedPullNumberResolverIgnoresRepoCase(t *testing.T) { + base := &countingResolver{bySHA: map[string]uint64{"abc": 4}} + c := newCachedPullNumberResolver(base, time.Minute, time.Now) + for _, repo := range []string{"Owner/Repo", "owner/repo"} { + if n, err := c.PullNumberForSHA(context.Background(), repo, "abc"); err != nil || n != 4 { + t.Fatalf("PullNumberForSHA(%q) = %d, %v", repo, n, err) + } + } + if base.calls != 1 { + t.Fatalf("base calls = %d, want 1", base.calls) + } +} diff --git a/go/internal/ingest/pull_request_sink.go b/go/internal/ingest/pull_request_sink.go index 396f8f0d8..a21d0094c 100644 --- a/go/internal/ingest/pull_request_sink.go +++ b/go/internal/ingest/pull_request_sink.go @@ -1,6 +1,8 @@ package ingest import ( + "context" + "fmt" "time" compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" @@ -14,3 +16,55 @@ type IngestedPullRequest struct { UpdatedAt time.Time ClosingRefs []forge.IssueRef } + +// pullRequestReader is the PR read surface, satisfied by *forge.GitHub. +type pullRequestReader interface { + GetPullRequest(ctx context.Context, repo string, number uint64) (forge.PullRequest, error) + ListOpenPullRequests(ctx context.Context, repo string) ([]forge.UpdatedPull, error) +} + +// prSink receives each hydrated PR; the board's IssueProjection satisfies it. +type prSink interface { + PublishPullRequestUpdate(ctx context.Context, in IngestedPullRequest) error +} + +// PullRequestHydrator reads one PR, translates it the way the create path does, +// and sinks it to the board. +type PullRequestHydrator struct { + reader pullRequestReader + sink prSink + forgeRef *compassv1.ForgeRef +} + +// NewPullRequestHydrator returns a hydrator reading from r, stamping ref and +// sinking to s. +func NewPullRequestHydrator(r pullRequestReader, s prSink, ref *compassv1.ForgeRef) *PullRequestHydrator { + return &PullRequestHydrator{reader: r, sink: s, forgeRef: ref} +} + +// Hydrate reads PR number in repo and sinks it. Errors wrap the cause, so a +// caller can errors.Is forge.ErrBudgetExhausted. +func (h *PullRequestHydrator) Hydrate(ctx context.Context, repo string, number uint64) error { + raw, err := h.reader.GetPullRequest(ctx, repo, number) + if err != nil { + return fmt.Errorf("ingest: read pull request #%d for %q: %w", number, repo, err) + } + _, author, ok := forge.StripOwner(raw.Body) + var attr *compassv1.AgentAttribution + if ok && author.AgentHandle != "" { + attr = &compassv1.AgentAttribution{AgentHandle: author.AgentHandle} + } + pr := forge.TranslatePullRequest(raw, attr) + pr.Forge = &compassv1.ForgeRef{Provider: h.forgeRef.GetProvider(), Host: h.forgeRef.GetHost()} + pr.Repo = repo + in := IngestedPullRequest{PR: pr, CreatedAt: raw.CreatedAt, UpdatedAt: raw.UpdatedAt, ClosingRefs: raw.ClosingRefs} + if err := h.sink.PublishPullRequestUpdate(ctx, in); err != nil { + return fmt.Errorf("ingest: publish pull request #%d for %q: %w", number, repo, err) + } + return nil +} + +// listOpen lists the repo's open PRs for the backfill pass. +func (h *PullRequestHydrator) listOpen(ctx context.Context, repo string) ([]forge.UpdatedPull, error) { + return h.reader.ListOpenPullRequests(ctx, repo) +} diff --git a/go/internal/ingest/pull_request_sink_test.go b/go/internal/ingest/pull_request_sink_test.go new file mode 100644 index 000000000..a71bf6278 --- /dev/null +++ b/go/internal/ingest/pull_request_sink_test.go @@ -0,0 +1,448 @@ +package ingest + +import ( + "context" + "errors" + "slices" + "sync" + "testing" + "time" + + "github.com/RigelBuild/compass/go/internal/forge" + compassv1internal "github.com/RigelBuild/compass/go/internal/gen/compass/v1" +) + +// fakePulls is the PR read surface: GetPullRequest records each number read +// and fails with errFor[number] when set. +type fakePulls struct { + mu sync.Mutex + reads []uint64 + errFor map[uint64]error + open []forge.UpdatedPull +} + +func (f *fakePulls) GetPullRequest(_ context.Context, _ string, number uint64) (forge.PullRequest, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.reads = append(f.reads, number) + if err := f.errFor[number]; err != nil { + return forge.PullRequest{}, err + } + return forge.PullRequest{ + Number: number, + State: "open", + Body: "\nbody", + CreatedAt: ts(1), + UpdatedAt: ts(2), + ClosingRefs: []forge.IssueRef{{Repo: "o/r", Number: 3}}, + }, nil +} + +func (f *fakePulls) ListOpenPullRequests(_ context.Context, _ string) ([]forge.UpdatedPull, error) { + return f.open, nil +} + +func (f *fakePulls) readNumbers() []uint64 { + f.mu.Lock() + defer f.mu.Unlock() + return slices.Clone(f.reads) +} + +type recordingPRSink struct{ got []IngestedPullRequest } + +func (s *recordingPRSink) PublishPullRequestUpdate(_ context.Context, in IngestedPullRequest) error { + s.got = append(s.got, in) + return nil +} + +// fakeNumbers resolves head SHAs; an unknown SHA has no PR. +type fakeNumbers map[string]uint64 + +func (f fakeNumbers) PullNumberForSHA(_ context.Context, _, sha string) (uint64, error) { + if n, ok := f[sha]; ok { + return n, nil + } + return 0, forge.ErrNoPullRequestForSHA +} + +func newPRArm(t *testing.T, pulls *fakePulls, numbers PullNumberResolver) (*BoardWebhookArm, *recordingPRSink) { + t.Helper() + return newPRArmFor(t, pulls, numbers, newFakeTargets("owner/repo")) +} + +func newPRArmFor(t *testing.T, pulls *fakePulls, numbers PullNumberResolver, targets TargetChecker) (*BoardWebhookArm, *recordingPRSink) { + t.Helper() + sink := &recordingPRSink{} + h := NewPullRequestHydrator(pulls, sink, testForgeRef()) + ing := NewIngester(forge.NewFakeProvider("gh"), &recordingSink{}, testForgeRef()) + arm := NewBoardWebhookArm(&fakeHydrator{}, ing, targets, BoardArmConfig{QueueSize: 64, Pulls: h, PullNumbers: numbers}) + return arm, sink +} + +// checksEvent is a SHA-only CHECKS event for head SHA "abc". +func checksEvent() forge.ForgeEvent { + ev := prEvent(0, compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_CHECKS) + ev.HeadSHA = "abc" + return ev +} + +// TestArmChecksOnDisabledRepoSpendsNoLookup: a CHECKS SHA event on a repo not +// enabled on the board makes no resolver call and no read. +func TestArmChecksOnDisabledRepoSpendsNoLookup(t *testing.T) { + pulls := &fakePulls{} + res := &countingResolver{bySHA: map[string]uint64{"abc": 9}} + arm, _ := newPRArmFor(t, pulls, res, newFakeTargets()) + arm.Enqueue(context.Background(), checksEvent()) + drainAll(context.Background(), arm) + if res.calls != 0 || len(pulls.readNumbers()) != 0 { + t.Fatalf("resolver calls=%d reads=%v, want none", res.calls, pulls.readNumbers()) + } +} + +// TestArmResolverBudgetPausesDrain: a budget error resolving a SHA abandons the +// rest of the batch. +func TestArmResolverBudgetPausesDrain(t *testing.T) { + pulls := &fakePulls{} + arm, _ := newPRArm(t, pulls, &countingResolver{err: forge.ErrBudgetExhausted}) + arm.Enqueue(context.Background(), checksEvent()) + arm.Enqueue(context.Background(), prEvent(11, changeUpdate)) + drainAll(context.Background(), arm) + if got := pulls.readNumbers(); len(got) != 0 { + t.Fatalf("reads = %v, want none after the budget pause", got) + } +} + +// TestArmNumberedBeforeChecksCoalesces: a numbered event ahead of a CHECKS +// event for the same PR still costs one hydrate. +func TestArmNumberedBeforeChecksCoalesces(t *testing.T) { + pulls := &fakePulls{} + arm, _ := newPRArm(t, pulls, fakeNumbers{"abc": 9}) + arm.Enqueue(context.Background(), prEvent(9, changeUpdate)) + arm.Enqueue(context.Background(), checksEvent()) + drainAll(context.Background(), arm) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{9}) { + t.Fatalf("reads = %v, want [9]", got) + } +} + +// TestArmDropsChecksWithoutResolver: with no resolver a SHA-only CHECKS event +// is filtered at Enqueue. +func TestArmDropsChecksWithoutResolver(t *testing.T) { + arm, _ := newPRArm(t, &fakePulls{}, nil) + arm.Enqueue(context.Background(), checksEvent()) + if l := len(arm.queue); l != 0 { + t.Fatalf("queue len = %d, want 0", l) + } +} + +func prEvent(number uint64, change compassv1internal.ForgeNotificationKind) forge.ForgeEvent { + ev := issueEvent("Owner/Repo", number, change) + ev.Kind = boardKindPR + return ev +} + +// TestArmHydratesEveryPRChangeKind: each PR change kind reaches the sink with +// the stamped coordinate, stripped attribution, forge times and closing refs. +func TestArmHydratesEveryPRChangeKind(t *testing.T) { + kinds := []compassv1internal.ForgeNotificationKind{ + compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_OPENED, + changeState, changeUpdate, changeComment, + compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_REVIEW, + compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_CHECKS, + } + for _, k := range kinds { + t.Run(k.String(), func(t *testing.T) { + pulls := &fakePulls{} + arm, sink := newPRArm(t, pulls, fakeNumbers{}) + arm.Enqueue(context.Background(), prEvent(9, k)) + drainAll(context.Background(), arm) + if len(sink.got) != 1 { + t.Fatalf("sank %d PRs, want 1", len(sink.got)) + } + got := sink.got[0] + if got.PR.GetNumber() != 9 || got.PR.GetRepo() != "owner/repo" || got.PR.GetForge().GetHost() != "github.com" { + t.Fatalf("coordinate = %d %q %v", got.PR.GetNumber(), got.PR.GetRepo(), got.PR.GetForge()) + } + if got.PR.GetAgent().GetAgentHandle() != "a" { + t.Fatalf("agent = %v, want a", got.PR.GetAgent()) + } + if !got.CreatedAt.Equal(ts(1)) || !got.UpdatedAt.Equal(ts(2)) || len(got.ClosingRefs) != 1 { + t.Fatalf("ingested = %+v", got) + } + }) + } +} + +// TestArmResolvesChecksByHeadSHA: a CHECKS event with only a head SHA hydrates +// its PR; a SHA with no PR is skipped without a read. +func TestArmResolvesChecksByHeadSHA(t *testing.T) { + pulls := &fakePulls{} + arm, sink := newPRArm(t, pulls, fakeNumbers{"abc": 12}) + for _, sha := range []string{"abc", "nopr"} { + ev := prEvent(0, compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_CHECKS) + ev.HeadSHA = sha + arm.Enqueue(context.Background(), ev) + } + drainAll(context.Background(), arm) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{12}) { + t.Fatalf("reads = %v, want [12]", got) + } + if len(sink.got) != 1 { + t.Fatalf("sank %d, want 1", len(sink.got)) + } +} + +// TestArmCoalescesPRBurst: review, comment and check events on one PR in one +// batch, including a CHECKS event resolved by SHA, cost one hydrate. +func TestArmCoalescesPRBurst(t *testing.T) { + pulls := &fakePulls{} + arm, _ := newPRArm(t, pulls, fakeNumbers{"abc": 9}) + arm.Enqueue(context.Background(), prEvent(9, compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_REVIEW)) + arm.Enqueue(context.Background(), prEvent(9, changeComment)) + checks := prEvent(0, compassv1internal.ForgeNotificationKind_FORGE_NOTIFICATION_KIND_CHECKS) + checks.HeadSHA = "abc" + arm.Enqueue(context.Background(), checks) + arm.Enqueue(context.Background(), prEvent(9, changeUpdate)) + drainAll(context.Background(), arm) + if got := pulls.readNumbers(); len(got) != 1 { + t.Fatalf("reads = %v, want one", got) + } +} + +// TestArmIssueAndPRSameNumberStayDistinct: an issue and a PR never coalesce. +func TestArmIssueAndPRSameNumberStayDistinct(t *testing.T) { + pulls := &fakePulls{} + arm, _ := newPRArm(t, pulls, nil) + arm.Enqueue(context.Background(), issueEvent("owner/repo", 9, changeUpdate)) + arm.Enqueue(context.Background(), prEvent(9, changeUpdate)) + drainAll(context.Background(), arm) + if got := pulls.readNumbers(); len(got) != 1 { + t.Fatalf("PR reads = %v, want one", got) + } +} + +func newPRReconciler(l updatedLister, st *fakeBoardStore, pulls *fakePulls) (*BoardReconciler, *recordingSink) { + sink := &recordingSink{} + h := NewPullRequestHydrator(pulls, &recordingPRSink{}, testForgeRef()) + return NewBoardReconciler(l, NewIngester(nil, sink, testForgeRef()), st, BoardReconcileConfig{Pace: -1, Pulls: h}), sink +} + +// rowsLister scripts one walk result per call and records each since. +type rowsLister struct { + results []forge.ConditionalResult[forge.UpdatedRows] + since []time.Time + etags []string +} + +func (l *rowsLister) ListUpdatedIssues(_ context.Context, _ string, since time.Time, etag string) (forge.ConditionalResult[forge.UpdatedRows], error) { + l.since = append(l.since, since) + l.etags = append(l.etags, etag) + i := min(len(l.since)-1, len(l.results)-1) + return l.results[i], nil +} + +// TestReconcileSkipsUnchangedPR: a PR row no newer than the stored row is not +// re-hydrated; a newer one is, and a never-stored one is. +func TestReconcileSkipsUnchangedPR(t *testing.T) { + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{{V: forge.UpdatedRows{Pulls: []forge.UpdatedPull{ + {Number: 1, State: "open", UpdatedAt: ts(10)}, + {Number: 2, State: "open", UpdatedAt: ts(20)}, + {Number: 3, State: "open", UpdatedAt: ts(30)}, + }}, ETag: `"e"`}}} + st := newBoardStore("o/r") + st.marks["o/r"] = storedMark{mark: ts(5)} + st.backfilled["o/r"] = ts(0) + st.prUpdated[1] = ts(10) + st.prUpdated[2] = ts(15) + pulls := &fakePulls{} + rc, _ := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{2, 3}) { + t.Fatalf("reads = %v, want [2 3]", got) + } + if m := st.marks["o/r"]; !m.mark.Equal(ts(30)) || m.etag != `"e"` { + t.Fatalf("watermark = %+v, want ts(30) with the list ETag", m) + } +} + +// TestReconcileBudgetOnPRAbortsWindow: a budget error on a PR row aborts the +// sweep, keeps the watermark at since with no ETag, and the next sweep re-lists +// the window unconditionally so an older issue row is seen again. +func TestReconcileBudgetOnPRAbortsWindow(t *testing.T) { + window := forge.ConditionalResult[forge.UpdatedRows]{V: forge.UpdatedRows{ + Issues: []forge.Issue{{Number: 4, UpdatedAt: ts(10)}}, + Pulls: []forge.UpdatedPull{{Number: 5, State: "open", UpdatedAt: ts(20)}}, + }, ETag: `"e"`} + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{window}} + st := newBoardStore("o/r", "o/next") + st.marks["o/r"] = storedMark{mark: ts(5), etag: `"old"`} + st.backfilled["o/r"] = ts(0) + pulls := &fakePulls{errFor: map[uint64]error{5: forge.ErrBudgetExhausted}} + rc, sink := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + + if m := st.marks["o/r"]; !m.mark.Equal(ts(5)) || m.etag != "" { + t.Fatalf("watermark = %+v, want since ts(5) with no ETag", m) + } + if len(l.since) != 1 { + t.Fatalf("lists = %d, want 1 (sweep aborted before o/next)", len(l.since)) + } + pulls.errFor = nil + st.repos = []string{"o/r"} + rc.sweep(context.Background()) + if !l.since[1].Equal(ts(5)) || l.etags[1] != "" { + t.Fatalf("second list since=%v etag=%q, want ts(5) and no ETag", l.since[1], l.etags[1]) + } + if n := countNumber(sink, 4); n != 2 { + t.Fatalf("issue #4 sank %d times, want 2 (re-listed)", n) + } +} + +func countNumber(s *recordingSink, n uint32) int { + c := 0 + for _, iss := range s.got { + if iss.GetNumber() == n { + c++ + } + } + return c +} + +// TestReconcileBackfillsOpenPRsOnce: a repo with a watermark but no backfill +// mark hydrates its open PRs and recent rows once, then is marked. +func TestReconcileBackfillsOpenPRsOnce(t *testing.T) { + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{ + {NotModified: true}, + {V: forge.UpdatedRows{Pulls: []forge.UpdatedPull{{Number: 8, State: "closed", UpdatedAt: ts(3)}}}}, + {NotModified: true}, + }} + st := newBoardStore("o/r") + st.marks["o/r"] = storedMark{mark: ts(50), etag: `"e"`} + pulls := &fakePulls{open: []forge.UpdatedPull{{Number: 7, State: "open", UpdatedAt: ts(40)}}} + rc, _ := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + rc.sweep(context.Background()) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{7, 8}) { + t.Fatalf("reads = %v, want [7 8] once", got) + } + if _, ok := st.backfilled["o/r"]; !ok { + t.Fatal("repo not marked backfilled") + } + if !l.since[1].Equal(ts(50).Add(-prBackfillWindow)) { + t.Fatalf("backfill window since = %v", l.since[1]) + } +} + +// TestReconcileHydratesCreatedPRWhoseWebhookDropped: the create path stores the +// epoch as forge_updated_at, so the next sweep hydrates the PR. +func TestReconcileHydratesCreatedPRWhoseWebhookDropped(t *testing.T) { + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{{V: forge.UpdatedRows{Pulls: []forge.UpdatedPull{ + {Number: 6, State: "open", UpdatedAt: ts(10)}, + }}}}} + st := newBoardStore("o/r") + st.marks["o/r"] = storedMark{mark: ts(5)} + st.backfilled["o/r"] = ts(0) + st.prUpdated[6] = time.Unix(0, 0) + pulls := &fakePulls{} + rc, _ := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{6}) { + t.Fatalf("reads = %v, want [6]", got) + } +} + +// TestReconcileColdStartSkipsOldClosedPR: with no watermark, a closed PR older +// than the backfill window is not hydrated; an open one and a recent one are. +func TestReconcileColdStartSkipsOldClosedPR(t *testing.T) { + now := time.Now() + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{{V: forge.UpdatedRows{Pulls: []forge.UpdatedPull{ + {Number: 1, State: "closed", UpdatedAt: now.Add(-40 * 24 * time.Hour)}, + {Number: 2, State: "open", UpdatedAt: now.Add(-90 * 24 * time.Hour)}, + {Number: 3, State: "closed", UpdatedAt: now.Add(-time.Hour)}, + }}}}} + st := newBoardStore("o/r") + pulls := &fakePulls{} + rc, _ := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{2, 3}) { + t.Fatalf("reads = %v, want [2 3]", got) + } + if len(l.since) != 1 { + t.Fatalf("lists = %d, want 1 (cold start needs no backfill window walk)", len(l.since)) + } +} + +// TestReconcileColdStartRowFailureLeavesRepoUnmarked: a recent closed PR that +// fails in the cold-start walk keeps the repo unmarked, so the next sweep's +// backfill window re-lists and hydrates it. +func TestReconcileColdStartRowFailureLeavesRepoUnmarked(t *testing.T) { + now := time.Now() + closed := forge.UpdatedPull{Number: 8, State: "closed", UpdatedAt: now.Add(-time.Hour)} + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{ + {V: forge.UpdatedRows{ + Issues: []forge.Issue{{Number: 4, UpdatedAt: now.Add(-time.Minute)}}, + Pulls: []forge.UpdatedPull{closed}, + }}, + {NotModified: true}, + {V: forge.UpdatedRows{Pulls: []forge.UpdatedPull{closed}}}, + }} + st := newBoardStore("o/r") + pulls := &fakePulls{errFor: map[uint64]error{8: errors.New("boom")}} + rc, _ := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + if _, ok := st.backfilled["o/r"]; ok { + t.Fatal("repo marked backfilled after a cold-start row failure") + } + pulls.errFor = nil + rc.sweep(context.Background()) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{8, 8}) { + t.Fatalf("reads = %v, want [8 8]", got) + } + if _, ok := st.backfilled["o/r"]; !ok { + t.Fatal("repo not marked after the clean retry") + } +} + +// TestReconcileBackfillStopsOnBudget: a budget error during backfill leaves the +// repo unmarked so the next sweep retries it. +func TestReconcileBackfillStopsOnBudget(t *testing.T) { + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{{NotModified: true}}} + st := newBoardStore("o/r") + st.marks["o/r"] = storedMark{mark: ts(50)} + pulls := &fakePulls{ + open: []forge.UpdatedPull{{Number: 7, State: "open", UpdatedAt: ts(40)}}, + errFor: map[uint64]error{7: forge.ErrBudgetExhausted}, + } + rc, _ := newPRReconciler(l, st, pulls) + if err := rc.reconcileRepo(context.Background(), "o/r"); !errors.Is(err, forge.ErrBudgetExhausted) { + t.Fatalf("err = %v, want budget exhausted", err) + } + if _, ok := st.backfilled["o/r"]; ok { + t.Fatal("repo marked backfilled after a budget abort") + } +} + +// TestReconcileBackfillRetriesAfterRowFailure: a non-budget hydrate failure +// leaves the repo unmarked, and the next sweep retries only the failed row. +func TestReconcileBackfillRetriesAfterRowFailure(t *testing.T) { + l := &rowsLister{results: []forge.ConditionalResult[forge.UpdatedRows]{{NotModified: true}}} + st := newBoardStore("o/r") + pulls := &fakePulls{ + open: []forge.UpdatedPull{{Number: 7, State: "open", UpdatedAt: ts(40)}, {Number: 8, State: "open", UpdatedAt: ts(41)}}, + errFor: map[uint64]error{8: errors.New("boom")}, + } + rc, _ := newPRReconciler(l, st, pulls) + rc.sweep(context.Background()) + if _, ok := st.backfilled["o/r"]; ok { + t.Fatal("repo marked backfilled after a failed row") + } + st.prUpdated[7] = ts(40) + pulls.errFor = nil + rc.sweep(context.Background()) + if got := pulls.readNumbers(); !slices.Equal(got, []uint64{7, 8, 8}) { + t.Fatalf("reads = %v, want [7 8 8]", got) + } + if _, ok := st.backfilled["o/r"]; !ok { + t.Fatal("repo not marked after the clean retry") + } +} diff --git a/go/server/board_pr_ingest_pgtest_test.go b/go/server/board_pr_ingest_pgtest_test.go new file mode 100644 index 000000000..38f70923d --- /dev/null +++ b/go/server/board_pr_ingest_pgtest_test.go @@ -0,0 +1,98 @@ +//go:build pgtest && unix + +package server + +import ( + "context" + "testing" + "time" + + "github.com/RigelBuild/compass/go/events" + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/board" + "github.com/RigelBuild/compass/go/internal/forge" + "github.com/RigelBuild/compass/go/internal/ingest" + "github.com/RigelBuild/compass/go/internal/store" +) + +// fakePRReads serves one PR closing an issue and lists it as open. +type fakePRReads struct{ pr forge.PullRequest } + +func (f *fakePRReads) GetPullRequest(context.Context, string, uint64) (forge.PullRequest, error) { + return f.pr, nil +} + +func (f *fakePRReads) ListOpenPullRequests(context.Context, string) ([]forge.UpdatedPull, error) { + return []forge.UpdatedPull{{Number: f.pr.Number, State: "open", UpdatedAt: f.pr.UpdatedAt}}, nil +} + +type notModifiedLister struct{} + +func (notModifiedLister) ListUpdatedIssues(context.Context, string, time.Time, string) (forge.ConditionalResult[forge.UpdatedRows], error) { + return forge.ConditionalResult[forge.UpdatedRows]{NotModified: true}, nil +} + +// TestBoardReconcileBackfillAttachesOpenPR drives the backfill through the real +// store adapter: an enabled repo with a watermark but no backfill mark hydrates +// its open PR onto the closing-ref issue, then is marked once. +func TestBoardReconcileBackfillAttachesOpenPR(t *testing.T) { + ctx := context.Background() + st := forgeTestStore(t) + const repo = "owner/backfill" + if err := st.EnsureForgeRepoSubscription(ctx, store.ForgeRepoSubscription{ + Provider: store.ForgeProviderGitHub, Host: forgeTestHost, Repo: repo, Enabled: true, + }); err != nil { + t.Fatalf("seed subscription: %v", err) + } + adapter := &boardReconcileStore{st: st, provider: store.ForgeProviderGitHub, host: forgeTestHost} + if err := adapter.StoreRepoWatermark(ctx, repo, time.Date(2026, 5, 1, 0, 0, 0, 0, time.UTC), `"e"`); err != nil { + t.Fatalf("seed watermark: %v", err) + } + bus := events.NewBus[busPayload]() + t.Cleanup(bus.Close) + brd := board.NewIssueProjection(bus, st) + ref := &compassv1.ForgeRef{Provider: compassv1.ForgeProvider_FORGE_PROVIDER_GITHUB, Host: forgeTestHost} + if err := brd.PublishIssueUpdate(ctx, &compassv1.Issue{Forge: ref, Repo: repo, Number: 2, Title: "t", ForgeState: "open"}); err != nil { + t.Fatalf("PublishIssueUpdate: %v", err) + } + at := time.Date(2026, 4, 1, 0, 0, 0, 0, time.UTC) + reads := &fakePRReads{pr: forge.PullRequest{Number: 40, State: "open", CreatedAt: at, UpdatedAt: at, + ClosingRefs: []forge.IssueRef{{Repo: repo, Number: 2}}}} + marked := &markSignal{boardReconcileStore: adapter, done: make(chan struct{})} + rc := ingest.NewBoardReconciler(notModifiedLister{}, ingest.NewIngester(nil, brd, ref), marked, ingest.BoardReconcileConfig{ + Pace: -1, Pulls: ingest.NewPullRequestHydrator(reads, brd, ref), + }) + runCtx, cancel := context.WithCancel(ctx) + defer cancel() + go func() { _ = rc.Run(runCtx) }() + select { + case <-marked.done: + case <-time.After(10 * time.Second): + t.Fatal("backfill never marked the repo") + } + + var prs []*compassv1.PullRequest + for _, iss := range brd.Snapshot() { + if iss.GetNumber() == 2 { + prs = iss.GetPrs() + } + } + if len(prs) != 1 || prs[0].GetNumber() != 40 { + t.Fatalf("issue #2 prs = %v, want PR #40", prs) + } + if _, ok, err := adapter.PRsBackfilledAt(ctx, repo); err != nil || !ok { + t.Fatalf("PRsBackfilledAt ok=%v err=%v, want marked", ok, err) + } +} + +// markSignal closes done once the backfill mark is written. +type markSignal struct { + *boardReconcileStore + done chan struct{} +} + +func (m *markSignal) MarkPRsBackfilled(ctx context.Context, repo string, at time.Time) error { + err := m.boardReconcileStore.MarkPRsBackfilled(ctx, repo, at) + close(m.done) + return err +} diff --git a/go/server/serve.go b/go/server/serve.go index b183e011a..b3123d92d 100644 --- a/go/server/serve.go +++ b/go/server/serve.go @@ -1334,11 +1334,13 @@ func buildBoardWebhookWiring( } client := forge.NewGitHub(forge.GitHubConfig{Host: rc.Host, Token: tok, Client: httpClient}) - lane, err := buildBoardIngestLane(ctx, cfg, st, issueBrd, client, log) + // The board lane reuses the notify lane's head-SHA cache: every check_suite + // reaches both lanes, so one cache keeps it to one lookup. + notifyLane := buildForgeNotifyLane(cfg, st, hub, client, log) + lane, err := buildBoardIngestLane(ctx, cfg, st, issueBrd, client, notifyLane.pulls, log) if err != nil { return nil, nil, nil, nil, nil, err } - notifyLane := buildForgeNotifyLane(cfg, st, hub, client, log) sink := &fanoutSink{sinks: []ForgeEventSink{lane.sink, notifyLane.sink}} secret := newCachedWebhookSecret(resolver, rc.App.AppWebhookSecretName) // client is returned as the shared primary App client the author write leg @@ -1406,6 +1408,7 @@ func buildBoardIngestLane( st *store.Store, issueBrd *board.IssueProjection, client *forge.GitHub, + pullNumbers ingest.PullNumberResolver, log *slog.Logger, ) (*boardIngestLane, error) { fc := cfg.Forge.resolved() @@ -1423,14 +1426,21 @@ func buildBoardIngestLane( // (2) Assemble the pipeline over the shared client: the Ingester (sharing the // existing issue projection), and the two arms over store adapters binding // (provider, host). - ing := ingest.NewIngester(client, issueBrd, &compassv1.ForgeRef{ + ref := &compassv1.ForgeRef{ Provider: compassv1.ForgeProvider_FORGE_PROVIDER_GITHUB, Host: fc.Host, + } + ing := ingest.NewIngester(client, issueBrd, ref) + pulls := ingest.NewPullRequestHydrator(client, issueBrd, ref) + arm := ingest.NewBoardWebhookArm(client, ing, &boardTargetStore{st: st}, ingest.BoardArmConfig{ + Log: log, + Pulls: pulls, + PullNumbers: pullNumbers, }) - arm := ingest.NewBoardWebhookArm(client, ing, &boardTargetStore{st: st}, ingest.BoardArmConfig{Log: log}) reconciler := ingest.NewBoardReconciler(client, ing, &boardReconcileStore{st: st, provider: provider, host: fc.Host}, ingest.BoardReconcileConfig{ Backstop: fc.App.ReconcileBackstop, Log: log, + Pulls: pulls, }) return &boardIngestLane{arm: arm, reconciler: reconciler, sink: arm, client: client}, nil } @@ -1478,6 +1488,21 @@ func (a *boardReconcileStore) StoreRepoWatermark(ctx context.Context, repo strin return a.st.StoreForgeRepoWatermark(ctx, a.provider, a.host, repo, mark, etag) } +// PullRequestUpdatedAt returns the stored PR row's forge updated_at. +func (a *boardReconcileStore) PullRequestUpdatedAt(ctx context.Context, repo string, number uint64) (time.Time, bool, error) { + return a.st.PullRequestUpdatedAt(ctx, store.ForgeCoord{Provider: a.provider, Host: a.host, Repo: repo, Number: number}) +} + +// PRsBackfilledAt returns when the repo's PR backfill ran. +func (a *boardReconcileStore) PRsBackfilledAt(ctx context.Context, repo string) (time.Time, bool, error) { + return a.st.PRsBackfilledAt(ctx, a.provider, a.host, repo) +} + +// MarkPRsBackfilled records that the repo's PR backfill ran. +func (a *boardReconcileStore) MarkPRsBackfilled(ctx context.Context, repo string, at time.Time) error { + return a.st.MarkPRsBackfilled(ctx, a.provider, a.host, repo, at) +} + // forgeNotifyLane is the assembled GitHub agent-notification lane (RIG-2732 T7): // the notify webhook arm (whose Run drains the ingress queue and whose Enqueue is // the accepted-event sink) and the notify reconciler (whose Run sweeps the diff --git a/go/server/serve_forge_budget_test.go b/go/server/serve_forge_budget_test.go index 5ed24a022..5ad1f4396 100644 --- a/go/server/serve_forge_budget_test.go +++ b/go/server/serve_forge_budget_test.go @@ -95,7 +95,7 @@ func TestForgeLanesShareOneBudgetGate(t *testing.T) { // Both lanes are built over the SAME client (nil store is safe: an empty seed // means reconcileForgeSeed never touches it, and this test never sweeps). - boardLane, err := buildBoardIngestLane(ctx, cfg, nil, nil, client, slog.Default()) + boardLane, err := buildBoardIngestLane(ctx, cfg, nil, nil, client, nil, slog.Default()) if err != nil { t.Fatalf("buildBoardIngestLane: %v", err) }