Repository navigation
fix(server): PR sync waits out a GitHub rate limit pause instead of failing every PR #16203
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
5cdc528
e58072f
0ab622a
5d4c748
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -16,16 +16,19 @@ import { | |
| visibleThreadPullRequests, | ||
| } from "@t3tools/shared/threadPullRequests"; | ||
| import * as Cause from "effect/Cause"; | ||
| import * as Clock from "effect/Clock"; | ||
| import * as Context from "effect/Context"; | ||
| import * as Crypto from "effect/Crypto"; | ||
| import * as DateTime from "effect/DateTime"; | ||
| import * as Effect from "effect/Effect"; | ||
| import * as Layer from "effect/Layer"; | ||
| import * as Schema from "effect/Schema"; | ||
| import * as Schedule from "effect/Schedule"; | ||
| import type * as Scope from "effect/Scope"; | ||
| import * as Semaphore from "effect/Semaphore"; | ||
| import * as Stream from "effect/Stream"; | ||
|
|
||
| import { PullRequestProviderError } from "../pullRequest/PullRequestProvider.ts"; | ||
| import * as PullRequestService from "../pullRequest/PullRequestService.ts"; | ||
| import { forkParked } from "../serverActivation.ts"; | ||
| import * as Orchestrator from "./Orchestrator.ts"; | ||
|
|
@@ -36,6 +39,8 @@ const SLOW_SYNC_INTERVAL_MS = 15 * 60 * 1_000; | |
| /** Shell commands that can merge or close a pull request without a merge notification. */ | ||
| const PULL_REQUEST_CLOSE_COMMAND = /\b(?:gh\s+pr|glab\s+mr)\s+(?:merge|close)\b/u; | ||
|
|
||
| const isPullRequestProviderError = Schema.is(PullRequestProviderError); | ||
|
|
||
| type SnapshotFields = Omit<ThreadPullRequestSnapshot, "syncedAt">; | ||
|
|
||
| interface LinkEntry { | ||
|
|
@@ -107,6 +112,21 @@ function stacksEqual( | |
| ); | ||
| } | ||
|
|
||
| function skipReason(cause: Cause.Cause<unknown>): string { | ||
| const error = Cause.squash(cause); | ||
| return error instanceof Error ? error.message : String(error); | ||
| } | ||
|
|
||
| /** When a host read failed because the host is rate limited, the time that pause ends. */ | ||
| function rateLimitRetryAt(cause: Cause.Cause<unknown>): number | undefined { | ||
| let error: unknown = Cause.squash(cause); | ||
| while (error instanceof Error) { | ||
| if (isPullRequestProviderError(error) && error.reason === "rate-limited") return error.retryAt; | ||
| error = error.cause; | ||
| } | ||
| return undefined; | ||
| } | ||
|
|
||
| function isUnsettled(thread: ProjectionStore.ProjectionThreadPullRequests): boolean { | ||
| return thread.settledOverride !== "settled" && thread.settledAt === null; | ||
| } | ||
|
|
@@ -122,7 +142,10 @@ export class PullRequestSyncReactor extends Context.Service< | |
| { | ||
| readonly start: () => Effect.Effect<void, never, Scope.Scope>; | ||
| readonly drain: Effect.Effect<void>; | ||
| /** Force the next sweep to re-read this pull request, even when its snapshot is terminal. */ | ||
| /** | ||
| * Force the next sweep to re-read this pull request, even when its snapshot is terminal. | ||
| * While its host is rate limited, the read waits for the first sweep after the pause. | ||
| */ | ||
| readonly requestSync: (key: ThreadPullRequestKey) => Effect.Effect<void>; | ||
| } | ||
| >()("t3/orchestration-v2/PullRequestSyncReactor") {} | ||
|
|
@@ -141,6 +164,10 @@ export const make = Effect.gen(function* () { | |
| // linking dozens of pull requests) is read together and shares the summary batches. | ||
| let requestedSweepQueued = false; | ||
| const retryStacks = new Set<string>(); | ||
| // Rate limit pauses by project and host, since each project reads with its own credential. | ||
| // A paused host refuses every read without asking it, so the sweep leaves its pull requests | ||
| // due until the pause ends rather than failing each of them every minute. | ||
| const pausedUntil = new Map<string, number>(); | ||
|
|
||
| const isDue = (key: string, entries: ReadonlyArray<LinkEntry>, nowMs: number): boolean => { | ||
| if (requested.has(key) || retryStacks.has(key)) return true; | ||
|
|
@@ -279,10 +306,15 @@ export const make = Effect.gen(function* () { | |
| })), | ||
| Effect.catchCauseIf( | ||
| (cause) => !Cause.hasInterruptsOnly(cause), | ||
| () => | ||
| Effect.logWarning("pull request stack lookup failed", { | ||
| key, | ||
| }).pipe(Effect.as(null)), | ||
| (cause) => | ||
| rateLimitRetryAt(cause) !== undefined | ||
| ? // The sweep records the pause and holds the host's other reads until it ends. | ||
| Effect.sync(() => retryStacks.add(key)).pipe( | ||
| Effect.andThen(Effect.failCause(cause)), | ||
| ) | ||
| : Effect.logWarning("pull request stack lookup failed", { | ||
| key, | ||
| }).pipe(Effect.as(null)), | ||
| ), | ||
| ); | ||
| if (needsStack) { | ||
|
|
@@ -312,18 +344,50 @@ export const make = Effect.gen(function* () { | |
| ); | ||
| }); | ||
|
|
||
| // Failed host reads by reason: how many, and the first key that failed that way. | ||
| const skips = new Map<string, { count: number; readonly key: string }>(); | ||
| const readGroup = (key: string, entries: ReadonlyArray<LinkEntry>, pauseKey: string) => | ||
| syncGroup(key, entries).pipe( | ||
| Effect.catchCause((cause) => { | ||
| if (Cause.hasInterruptsOnly(cause)) return Effect.failCause(cause); | ||
| const retryAt = rateLimitRetryAt(cause); | ||
| if (retryAt !== undefined) { | ||
| pausedUntil.set(pauseKey, Math.max(retryAt, pausedUntil.get(pauseKey) ?? 0)); | ||
| } | ||
| const reason = skipReason(cause); | ||
| const skip = skips.get(reason); | ||
| if (skip === undefined) skips.set(reason, { count: 1, key }); | ||
| else skip.count += 1; | ||
| return Effect.void; | ||
| }), | ||
| ); | ||
| yield* Effect.forEach( | ||
| groups, | ||
| ([key, entries]) => | ||
| (scope === "all" || requested.has(key)) && isDue(key, entries, nowMs) | ||
| ? syncGroup(key, entries).pipe( | ||
| Effect.catchCause(logSkipped("pull request sync skipped", { key })), | ||
| ) | ||
| : Effect.void, | ||
| ([key, entries]) => { | ||
| if (!((scope === "all" || requested.has(key)) && isDue(key, entries, nowMs))) { | ||
| return Effect.void; | ||
| } | ||
| const first = entries[0]!; | ||
| const pauseKey = `${first.thread.projectId}\0${normalizeThreadPullRequestKey(first.link).host}`; | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🟡 Medium
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Note 🤖 Claude Opus 5.5 responding on behalf of Theo Leaving this as is. Every project records its own pause on its first refused read: |
||
| // Checked against the clock as each read starts, so a pause found earlier in this sweep | ||
| // holds the rest, and one that ends during the sweep lets the rest through. | ||
| return Clock.currentTimeMillis.pipe( | ||
| Effect.flatMap((startedAtMs) => | ||
| (pausedUntil.get(pauseKey) ?? 0) > startedAtMs | ||
| ? Effect.void | ||
| : readGroup(key, entries, pauseKey), | ||
| ), | ||
| ); | ||
| }, | ||
| // As wide as one batched summary read, so the sweep's reads on a host arrive together and | ||
| // GitHub answers them in one request rather than one `gh pr view` apiece. | ||
| { concurrency: 25, discard: true }, | ||
| ); | ||
| // A host failure such as a signed-out CLI fails every due pull request the same way, so a | ||
| // sweep reports one line per reason rather than one per pull request. | ||
| for (const [reason, { count, key }] of skips) { | ||
| yield* Effect.logWarning("pull request sync skipped", { count, key, reason }); | ||
| } | ||
| }); | ||
|
|
||
| const worker = yield* makeDrainableWorker((scope: "all" | "requested") => | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.