diff --git a/apps/server/src/orchestration-v2/PullRequestSyncReactor.test.ts b/apps/server/src/orchestration-v2/PullRequestSyncReactor.test.ts index 6febdb9bbce3..256544e99312 100644 --- a/apps/server/src/orchestration-v2/PullRequestSyncReactor.test.ts +++ b/apps/server/src/orchestration-v2/PullRequestSyncReactor.test.ts @@ -19,15 +19,18 @@ import { type ThreadPullRequestSnapshot, } from "@t3tools/contracts"; import { assert, describe, it } from "@effect/vitest"; +import * as Clock from "effect/Clock"; import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Logger from "effect/Logger"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import { TestClock } from "effect/testing"; +import { PullRequestProviderError } from "../pullRequest/PullRequestProvider.ts"; import * as PullRequestService from "../pullRequest/PullRequestService.ts"; import * as ServerActivation from "../serverActivation.ts"; import * as Orchestrator from "./Orchestrator.ts"; @@ -908,6 +911,134 @@ describe("PullRequestSyncReactor", () => { ), ); + it.effect("leaves a rate limited host unread until its pause ends", () => { + const skips: Array> = []; + const logger = Logger.make(({ logLevel, message }) => { + const parts = Array.isArray(message) ? message : [message]; + if (logLevel === "Warn" && parts[0] === "pull request sync skipped") skips.push(parts); + }); + const retryAt = Date.parse(NOW) + 3 * 60_000; + return Effect.scoped( + Effect.gen(function* () { + yield* TestClock.setTime(Date.parse(NOW)); + const fixture = yield* makeHarness({ + snapshot: makeSnapshot([ + makeThread("first", { pullRequests: [makeLink(7, { state: "open" })] }), + makeThread("second", { pullRequests: [makeLink(8, { state: "open" })] }), + makeThread("third", { + pullRequests: [ + makeLink( + 9, + { state: "open" }, + { host: "forge.example", url: "https://forge.example/owner/repository/pulls/9" }, + ), + ], + }), + ]), + summary: (input) => + Effect.gen(function* () { + if (input.host === "github.com" && (yield* Clock.currentTimeMillis) < retryAt) { + const paused = new PullRequestProviderError({ + provider: "github", + operation: "getChangeRequestSummary", + reason: "rate-limited", + detail: "paused", + retryAt, + }); + return yield* new PullRequestOperationError({ + operation: "summary", + detail: "paused", + cause: paused, + }); + } + return makeSummary(input, { state: "open" }); + }), + }); + const githubReads = Ref.get(fixture.summaryCalls).pipe( + Effect.map((calls) => + calls.filter((call) => call.host === "github.com").map((call) => call.number), + ), + ); + + yield* Effect.gen(function* () { + // The first refused read pauses the host, so the sweep does not try the other. + const reactor = yield* startAndSweep(fixture); + assert.deepStrictEqual(yield* githubReads, [7]); + assert.strictEqual(skips.length, 1); + assert.deepInclude(skips[0]![1], { count: 1 }); + + // A requested refresh waits for the pause like the sweep does. + yield* reactor.requestSync({ + host: "github.com", + repository: "owner/repository", + number: 7, + }); + yield* Queue.take(fixture.snapshotReads); + yield* reactor.drain; + yield* sweepAgain(fixture, reactor); + yield* sweepAgain(fixture, reactor); + assert.deepStrictEqual(yield* githubReads, [7]); + assert.strictEqual(skips.length, 1); + // Other hosts keep their cadence. + assert.strictEqual((yield* Ref.get(fixture.summaryCalls)).length, 4); + + yield* sweepAgain(fixture, reactor); + assert.sameMembers((yield* githubReads).slice(1), [7, 8]); + }).pipe( + // The reactor forks its worker while its layer builds, so the logger must reach it there. + Effect.provide( + fixture.layer.pipe(Layer.provide(Logger.layer([logger], { mergeWithExisting: false }))), + ), + ); + }), + ); + }); + + it.effect("pauses the host when only its stack read is rate limited", () => + Effect.scoped( + Effect.gen(function* () { + yield* TestClock.setTime(Date.parse(NOW)); + const retryAt = Date.parse(NOW) + 3 * 60_000; + const fixture = yield* makeHarness({ + snapshot: makeSnapshot([makeThread("one", { pullRequests: [makeLink(7, null)] })]), + summary: (input) => Effect.succeed(makeSummary(input, { state: "open" })), + stack: () => + Effect.gen(function* () { + if ((yield* Clock.currentTimeMillis) >= retryAt) return null; + return yield* new PullRequestOperationError({ + operation: "stack", + detail: "paused", + cause: new PullRequestProviderError({ + provider: "github", + operation: "getChangeRequestStack", + reason: "rate-limited", + detail: "paused", + retryAt, + }), + }); + }), + }); + const reads = Effect.all([ + Ref.get(fixture.summaryCalls).pipe(Effect.map((calls) => calls.length)), + Ref.get(fixture.stackCalls).pipe(Effect.map((calls) => calls.length)), + ]); + + yield* Effect.gen(function* () { + const reactor = yield* startAndSweep(fixture); + assert.deepStrictEqual(yield* reads, [1, 1]); + yield* sweepAgain(fixture, reactor); + yield* sweepAgain(fixture, reactor); + assert.deepStrictEqual(yield* reads, [1, 1]); + + // At retryAt the pull request is read again, stack included. + yield* sweepAgain(fixture, reactor); + assert.deepStrictEqual(yield* reads, [2, 2]); + assert.strictEqual((yield* Ref.get(fixture.syncCommands)).length, 1); + }).pipe(Effect.provide(fixture.layer)); + }), + ), + ); + it.effect("reads open links fresh only when a run that ran a merge command ends", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/PullRequestSyncReactor.ts b/apps/server/src/orchestration-v2/PullRequestSyncReactor.ts index 1b526d8f532e..c2862fa7f1cd 100644 --- a/apps/server/src/orchestration-v2/PullRequestSyncReactor.ts +++ b/apps/server/src/orchestration-v2/PullRequestSyncReactor.ts @@ -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; interface LinkEntry { @@ -107,6 +112,21 @@ function stacksEqual( ); } +function skipReason(cause: Cause.Cause): 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): 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; readonly drain: Effect.Effect; - /** 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; } >()("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(); + // 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(); const isDue = (key: string, entries: ReadonlyArray, 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(); + const readGroup = (key: string, entries: ReadonlyArray, 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}`; + // 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") =>