Skip to content
Merged
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
131 changes: 131 additions & 0 deletions apps/server/src/orchestration-v2/PullRequestSyncReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -908,6 +911,134 @@ describe("PullRequestSyncReactor", () => {
),
);

it.effect("leaves a rate limited host unread until its pause ends", () => {
const skips: Array<ReadonlyArray<unknown>> = [];
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* () {
Expand Down
86 changes: 75 additions & 11 deletions apps/server/src/orchestration-v2/PullRequestSyncReactor.ts
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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 {
Expand Down Expand Up @@ -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;
}
Expand All @@ -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") {}
Expand All @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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}`;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Medium orchestration-v2/PullRequestSyncReactor.ts:371

pausedUntil only suppresses reads for the current projectId, so when two projects share a GitHub credential, a rate limit discovered in one project does not suppress the other project's due groups. Those groups continue calling summary, are rejected by the shared limiter, and emit skipped-read warnings throughout the pause. Key this cache by the credential scope (or conservatively by provider/host) instead of projectId.

🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @apps/server/src/orchestration-v2/PullRequestSyncReactor.ts around line 371:

`pausedUntil` only suppresses reads for the current `projectId`, so when two projects share a GitHub credential, a rate limit discovered in one project does not suppress the other project's due groups. Those groups continue calling `summary`, are rejected by the shared limiter, and emit skipped-read warnings throughout the pause. Key this cache by the credential scope (or conservatively by provider/host) instead of `projectId`.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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: readGroup sets pausedUntil for that projectId + host. So a second project on the same credential makes one more read, which the shared limiter refuses locally without a GitHub request, and then skips until retryAt like the first. That is one fast local failure per project per pause, not one per sweep. Keying by credential would mean depending on how SourceControlRateLimit scopes credentials, and that coupling costs more than the one extra local read it saves.

// 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") =>
Expand Down
Loading