From b0f2660ca527e75a992b6d6cdb4fb5c28907b2df Mon Sep 17 00:00:00 2001 From: Landyn Date: Sat, 13 Jun 2026 19:05:24 -0500 Subject: [PATCH] fix(webhook): self-heal PR state drift from missed merge webhooks PR state (OPEN -> MERGED/CLOSED) was written only by the pull_request webhook handler. A single missed/dropped `pull_request.closed` delivery left a merged PR stuck OPEN forever: the frequent metadata-fetch path deliberately never touched state, there was no reconciliation, and backfill only ran on manual trigger. Observed on phase-rs/phase (#2745/2751/2752/2753/2756/3095 merged on GitHub, OPEN in the mirror), which makes the validator under-credit the miner since it reads PR state verbatim from the mirror. - Fetcher: fetchPrMetadata now also returns authoritative GraphQL state/mergedAt/closedAt/mergedBy. - Metadata handler: re-asserts that state (MERGED is terminal, so an in-flight stale fetch can't revert it) and logs corrected drift. - PrReconcileService: hourly sweep enqueues a metadata fetch for every still-open PR in registered repos within the scoring window, so missed merge events self-heal. Window/interval env-tunable. - RepoBackfillScheduleService: daily full backfill per registered repo as a coarse safety net; env kill-switch (NIGHTLY_BACKFILL_ENABLED=false). - Webhook handler: derive merged state from `merged` alone (synthesize merged_at from closed_at) so a closed event with a not-yet-populated merged_at isn't pinned to OPEN. --- .env.example | 10 ++ packages/das/src/queue/fetch.processor.ts | 49 ++++++++-- .../das/src/webhook/github-fetcher.service.ts | 18 ++++ .../webhook/handlers/pull-request.handler.ts | 11 ++- .../das/src/webhook/pr-reconcile.service.ts | 90 +++++++++++++++++ .../webhook/repo-backfill-schedule.service.ts | 97 +++++++++++++++++++ packages/das/src/webhook/webhook.module.ts | 4 + 7 files changed, 267 insertions(+), 12 deletions(-) create mode 100644 packages/das/src/webhook/pr-reconcile.service.ts create mode 100644 packages/das/src/webhook/repo-backfill-schedule.service.ts diff --git a/.env.example b/.env.example index 40600a6..0de88aa 100644 --- a/.env.example +++ b/.env.example @@ -23,3 +23,13 @@ REDIS_PORT=6379 # Validator API Keys (comma-separated) API_KEYS= + +# PR state reconciliation (self-heals missed pull_request.closed webhooks). +# Hourly sweep re-checks every still-open PR within the window against GitHub. +PR_RECONCILE_INTERVAL_MS=3600000 +PR_RECONCILE_WINDOW_DAYS=45 + +# Nightly full repo backfill (coarse safety net; heavier — set false to disable). +NIGHTLY_BACKFILL_ENABLED=true +NIGHTLY_BACKFILL_INTERVAL_MS=86400000 +NIGHTLY_BACKFILL_DAYS=40 diff --git a/packages/das/src/queue/fetch.processor.ts b/packages/das/src/queue/fetch.processor.ts index 9600559..2d0f398 100644 --- a/packages/das/src/queue/fetch.processor.ts +++ b/packages/das/src/queue/fetch.processor.ts @@ -2,6 +2,7 @@ import { Processor, WorkerHost, InjectQueue } from "@nestjs/bullmq"; import { Logger } from "@nestjs/common"; import { InjectRepository } from "@nestjs/typeorm"; import { IsNull, Repository } from "typeorm"; +import { QueryDeepPartialEntity } from "typeorm/query-builder/QueryPartialEntity"; import { Job, Queue } from "bullmq"; import { Issue, PullRequest } from "../entities"; import { GitHubFetcherService } from "../webhook/github-fetcher.service"; @@ -123,19 +124,47 @@ export class FetchProcessor extends WorkerHost { ): Promise { this.logger.log(`Fetching PR metadata for ${repoFullName}#${prNumber}`); - const { closingIssueNumbers, body, lastEditedAt } = - await this.fetcher.fetchPrMetadata(repoFullName, prNumber); + const { + closingIssueNumbers, + body, + lastEditedAt, + state, + mergedAt, + closedAt, + mergedByLogin, + } = await this.fetcher.fetchPrMetadata(repoFullName, prNumber); const currentClosingIssueNumbers = this.uniqueIssueNumbers(closingIssueNumbers); - await this.prRepo.update( - { repoFullName, prNumber }, - { - closingIssueNumbers: currentClosingIssueNumbers, - body, - lastEditedAt, - }, - ); + // Re-assert authoritative state from GraphQL so a missed + // `pull_request.closed` webhook self-heals (see PrReconcileService, which + // enqueues this job for every still-open PR on a schedule). MERGED is + // terminal: never let an in-flight stale fetch revert a merged PR back to + // OPEN/CLOSED — only forward transitions are applied. + const existing = await this.prRepo.findOne({ + where: { repoFullName, prNumber }, + select: { state: true }, + }); + const applyState = !(existing?.state === "MERGED" && state !== "MERGED"); + + if (applyState && existing && existing.state !== state) { + this.logger.warn( + `State drift corrected for ${repoFullName}#${prNumber}: ` + + `${existing.state} → ${state} (missed webhook)`, + ); + } + + // Cast past the entity's non-null column types: merged_at/closed_at/ + // merged_by_login are nullable in the DB, and writing null is correct + // (clears them on a reopened PR). + const update = { + closingIssueNumbers: currentClosingIssueNumbers, + body, + lastEditedAt, + ...(applyState ? { state, mergedAt, closedAt, mergedByLogin } : {}), + } as QueryDeepPartialEntity; + + await this.prRepo.update({ repoFullName, prNumber }, update); // Issue solver attribution is closure-driven (ISSUE_CLOSURE jobs read // ClosedEvent.closer). PR metadata only refreshes the PR-side text view diff --git a/packages/das/src/webhook/github-fetcher.service.ts b/packages/das/src/webhook/github-fetcher.service.ts index ba8a666..89f9568 100644 --- a/packages/das/src/webhook/github-fetcher.service.ts +++ b/packages/das/src/webhook/github-fetcher.service.ts @@ -287,16 +287,30 @@ export class GitHubFetcherService implements OnModuleInit { closingIssueNumbers: number[]; body: string | null; lastEditedAt: string | null; + state: string; + mergedAt: string | null; + closedAt: string | null; + mergedByLogin: string | null; }> { const [owner, repo] = repoFullName.split("/"); const token = await this.getTokenForRepo(repoFullName); + // `state`/`mergedAt`/`closedAt`/`mergedBy` are returned alongside the body + // so the metadata-fetch path can re-assert authoritative PR state — this is + // what lets a missed `pull_request.closed` webhook self-heal (the webhook + // handler is otherwise the only writer of state). GraphQL `state` is the + // source of truth (OPEN / CLOSED / MERGED), unlike REST which reports a + // merged PR as `closed` + `merged: true`. const query = ` query($owner: String!, $repo: String!, $pr: Int!) { repository(owner: $owner, name: $repo) { pullRequest(number: $pr) { bodyText lastEditedAt + state + mergedAt + closedAt + mergedBy { login } closingIssuesReferences(first: 10) { nodes { number @@ -343,6 +357,10 @@ export class GitHubFetcherService implements OnModuleInit { ), body: pr.bodyText ?? null, lastEditedAt: pr.lastEditedAt ?? null, + state: pr.state, + mergedAt: pr.mergedAt ?? null, + closedAt: pr.closedAt ?? null, + mergedByLogin: pr.mergedBy?.login ?? null, }; } diff --git a/packages/das/src/webhook/handlers/pull-request.handler.ts b/packages/das/src/webhook/handlers/pull-request.handler.ts index e99a531..88431ad 100644 --- a/packages/das/src/webhook/handlers/pull-request.handler.ts +++ b/packages/das/src/webhook/handlers/pull-request.handler.ts @@ -26,6 +26,13 @@ export class PullRequestHandler { const prNumber: number = pr.number; const action: string = payload.action; + // A merged PR is reported by REST as state=closed + merged=true. Key off + // `merged` alone — requiring a non-null `merged_at` too risks pinning the + // PR to OPEN if GitHub sends the closed event before merged_at is populated. + // Synthesize merged_at from closed_at when absent so the merge gate (which + // requires merged_at) downstream still passes. + const isMerged = Boolean(pr.merged); + const data: Partial = { repoFullName, prNumber, @@ -33,10 +40,10 @@ export class PullRequestHandler { authorLogin: pr.user.login, authorAssociation: pr.author_association, title: pr.title, - state: pr.merged && pr.merged_at ? "MERGED" : pr.state.toUpperCase(), + state: isMerged ? "MERGED" : pr.state.toUpperCase(), createdAt: pr.created_at, closedAt: pr.closed_at ?? null, - mergedAt: pr.merged_at ?? null, + mergedAt: isMerged ? (pr.merged_at ?? pr.closed_at ?? null) : null, // last_edited_at is populated by the fetch-pr-metadata job via GraphQL — // REST's updated_at changes on any interaction, not just body edits. mergedByLogin: pr.merged_by?.login ?? null, diff --git a/packages/das/src/webhook/pr-reconcile.service.ts b/packages/das/src/webhook/pr-reconcile.service.ts new file mode 100644 index 0000000..ea755af --- /dev/null +++ b/packages/das/src/webhook/pr-reconcile.service.ts @@ -0,0 +1,90 @@ +import { + Injectable, + Logger, + OnModuleDestroy, + OnModuleInit, +} from "@nestjs/common"; +import { InjectQueue } from "@nestjs/bullmq"; +import { InjectRepository } from "@nestjs/typeorm"; +import { Queue } from "bullmq"; +import { Repository } from "typeorm"; +import { PullRequest } from "../entities"; +import { FETCH_QUEUE, FETCH_JOBS } from "../queue/constants"; + +// PR state (OPEN → MERGED/CLOSED) is written only by the pull_request webhook +// handler. A single missed/dropped `pull_request.closed` delivery therefore +// leaves a merged PR stuck OPEN forever, with no path to recover. This sweep +// closes that gap: on a schedule it re-enqueues a metadata fetch for every +// still-open PR in registered repos, and the metadata handler re-asserts +// authoritative GraphQL state — so missed merge events self-heal. +const RECONCILE_INTERVAL_MS = Number( + process.env.PR_RECONCILE_INTERVAL_MS ?? 60 * 60 * 1000, // hourly +); +// Bound the sweep to the validator's scoring window — older PRs are no longer +// scored, so refreshing them buys nothing and only spends GitHub API budget. +const RECONCILE_WINDOW_DAYS = Number( + process.env.PR_RECONCILE_WINDOW_DAYS ?? 45, +); + +@Injectable() +export class PrReconcileService implements OnModuleInit, OnModuleDestroy { + private readonly logger = new Logger(PrReconcileService.name); + private timer: NodeJS.Timeout | null = null; + + constructor( + @InjectRepository(PullRequest) + private readonly prRepo: Repository, + @InjectQueue(FETCH_QUEUE) + private readonly fetchQueue: Queue, + ) {} + + onModuleInit(): void { + // Run once at startup, then on the interval. + void this.reconcile(); + this.timer = setInterval( + () => void this.reconcile(), + RECONCILE_INTERVAL_MS, + ); + } + + onModuleDestroy(): void { + if (this.timer) clearInterval(this.timer); + } + + private async reconcile(): Promise { + try { + const rows: { repo_full_name: string; pr_number: number }[] = + await this.prRepo.query( + `SELECT p.repo_full_name, p.pr_number + FROM pull_requests p + JOIN repos r ON r.repo_full_name = p.repo_full_name + WHERE p.state = 'OPEN' + AND r.registered = true + AND p.created_at > NOW() - INTERVAL '${RECONCILE_WINDOW_DAYS} days'`, + ); + + this.logger.log( + `Reconciling ${rows.length} open PRs against GitHub ` + + `(window ${RECONCILE_WINDOW_DAYS}d)`, + ); + + for (const row of rows) { + await this.fetchQueue.add( + FETCH_JOBS.PR_METADATA, + { repoFullName: row.repo_full_name, prNumber: row.pr_number }, + { + // Same stable per-PR jobId as the webhook path — a reconcile job + // dedupes against an already-pending webhook-triggered fetch. + jobId: `meta-${row.repo_full_name}-${row.pr_number}`, + removeOnComplete: true, + removeOnFail: true, + attempts: 3, + backoff: { type: "exponential", delay: 5000 }, + }, + ); + } + } catch (err) { + this.logger.error(`Reconcile failed: ${String(err)}`); + } + } +} diff --git a/packages/das/src/webhook/repo-backfill-schedule.service.ts b/packages/das/src/webhook/repo-backfill-schedule.service.ts new file mode 100644 index 0000000..2494605 --- /dev/null +++ b/packages/das/src/webhook/repo-backfill-schedule.service.ts @@ -0,0 +1,97 @@ +import { + Injectable, + Logger, + OnModuleDestroy, + OnModuleInit, +} from "@nestjs/common"; +import { InjectQueue } from "@nestjs/bullmq"; +import { InjectRepository } from "@nestjs/typeorm"; +import { Queue } from "bullmq"; +import { Repository } from "typeorm"; +import { Repo } from "../entities"; +import { + FETCH_QUEUE, + FETCH_JOBS, + DEFAULT_BACKFILL_DAYS, +} from "../queue/constants"; + +// Coarse safety net beneath the per-PR reconcile sweep: periodically re-backfill +// every registered repo from GitHub via GraphQL (authoritative state for PRs + +// issues + labels), catching any drift the targeted open-PR sweep doesn't — +// e.g. issue state, labels, or a PR that was already non-OPEN when last seen. +// Heavier than the reconcile sweep (re-touches every PR in the window), so it +// runs daily and can be disabled on critical infra via env. +const BACKFILL_ENABLED = process.env.NIGHTLY_BACKFILL_ENABLED !== "false"; +const BACKFILL_INTERVAL_MS = Number( + process.env.NIGHTLY_BACKFILL_INTERVAL_MS ?? 24 * 60 * 60 * 1000, // daily +); +const BACKFILL_DAYS = Number( + process.env.NIGHTLY_BACKFILL_DAYS ?? DEFAULT_BACKFILL_DAYS, +); + +@Injectable() +export class RepoBackfillScheduleService + implements OnModuleInit, OnModuleDestroy +{ + private readonly logger = new Logger(RepoBackfillScheduleService.name); + private timer: NodeJS.Timeout | null = null; + + constructor( + @InjectRepository(Repo) + private readonly repoRepo: Repository, + @InjectQueue(FETCH_QUEUE) + private readonly fetchQueue: Queue, + ) {} + + onModuleInit(): void { + if (!BACKFILL_ENABLED) { + this.logger.log( + "Nightly repo backfill disabled (NIGHTLY_BACKFILL_ENABLED=false)", + ); + return; + } + // Unlike the reconcile sweep, don't run at startup — a deploy already + // implies fresh data, and this is the heavy job. Start on the interval. + this.timer = setInterval( + () => void this.backfillAll(), + BACKFILL_INTERVAL_MS, + ); + } + + onModuleDestroy(): void { + if (this.timer) clearInterval(this.timer); + } + + private async backfillAll(): Promise { + try { + const repos = await this.repoRepo.find({ + where: { registered: true }, + select: { repoFullName: true }, + }); + + this.logger.log( + `Enqueuing nightly backfill for ${repos.length} repos ` + + `(last ${BACKFILL_DAYS}d)`, + ); + + for (const repo of repos) { + await this.fetchQueue.add( + FETCH_JOBS.BACKFILL_REPO, + { repoFullName: repo.repoFullName, days: BACKFILL_DAYS }, + { + // Static per-repo jobId so a still-running nightly backfill isn't + // stacked on by the next tick. Distinct from the admin endpoint's + // timestamped ids, so manual backfills are never blocked. + jobId: `backfill-${repo.repoFullName}-nightly`, + removeOnComplete: true, + removeOnFail: true, + attempts: 2, + backoff: { type: "exponential", delay: 30000 }, + }, + ); + } + } catch (err) { + this.logger.error(`Nightly backfill enqueue failed: ${String(err)}`); + } + } +} diff --git a/packages/das/src/webhook/webhook.module.ts b/packages/das/src/webhook/webhook.module.ts index a5e5144..cb7fc10 100644 --- a/packages/das/src/webhook/webhook.module.ts +++ b/packages/das/src/webhook/webhook.module.ts @@ -14,6 +14,8 @@ import { FETCH_QUEUE } from "../queue/constants"; import { WebhookController } from "./webhook.controller"; import { WebhookService } from "./webhook.service"; import { WebhookPruneService } from "./webhook-prune.service"; +import { PrReconcileService } from "./pr-reconcile.service"; +import { RepoBackfillScheduleService } from "./repo-backfill-schedule.service"; import { PullRequestHandler } from "./handlers/pull-request.handler"; import { IssueHandler } from "./handlers/issue.handler"; import { ReviewHandler } from "./handlers/review.handler"; @@ -39,6 +41,8 @@ import { InstallationHandler } from "./handlers/installation.handler"; providers: [ WebhookService, WebhookPruneService, + PrReconcileService, + RepoBackfillScheduleService, PullRequestHandler, IssueHandler, ReviewHandler,