Skip to content
Open
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
176 changes: 142 additions & 34 deletions go/internal/ingest/board_reconcile.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand All @@ -65,6 +77,7 @@ type BoardReconciler struct {
lister updatedLister
ingester *Ingester
store BoardStore
pulls *PullRequestHydrator
backstop time.Duration
pace time.Duration
log *slog.Logger
Expand All @@ -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,
Expand Down Expand Up @@ -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 {
Expand All @@ -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
}
19 changes: 18 additions & 1 deletion go/internal/ingest/board_reconcile_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading
Loading