diff --git a/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts b/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts index 291105e25cd9..5260f3ca3f02 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts @@ -65,6 +65,7 @@ const layer = it.layer( }), ), Layer.provideMerge(GitHubGraphQlBudget.layer), + Layer.provideMerge(SourceControlRateLimit.layer), ), ); @@ -247,7 +248,7 @@ it.effect( ); const cli = yield* GitHubPullRequestCli.make.pipe( Effect.provideService(GitHubCli.GitHubCli, github), - Effect.provide(GitHubGraphQlBudget.layer), + Effect.provide(Layer.merge(GitHubGraphQlBudget.layer, SourceControlRateLimit.layer)), ); const input = { cwd: "/repo", host: "github.com" }; const first = yield* cli.withVerifiedCredential(input, (identity) => @@ -285,6 +286,152 @@ it.effect( }), ); +it.effect("does not verify a paused credential over the network and resumes after cooldown", () => + Effect.gen(function* () { + const limits = yield* SourceControlRateLimit.make; + let token = "token-a"; + const verified: string[] = []; + const cli = yield* GitHubPullRequestCli.make.pipe( + Effect.provideService(SourceControlRateLimit.SourceControlRateLimit, limits), + Effect.provide( + Layer.merge( + GitHubGraphQlBudget.layer, + Layer.mock(GitHubCli.GitHubCli)({ + execute: (input) => + Effect.sync(() => { + if (input.args[0] === "auth") return output(token); + verified.push(token); + return output('{"id":123,"login":"viewer"}'); + }), + }), + ), + ), + ); + const input = { cwd: "/repo", host: "github.com" }; + yield* cli.withVerifiedCredential(input, () => + limits.recordRateLimit({ + provider: "github", + host: "github.com", + lease: 0, + retryAt: 20 * 60_000, + }), + ); + // A cached identity does not spend quota and remains available to interactive routing. + assert.strictEqual(yield* cli.getViewerLogin(input), "viewer"); + // Expire the ten-minute identity cache while this account remains paused. + yield* TestClock.adjust("11 minutes"); + const paused = yield* Effect.flip(cli.getViewerLogin(input)); + assert.strictEqual(paused._tag, "SourceControlRateLimitPausedError"); + assert.deepStrictEqual(verified, ["token-a"]); + token = "token-b"; + assert.strictEqual(yield* cli.getViewerLogin(input), "viewer"); + assert.deepStrictEqual(verified, ["token-a", "token-b"]); + token = "token-a"; + yield* TestClock.adjust("10 minutes"); + assert.strictEqual(yield* cli.getViewerLogin(input), "viewer"); + assert.deepStrictEqual(verified, ["token-a", "token-b", "token-a"]); + }), +); + +it.effect.each([false, true])( + "pins workspace credentials for concurrent summaries, fallback=%s", + (fallback) => + Effect.gen(function* () { + const commands: VcsProcess.VcsProcessInput[] = []; + const github = yield* GitHubCli.make.pipe( + Effect.provide(Layer.merge(GitHubGraphQlBudget.layer, SourceControlRateLimit.layer)), + Effect.provideService(VcsProcess.VcsProcess, { + run: (input) => + Effect.sync(() => { + commands.push(input); + const token = input.cwd === "/b" ? "token-b" : "token-a"; + if (input.args[0] === "auth") return output(token); + if (input.args[1] === "user") return output('{"id":123,"login":"same-viewer"}'); + if (input.args.includes("rate_limit")) { + expect(["token-a", "token-b"]).toContain(input.env?.GH_TOKEN); + return output( + encodeJson({ + data: { + rateLimit: { + cost: 1, + limit: 5000, + remaining: 4999, + resetAt: "2099-01-01T00:00:00Z", + }, + }, + }), + ); + } + expect(input.env?.GH_TOKEN).toBe(token); + if (input.args[0] === "pr") + return output( + encodeJson({ + ...coreResponse().data.repository.pullRequest, + number: Number(input.args[2]), + author: null, + isDraft: false, + additions: 1, + deletions: 0, + changedFiles: 1, + mergedAt: null, + closedAt: null, + reviewDecision: null, + mergeable: "MERGEABLE", + reviewRequests: [], + labels: [], + statusCheckRollup: [], + body: "", + }), + ); + const query = input.args.find((arg) => arg.startsWith("query=")) ?? ""; + const data: Record = {}; + for (const match of query.matchAll( + /(s\d+): repository\(owner: "([^"]+)", name: "([^"]+)"\) \{ pullRequest\(number: (\d+)\)/g, + )) { + expect(match[2] === "b" ? "token-b" : "token-a").toBe(token); + data[match[1]!] = { + pullRequest: { + ...coreResponse().data.repository.pullRequest, + number: Number(match[4]), + author: null, + isDraft: false, + additions: 1, + deletions: 0, + changedFiles: 1, + mergedAt: null, + closedAt: null, + reviewDecision: null, + mergeable: "MERGEABLE", + }, + }; + } + return output(encodeJson({ data: fallback ? {} : data })); + }).pipe(Effect.delay(input.args[0] === "auth" ? (input.cwd === "/b" ? 40 : 20) : 0)), + }), + ); + const cli = yield* GitHubPullRequestCli.make.pipe( + Effect.provideService(GitHubCli.GitHubCli, github), + Effect.provide(Layer.merge(GitHubGraphQlBudget.layer, SourceControlRateLimit.layer)), + ); + const reads = yield* Effect.forEach( + ["a", "b", "c", "a"], + (name, index) => + cli.getPullRequestSummary({ + cwd: `/${name}`, + repository: `${name}/repo`, + host: "github.com", + number: index + 1, + }), + { concurrency: "unbounded" }, + ).pipe(Effect.forkChild); + yield* TestClock.adjust("100 millis"); + expect((yield* Fiber.join(reads)).map((row) => row.number)).toEqual([1, 2, 3, 4]); + expect(commands.filter((command) => command.args[1] === "graphql")).toHaveLength(2); + expect(commands.filter((command) => command.args[0] === "auth")).toHaveLength(3); + expect(commands.filter((command) => command.args[0] === "pr")).toHaveLength(fallback ? 4 : 0); + }), +); + layer("GitHubPullRequestCli.layer", (it) => { it.effect("admits only one concurrent preview above the reserve and resumes after reset", () => Effect.gen(function* () { @@ -474,13 +621,16 @@ layer("GitHubPullRequestCli.layer", (it) => { closedAt: null, commits: { nodes: [{ commit: { statusCheckRollup: { state: "SUCCESS" } } }] }, }); - mockedExecute.mockReturnValueOnce( + mockedExecute.mockImplementation((input) => Effect.succeed( output( - // @effect-diagnostics-next-line preferSchemaOverJson:off - JSON.stringify({ - data: { s0: { pullRequest: node(7) }, s1: { pullRequest: node(8) } }, - }), + input.args[0] === "auth" + ? "summary-credential" + : input.args[1] === "user" + ? '{"id":123,"login":"viewer"}' + : encodeJson({ + data: { s0: { pullRequest: node(7) }, s1: { pullRequest: node(8) } }, + }), ), ), ); @@ -523,8 +673,9 @@ layer("GitHubPullRequestCli.layer", (it) => { }, ); assert.strictEqual(eight?.headBranch, "feat/8"); - expect(mockedExecute).toHaveBeenCalledOnce(); - const document = callAt(0).args.at(-1) ?? ""; + const graphql = mockedExecute.mock.calls.filter(([input]) => input.args[1] === "graphql"); + expect(graphql).toHaveLength(1); + const document = graphql[0]![0].args.at(-1) ?? ""; expect(document).toContain( 's0: repository(owner: "acme", name: "web") { pullRequest(number: 7)', ); @@ -534,43 +685,46 @@ layer("GitHubPullRequestCli.layer", (it) => { it.effect("reads a pull request the batch said nothing about on its own", () => Effect.gen(function* () { - mockedExecute - .mockReturnValueOnce(Effect.succeed(output('{"data":{"s0":{"pullRequest":null}}}'))) - .mockReturnValueOnce( - Effect.succeed( - output( - // @effect-diagnostics-next-line preferSchemaOverJson:off - JSON.stringify({ - number: 7, - title: "Reuse the summary", - url: "https://github.com/acme/web/pull/7", - author: { login: "octocat", name: "Octo Cat" }, - baseRefName: "main", - headRefName: "feat/summary", - state: "OPEN", - isDraft: false, - mergeable: "MERGEABLE", - reviewDecision: "APPROVED", - additions: 12, - deletions: 3, - changedFiles: 2, - createdAt: "2026-08-20T00:00:00.000Z", - updatedAt: "2026-08-24T12:34:56.000Z", - reviewRequests: [], - labels: [], - statusCheckRollup: [ - { - __typename: "CheckRun", - status: "COMPLETED", - conclusion: "SUCCESS", - name: "ci", - }, - ], - body: "", - }), - ), + mockedExecute.mockImplementation((input) => + Effect.succeed( + output( + input.args[0] === "auth" + ? "summary-credential" + : input.args[1] === "user" + ? '{"id":123,"login":"viewer"}' + : input.args[1] === "graphql" + ? '{"data":{"s0":{"pullRequest":null}}}' + : encodeJson({ + number: 7, + title: "Reuse the summary", + url: "https://github.com/acme/web/pull/7", + author: { login: "octocat", name: "Octo Cat" }, + baseRefName: "main", + headRefName: "feat/summary", + state: "OPEN", + isDraft: false, + mergeable: "MERGEABLE", + reviewDecision: "APPROVED", + additions: 12, + deletions: 3, + changedFiles: 2, + createdAt: "2026-08-20T00:00:00.000Z", + updatedAt: "2026-08-24T12:34:56.000Z", + reviewRequests: [], + labels: [], + statusCheckRollup: [ + { + __typename: "CheckRun", + status: "COMPLETED", + conclusion: "SUCCESS", + name: "ci", + }, + ], + body: "", + }), ), - ); + ), + ); const cli = yield* GitHubPullRequestCli.GitHubPullRequestCli; const read = yield* cli @@ -581,8 +735,9 @@ layer("GitHubPullRequestCli.layer", (it) => { assert.strictEqual(summary.headBranch, "feat/summary"); assert.strictEqual(summary.checksState, "passing"); - assert.strictEqual(mockedExecute.mock.calls.length, 2); - expect(callAt(1).args).toEqual([ + const views = mockedExecute.mock.calls.filter(([input]) => input.args[0] === "pr"); + expect(views).toHaveLength(1); + expect(views[0]![0].args).toEqual([ "pr", "view", "7", diff --git a/apps/server/src/pullRequest/GitHubPullRequestCli.ts b/apps/server/src/pullRequest/GitHubPullRequestCli.ts index a609fd3d097b..a182dac862a5 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestCli.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestCli.ts @@ -1086,9 +1086,13 @@ function actionArgs( } } -/** @public Service construction is part of the canonical Effect module API. */ +/** + * Construct GitHub pull request reads and mutations with credential-pinned batching. + * @public Service construction is part of the canonical Effect module API. + */ export const make = Effect.gen(function* () { const github = yield* GitHubCli.GitHubCli; + const rateLimits = yield* SourceControlRateLimit.SourceControlRateLimit; const graphQlBudget = yield* GitHubGraphQlBudget.GitHubGraphQlBudget; const routingIdentities = new Map< string, @@ -1143,6 +1147,11 @@ export const make = Effect.gen(function* () { const cached = routingIdentities.get(key); if (cached !== undefined && now - cached.at < 10 * 60_000) return { ...credential, ...cached.value }; + // Cached verification spends no quota. A cold verification respects this + // account's pause before calling the network, just like repository reads. + yield* rateLimits + .check({ provider: "github", host }) + .pipe(Effect.provideService(SourceControlRateLimit.CredentialScope, key)); // Pin this read so an auth switch cannot poison its cache entry. const response = yield* github .execute({ @@ -1808,6 +1817,62 @@ export const make = Effect.gen(function* () { * `gh pr view` apiece is most of what it spends. Whatever the batch cannot answer — a selector * GraphQL cannot address, a pull request GitHub returned nothing for — is read on its own. */ + const resolveSummaryBatch = (entries: ReadonlyArray>) => { + const first = entries[0]!; + const batchable = entries.filter( + (entry) => buildPullRequestSummariesGraphQlQuery([entry.request]) !== null, + ); + const query = buildPullRequestSummariesGraphQlQuery(batchable.map((entry) => entry.request)); + const batched = + query === null + ? Effect.succeed(new Map()) + : graphqlRead({ + cwd: first.request.cwd, + host: first.request.host, + operation: "getPullRequestSummary", + query, + decode: decodePullRequestSummariesJson, + }); + return batched.pipe( + // A GraphQL error anywhere fails the whole document — one repository gone or out of + // reach — so a batch that could not be read leaves every entry to its own read. A paused + // budget is the exception: reading one at a time would only spend what is being saved. + Effect.catchCauseIf( + (cause) => + !Cause.hasInterruptsOnly(cause) && + !Cause.findErrorOption(cause).pipe( + Option.exists((error) => error._tag === "SourceControlRateLimitPausedError"), + ), + (cause) => + Effect.logDebug("batched pull request summary read failed", { cause }).pipe( + Effect.as(new Map()), + ), + ), + Effect.flatMap((summaries) => { + const unanswered = entries.filter((entry) => { + const summary = summaries.get(batchable.indexOf(entry)); + if (summary === undefined) return true; + entry.completeUnsafe(Exit.succeed(summary)); + return false; + }); + return Effect.forEach( + unanswered, + (entry) => + viewPullRequestSummary(entry.request).pipe( + Effect.exit, + Effect.map((exit) => entry.completeUnsafe(exit)), + ), + { concurrency: STAT_REQUEST_CONCURRENCY, discard: true }, + ); + }), + Effect.catchCause((cause) => + Effect.sync(() => { + for (const entry of entries) entry.completeUnsafe(Exit.failCause(cause)); + }), + ), + ); + }; + const summaryResolver = RequestResolver.makeGrouped({ key: ({ request, context }) => JSON.stringify([ @@ -1816,65 +1881,58 @@ export const make = Effect.gen(function* () { ?.credentialFingerprint ?? null, Context.getOrElse(context, SourceControlRateLimit.CredentialScope, () => ""), ]), - resolver: (entries) => { - const [first] = entries; - const batchable = entries.filter( - (entry) => buildPullRequestSummariesGraphQlQuery([entry.request]) !== null, - ); - const query = buildPullRequestSummariesGraphQlQuery(batchable.map((entry) => entry.request)); - const batched = - query === null - ? Effect.succeed(new Map()) - : graphqlRead({ - cwd: first.request.cwd, - host: first.request.host, - operation: "getPullRequestSummary", - query, - decode: decodePullRequestSummariesJson, - }); - return batched.pipe( - // A GraphQL error anywhere fails the whole document — one repository gone or out of - // reach — so a batch that could not be read leaves every entry to its own read. A paused - // budget is the exception: reading one at a time would only spend what is being saved. - Effect.catchCauseIf( - (cause) => - !Cause.hasInterruptsOnly(cause) && - !Cause.findErrorOption(cause).pipe( - Option.exists((error) => error._tag === "SourceControlRateLimitPausedError"), - ), - (cause) => - Effect.logDebug("batched pull request summary read failed", { cause }).pipe( - Effect.as(new Map()), + resolver: (entries) => + Effect.gen(function* () { + const byWorkspace = new Map>>(); + for (const entry of entries) { + const held = byWorkspace.get(entry.request.cwd); + if (held === undefined) byWorkspace.set(entry.request.cwd, [entry]); + else held.push(entry); + } + // Gather first, then verify once per workspace. Process latency must not split a + // sweep into one GraphQL request (and one auth process) per linked pull request. + const captured = yield* Effect.forEach( + [...byWorkspace.values()], + (group) => + captureVerifiedCredential(group[0]!.request).pipe( + Effect.provideContext(group[0]!.context), + Effect.map((credential) => ({ entries: group, credential })), + Effect.catchCause((cause) => + Effect.sync(() => { + for (const entry of group) entry.completeUnsafe(Exit.failCause(cause)); + return null; + }), + ), ), - ), - Effect.flatMap((summaries) => { - const unanswered = entries.filter((entry) => { - const summary = summaries.get(batchable.indexOf(entry)); - if (summary === undefined) return true; - entry.completeUnsafe(Exit.succeed(summary)); - return false; - }); - return Effect.forEach( - unanswered, - (entry) => - viewPullRequestSummary(entry.request).pipe( - Effect.exit, - Effect.map((exit) => entry.completeUnsafe(exit)), + { concurrency: STAT_REQUEST_CONCURRENCY }, + ); + const byCredential = new Map>(); + for (const group of captured) { + if (group === null) continue; + const key = group.credential.credentialFingerprint; + const held = byCredential.get(key); + if (held === undefined) byCredential.set(key, group); + else held.entries.push(...group.entries); + } + yield* Effect.forEach( + [...byCredential.values()], + ({ entries: group, credential }) => + resolveSummaryBatch(group).pipe( + Effect.provideService(GitHubCli.PinnedGitHubCredential, credential), + Effect.provideService( + SourceControlRateLimit.CredentialScope, + credential.credentialFingerprint, ), - { concurrency: STAT_REQUEST_CONCURRENCY, discard: true }, - ); - }), - Effect.catchCause((cause) => - Effect.sync(() => { - for (const entry of entries) entry.completeUnsafe(Exit.failCause(cause)); - }), - ), - ); - }, + Effect.provideContext(group[0]!.context), + ), + { concurrency: STAT_REQUEST_CONCURRENCY, discard: true }, + ); + }), }).pipe( RequestResolver.setDelay(SUMMARY_BATCH_WINDOW), RequestResolver.batchN(STAT_ALIASES_PER_REQUEST), ); + /** Queue summaries before capturing credentials so a sweep keeps its batching window. */ const getPullRequestSummary: GitHubPullRequestCli["Service"]["getPullRequestSummary"] = (input) => Effect.request(new PullRequestSummaryRead(input), summaryResolver); diff --git a/apps/server/src/pullRequest/PullRequestService.test.ts b/apps/server/src/pullRequest/PullRequestService.test.ts index 4ba438a97da9..9a61d28bf156 100644 --- a/apps/server/src/pullRequest/PullRequestService.test.ts +++ b/apps/server/src/pullRequest/PullRequestService.test.ts @@ -8,6 +8,7 @@ import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Stream from "effect/Stream"; +import * as Schema from "effect/Schema"; import * as TestClock from "effect/testing/TestClock"; import type { OrchestrationProjectShell, @@ -33,11 +34,25 @@ import { import { PullRequestProviderRegistry, fromProviders } from "./PullRequestProviderRegistry.ts"; import * as PullRequestService from "./PullRequestService.ts"; import * as PullRequestReadCache from "./PullRequestReadCache.ts"; +import * as GitHubCli from "../sourceControl/GitHubCli.ts"; +import * as GitHubGraphQlBudget from "../sourceControl/githubGraphQlBudget.ts"; +import * as VcsProcess from "../vcs/VcsProcess.ts"; +import * as GitHubPullRequestCli from "./GitHubPullRequestCli.ts"; +import * as GitHubPullRequestProvider from "./GitHubPullRequestProvider.ts"; import { FILE_REVISIONS_CACHE_CAPACITY, MAX_FILE_REVISION_PATHS, } from "./pullRequestViewedFiles.ts"; +const encodeProcessResponse = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown)); +const decodeSearchRequest = Schema.decodeSync( + Schema.fromJsonString( + Schema.Struct({ + variables: Schema.Struct({ q: Schema.String }), + }), + ), +); + function project(input: { readonly id: string; readonly title: string; @@ -95,6 +110,7 @@ function changeRequest(number: number, updatedAt: string): ProviderChangeRequest }; } +/** Build a readable pull request with action permissions for service workflow tests. */ function hostedChangeRequest(body: string, additions = 1) { return { ...changeRequest(1, "2026-07-02T00:00:00Z"), @@ -116,6 +132,608 @@ function hostedChangeRequest(body: string, additions = 1) { }; } +it.effect.each(["list", "stats"] as const)( + "keeps %s batches within each workspace credential", + (operation) => + Effect.gen(function* () { + const credential = SourceControlRateLimit.CredentialScope; + const batches: Array<{ credential: string; repositories: ReadonlyArray }> = []; + const projects = ["a", "b", "c"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `org-${name}/repo`, + }), + ); + const accountFor = (cwd: string) => (cwd === "/b" ? "bob" : "alice"); + const service = yield* makeService({ + projects, + providers: [ + fakeProvider("github", { + withVerifiedCredential: ({ cwd }, use) => { + const account = accountFor(cwd); + return use({ + accountId: account, + viewer: account, + credentialFingerprint: account, + }).pipe(Effect.provideService(credential, account)); + }, + getViewer: () => Effect.succeed("alice"), + listChangeRequestsAcross: (input) => + Effect.gen(function* () { + const pinned = yield* credential; + batches.push({ credential: pinned, repositories: input.repositories }); + assert.strictEqual(pinned, accountFor(input.cwd)); + assert.strictEqual(input.viewer, pinned); + assert.isTrue( + input.repositories.every((repo) => accountFor(`/${repo.slice(4, 5)}`) === pinned), + ); + return { + items: input.repositories.map((repository) => ({ + ...batchedChangeRequest(1, repository, "2026-07-02T00:00:00Z"), + author: { login: pinned, name: null, avatarUrl: null }, + })), + truncated: false, + }; + }), + listChangeRequestStats: (input) => + Effect.gen(function* () { + const pinned = yield* credential; + assert.strictEqual(pinned, accountFor(input.cwd)); + assert.isTrue( + input.changeRequests.every( + ({ repository }) => accountFor(`/${repository.slice(4, 5)}`) === pinned, + ), + ); + return input.changeRequests.map((ref) => ({ + ...ref, + additions: pinned === "bob" ? 2 : 1, + deletions: 0, + })); + }), + }), + ], + }); + if (operation === "list") { + const listed = yield* service.list({ state: "open", filters: { author: "@me" } }); + assert.strictEqual(listed.entries.length, 3); + assert.strictEqual(listed.viewers["a github.com"], "alice"); + assert.strictEqual(listed.viewers["b github.com"], "bob"); + assert.deepStrictEqual( + batches.toSorted((a, b) => a.credential.localeCompare(b.credential)), + [ + { credential: "alice", repositories: ["org-a/repo", "org-c/repo"] }, + { credential: "bob", repositories: ["org-b/repo"] }, + ], + ); + return; + } + const counts = yield* service.listStats({ + refs: projects.map((p) => ({ + projectId: p.id, + repository: `org-${p.id}/repo`, + number: 1, + })), + }); + assert.deepStrictEqual(counts.stats.map((s) => [s.projectId, s.additions]).sort(), [ + ["a", 1], + ["b", 2], + ["c", 1], + ]); + }), +); + +it.effect("prepares only cursor repositories while retaining workspace counts", () => + Effect.gen(function* () { + const captures: string[] = []; + const service = yield* makeService({ + projects: ["a", "b"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `${name}/repo`, + }), + ), + providers: [ + fakeProvider("github", { + withVerifiedCredential: ({ cwd }, use) => + Effect.suspend(() => { + captures.push(cwd); + if (cwd === "/a") return Effect.die("Unrelated credentials must not be requested"); + return use({ accountId: cwd, viewer: cwd, credentialFingerprint: cwd }); + }), + listChangeRequests: () => + Effect.succeed({ + items: [changeRequest(2, "2026-07-01T00:00:00Z")], + truncated: false, + continues: false, + }), + }), + ], + }); + const result = yield* service.list({ + state: "open", + cursors: { + "github.com b/repo": "2026-07-02T00:00:00Z|1|1", + }, + }); + assert.deepStrictEqual(captures, ["/b"]); + assert.deepStrictEqual( + result.entries.map((entry) => entry.projectId), + ["b"], + ); + assert.strictEqual(result.providers[0]?.projectCount, 2); + assert.strictEqual(result.viewers["b github.com"], "/b"); + }), +); + +it.effect("retains unrelated host summaries without credential lookups during continuation", () => + Effect.gen(function* () { + const captures: string[] = []; + const service = yield* makeService({ + projects: ["a", "b"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `${name}/repo`, + host: `${name}.example.com`, + }), + ), + providers: [ + fakeProvider("github", { + withVerifiedCredential: ({ cwd }, use) => + Effect.suspend(() => { + captures.push(cwd); + return use({ accountId: cwd, viewer: cwd, credentialFingerprint: cwd }); + }), + listChangeRequests: () => + Effect.succeed({ + items: [changeRequest(2, "2026-07-01T00:00:00Z")], + truncated: false, + continues: false, + }), + }), + ], + }); + const first = yield* service.list({ state: "open" }); + assert.deepStrictEqual( + first.providers.map((host) => host.configured), + [true, true], + ); + captures.length = 0; + const next = yield* service.list({ + state: "open", + cursors: { + "b.example.com b/repo": "2026-07-02T00:00:00Z|1|1", + }, + }); + assert.deepStrictEqual(captures, ["/b"]); + assert.deepStrictEqual(next.providers, first.providers); + yield* service.invalidate({}); + captures.length = 0; + const cold = yield* service.list({ + state: "open", + cursors: { + "b.example.com b/repo": "2026-07-02T00:00:00Z|1|1", + }, + }); + assert.deepStrictEqual(captures, ["/b"]); + assert.strictEqual( + cold.providers.find((host) => host.host === "a.example.com")?.detail, + "Sign-in was not checked on this continuation page.", + ); + }), +); + +it.effect("retains whole-host viewer failures for providers without pinned credentials", () => + Effect.gen(function* () { + const service = yield* makeService({ + projects: ["a", "b", "c"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `${name}/repo`, + host: name === "c" ? "github.example.com" : "github.com", + }), + ), + providers: [ + fakeProvider("github", { + getViewer: ({ host }) => + host === "github.com" ? Effect.fail(requestFailed) : Effect.succeed("viewer"), + }), + ], + }); + const result = yield* service.list({ + state: "open", + cursors: { + "github.com a/repo": "2026-07-02T00:00:00Z|1|1", + "github.example.com c/repo": "2026-07-02T00:00:00Z|1|1", + }, + }); + assert.strictEqual( + result.providers.find((host) => host.host === "github.com")?.detail, + requestFailed.detail, + ); + }), +); + +it.effect("marks a paused credential unreadable without pausing another account and recovers", () => + Effect.gen(function* () { + const reads: string[] = []; + let limited = true; + const service = yield* makeService({ + projects: ["a", "b"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `${name}/repo`, + }), + ), + providers: [ + fakeProvider("github", { + withVerifiedCredential: ({ cwd }, use) => + use({ + accountId: cwd, + viewer: cwd, + credentialFingerprint: cwd, + }).pipe(Effect.provideService(SourceControlRateLimit.CredentialScope, cwd)), + listChangeRequests: ({ cwd }) => + Effect.suspend(() => { + reads.push(cwd); + if (cwd === "/a" && limited) + return Effect.fail( + new PullRequestProviderError({ + provider: "github", + operation: "list", + reason: "rate-limited", + detail: "Try later.", + }), + ); + return Effect.succeed({ + items: [changeRequest(1, "2026-07-02T00:00:00Z")], + truncated: false, + continues: false, + }); + }), + }), + ], + }); + yield* service.list({ state: "open" }); + reads.length = 0; + yield* service.invalidate({}); + const paused = yield* service.list({ state: "open" }); + assert.deepStrictEqual(reads, ["/b"]); + assert.strictEqual(paused.viewers["github.com"], "/b"); + assert.notProperty(paused.viewers, "a github.com"); + assert.strictEqual(paused.viewers["b github.com"], "/b"); + assert.deepStrictEqual( + paused.errors.map((error) => error.projectId), + ["a"], + ); + limited = false; + yield* TestClock.adjust("31 seconds"); + yield* service.invalidate({}); + const recovered = yield* service.list({ state: "open" }); + assert.deepStrictEqual(recovered.entries.map((entry) => entry.projectId).sort(), ["a", "b"]); + assert.strictEqual(recovered.viewers["a github.com"], "/a"); + }), +); + +it.effect("refreshes viewer and statistics ownership after workspace credentials change", () => + Effect.gen(function* () { + let account = "alice"; + const service = yield* makeService({ + projects: [project({ id: "a", title: "a", workspaceRoot: "/a", repository: "a/repo" })], + providers: [ + fakeProvider("github", { + withVerifiedCredential: (_, use) => + use({ accountId: account, viewer: account, credentialFingerprint: account }).pipe( + Effect.provideService(SourceControlRateLimit.CredentialScope, account), + ), + listChangeRequests: () => + Effect.gen(function* () { + const viewer = yield* SourceControlRateLimit.CredentialScope; + return { + items: [ + { + ...changeRequest(1, "2026-07-02T00:00:00Z"), + author: { login: viewer, name: null, avatarUrl: null }, + }, + ], + truncated: false, + continues: false, + }; + }), + listChangeRequestStats: ({ changeRequests }) => + Effect.gen(function* () { + const viewer = yield* SourceControlRateLimit.CredentialScope; + return changeRequests.map((ref) => ({ + ...ref, + additions: viewer === "alice" ? 1 : 2, + deletions: 0, + })); + }), + }), + ], + }); + const first = yield* service.list({ state: "open" }); + assert.strictEqual(first.viewers["a github.com"], "alice"); + assert.strictEqual((yield* service.listStats({ refs: first.entries })).stats[0]?.additions, 1); + account = "bob"; + yield* service.invalidate({}); + const refreshed = yield* service.list({ state: "open" }); + assert.strictEqual(refreshed.viewers["a github.com"], "bob"); + assert.strictEqual(refreshed.entries[0]?.author?.login, "bob"); + assert.strictEqual( + (yield* service.listStats({ refs: refreshed.entries })).stats[0]?.additions, + 2, + ); + }), +); + +it.effect("keeps a failed workspace credential isolated and recovers after refresh", () => + Effect.gen(function* () { + let unavailable = true; + const readRoots: string[] = []; + const service = yield* makeService({ + projects: ["a", "b"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `${name}/repo`, + }), + ), + providers: [ + fakeProvider("github", { + withVerifiedCredential: ({ cwd }, use) => + cwd === "/b" && unavailable + ? Effect.fail(requestFailed) + : use({ accountId: cwd, viewer: cwd, credentialFingerprint: cwd }), + getViewer: () => Effect.die("A failed credential must not borrow a host viewer"), + listChangeRequests: ({ cwd }) => + Effect.sync(() => { + readRoots.push(cwd); + return { + items: [changeRequest(1, "2026-07-02T00:00:00Z")], + truncated: false, + continues: false, + }; + }), + }), + ], + }); + const first = yield* service.list({ state: "open" }); + assert.deepStrictEqual(readRoots, ["/a"]); + assert.deepStrictEqual( + first.entries.map((entry) => entry.projectId), + ["a"], + ); + assert.deepStrictEqual( + first.errors.map((error) => error.projectId), + ["b"], + ); + const continuationError = yield* Effect.flip( + service.list({ + state: "open", + cursors: { "github.com b/repo": "2026-07-02T00:00:00Z|1|1" }, + }), + ); + assert.strictEqual(continuationError._tag, "PullRequestOperationError"); + unavailable = false; + yield* service.invalidate({}); + const recovered = yield* service.list({ state: "open" }); + assert.deepStrictEqual(recovered.entries.map((entry) => entry.projectId).sort(), ["a", "b"]); + assert.deepStrictEqual(recovered.errors, []); + }), +); + +it.effect("does not capture workspace credentials while the host is paused", () => + Effect.gen(function* () { + let captures = 0; + const service = yield* makeService({ + projects: [project({ id: "a", title: "a", workspaceRoot: "/a", repository: "a/repo" })], + providers: [ + fakeProvider("github", { + withVerifiedCredential: (_, use) => + Effect.suspend(() => { + captures++; + return use({ accountId: "a", viewer: "a", credentialFingerprint: "a" }); + }), + listChangeRequests: () => + Effect.fail( + new PullRequestProviderError({ + provider: "github", + operation: "list", + reason: "rate-limited", + detail: "Try later.", + }), + ), + }), + ], + }); + yield* service.list({ state: "open" }); + assert.strictEqual(captures, 1); + yield* service.invalidate({}); + const paused = yield* Effect.flip(service.list({ state: "open" })); + assert.strictEqual(paused._tag, "PullRequestOperationError"); + assert.strictEqual(captures, 1); + }), +); + +it.effect( + "lists and loads statistics through the GitHub provider with separate workspace accounts", + () => + Effect.gen(function* () { + const commands: VcsProcess.VcsProcessInput[] = []; + const rateLimits = yield* SourceControlRateLimit.make; + const response = (value: unknown) => ({ + exitCode: ChildProcessSpawner.ExitCode(0), + stdout: typeof value === "string" ? value : encodeProcessResponse(value), + stderr: "", + stdoutTruncated: false, + stderrTruncated: false, + stdoutInvalidUtf8: false, + }); + const github = yield* GitHubCli.make.pipe( + Effect.provide(GitHubGraphQlBudget.layer), + Effect.provideService(SourceControlRateLimit.SourceControlRateLimit, rateLimits), + Effect.provideService(VcsProcess.VcsProcess, { + run: (input) => + Effect.sync(() => { + commands.push(input); + const account = input.cwd === "/b" ? "bob" : "alice"; + if (input.args[0] === "auth") return response(`token-${account}`); + assert.strictEqual(input.env?.GH_TOKEN, `token-${account}`); + if (input.args[1] === "user") + return response({ id: account === "alice" ? 1 : 2, login: account }); + if (input.stdin !== undefined) { + const body = decodeSearchRequest(input.stdin); + const repositories = [...body.variables.q.matchAll(/repo:([^ ]+)/g)].map( + (match) => match[1]!, + ); + assert.isTrue(repositories.length > 0); + assert.isTrue( + repositories.every((repo) => (repo === "b/repo" ? "bob" : "alice") === account), + ); + assert.include(body.variables.q, `author:"${account}"`); + return response({ + data: { + search: { + pageInfo: { hasNextPage: false }, + nodes: repositories.map((repository) => ({ + number: 1, + title: repository, + url: `https://github.com/${repository}/pull/1`, + author: { login: account, avatarUrl: null }, + headRefName: "branch", + baseRefName: "main", + state: "OPEN", + isDraft: false, + mergeable: "MERGEABLE", + createdAt: "2026-07-01T00:00:00Z", + updatedAt: "2026-07-02T00:00:00Z", + repository: { nameWithOwner: repository }, + reviewRequests: { nodes: [] }, + labels: { nodes: [] }, + })), + }, + }, + }); + } + const query = input.args.find((arg) => arg.startsWith("query=")) ?? ""; + const data: Record = {}; + for (const match of query.matchAll(/(s\d+): repository\(owner: "([^"]+)"/g)) { + assert.strictEqual(match[2] === "b" ? "bob" : "alice", account); + data[match[1]!] = { + pullRequest: { additions: account === "bob" ? 2 : 1, deletions: 0 }, + }; + } + assert.isTrue(Object.keys(data).length > 0); + return response({ data }); + }), + }), + ); + const cli = yield* GitHubPullRequestCli.make.pipe( + Effect.provideService(GitHubCli.GitHubCli, github), + Effect.provide(GitHubGraphQlBudget.layer), + Effect.provideService(SourceControlRateLimit.SourceControlRateLimit, rateLimits), + ); + const provider = yield* GitHubPullRequestProvider.make.pipe( + Effect.provideService(GitHubPullRequestCli.GitHubPullRequestCli, cli), + ); + const service = yield* makeService({ + projects: ["a", "b", "c"].map((name) => + project({ + id: name, + title: name, + workspaceRoot: `/${name}`, + repository: `${name}/repo`, + }), + ), + providers: [provider], + rateLimits, + }); + const listed = yield* service.list({ state: "open", filters: { author: "@me" } }); + assert.deepStrictEqual(listed.entries.map((entry) => entry.projectId).sort(), [ + "a", + "b", + "c", + ]); + assert.deepStrictEqual(listed.errors, []); + const counts = yield* service.listStats({ refs: listed.entries }); + assert.deepStrictEqual(counts.stats.map((stat) => [stat.projectId, stat.additions]).sort(), [ + ["a", 1], + ["b", 2], + ["c", 1], + ]); + assert.strictEqual(commands.filter((command) => command.stdin !== undefined).length, 2); + assert.strictEqual(commands.filter((command) => command.args[1] === "graphql").length, 4); + const beforeContinuation = commands.length; + const next = yield* service.list({ + state: "open", + filters: { author: "@me" }, + cursors: { "github.com b/repo": "2026-07-02T00:00:00Z|1|1" }, + }); + assert.deepStrictEqual(next.entries, []); + assert.deepStrictEqual( + commands + .slice(beforeContinuation) + .filter((command) => command.stdin !== undefined) + .map((command) => command.cwd), + ["/b"], + ); + yield* cli.withVerifiedCredential({ cwd: "/b", host: "github.com" }, () => + rateLimits.recordRateLimit({ + provider: "github", + host: "github.com", + lease: 0, + retryAt: 20 * 60_000, + }), + ); + yield* TestClock.adjust("11 minutes"); + yield* service.invalidate({}); + const beforePaused = commands.length; + const paused = yield* service.list({ state: "open", filters: { author: "@me" } }); + assert.deepStrictEqual(paused.entries.map((entry) => entry.projectId).sort(), ["a", "c"]); + const pausedStats = yield* service.listStats({ refs: listed.entries }); + assert.deepStrictEqual(pausedStats.stats.map((entry) => entry.projectId).sort(), ["a", "c"]); + const rejected = yield* Effect.flip( + service.withRoutingCredential( + { + projectId: "b" as ProjectId, + repository: "b/repo", + number: 1, + host: "github.com", + expectedAccountId: "2", + }, + Effect.die("An unverified action must not run"), + ), + ); + assert.strictEqual(rejected._tag, "PullRequestOperationError"); + if (rejected._tag === "PullRequestOperationError") assert.include(rejected.detail, "paused"); + assert.isTrue( + commands + .slice(beforePaused) + .filter((command) => command.cwd === "/b") + .every((command) => command.args[0] === "auth"), + ); + yield* TestClock.adjust("10 minutes"); + yield* service.invalidate({}); + const recovered = yield* service.list({ state: "open", filters: { author: "@me" } }); + assert.deepStrictEqual(recovered.entries.map((entry) => entry.projectId).sort(), [ + "a", + "b", + "c", + ]); + }), +); + it.effect("caches narrow previews and invalidates them after refresh or mutation", () => Effect.gen(function* () { let reads = 0; @@ -397,9 +1015,11 @@ function fakeProvider( }; } +/** Build an isolated service with optional shared rate limits for integrated provider tests. */ function makeService(input: { readonly projects: ReadonlyArray; readonly providers: ReadonlyArray; + readonly rateLimits?: SourceControlRateLimit.SourceControlRateLimit["Service"]; readonly resolveHandle?: SourceControlProviderRegistry.SourceControlProviderRegistry["Service"]["resolveHandle"]; }) { // Built into the test's own scope rather than provided call by call: the marks store owns a @@ -421,7 +1041,9 @@ function makeService(input: { getProjectShellById: (projectId) => Effect.succeed(Option.fromNullishOr(input.projects.find((p) => p.id === projectId))), }), - SourceControlRateLimit.layer, + input.rateLimits === undefined + ? SourceControlRateLimit.layer + : Layer.succeed(SourceControlRateLimit.SourceControlRateLimit, input.rateLimits), // The real store over a database of its own, so the environment-kept marks are exercised // through the SQL that holds them rather than through a stand-in that agrees with itself. PullRequestFilesViewed.layer.pipe(Layer.provide(SqlitePersistenceMemory)), @@ -1703,6 +2325,8 @@ it.effect("keeps two hosts of one provider kind as two accounts", () => assert.deepStrictEqual(result.viewers, { "github.com": "bilal", "github.acme.dev": "b.hassan", + "p1 github.com": "bilal", + "p2 github.acme.dev": "b.hassan", }); assert.deepStrictEqual(result.entries.map((entry) => entry.host).toSorted(), [ "github.acme.dev", diff --git a/apps/server/src/pullRequest/PullRequestService.ts b/apps/server/src/pullRequest/PullRequestService.ts index 866a45e49b09..3272a151c3c4 100644 --- a/apps/server/src/pullRequest/PullRequestService.ts +++ b/apps/server/src/pullRequest/PullRequestService.ts @@ -72,7 +72,7 @@ import { } from "@t3tools/contracts"; import { detectSourceControlProviderFromRemoteUrl } from "@t3tools/shared/sourceControl"; -import { AllowGitHubReserve } from "../sourceControl/GitHubCli.ts"; +import { AllowGitHubReserve, PinnedGitHubCredential } from "../sourceControl/GitHubCli.ts"; import * as ProjectionSnapshotQuery from "../orchestration/Services/ProjectionSnapshotQuery.ts"; import * as PullRequestFilesViewed from "../persistence/PullRequestFilesViewed.ts"; import * as SourceControlProviderRegistry from "../sourceControl/SourceControlProviderRegistry.ts"; @@ -312,7 +312,7 @@ export interface SupportedProject { readonly project: OrchestrationProjectShell; readonly api: PullRequestProviderApi; readonly repository: string; - /** The host the repository lives on, which is the account boundary rather than the kind. */ + /** The host the repository lives on; individual workspaces may use different accounts. */ readonly host: string; /** * The identity's canonical key, which is what this environment's own records are keyed by. @@ -321,6 +321,73 @@ export interface SupportedProject { readonly remote: string; } +const readGroupKey = Schema.encodeSync(Schema.fromJsonString(Schema.Array(Schema.String))); + +/** A listing's workspace credential, retained only for the duration of its reads. */ +interface ReadProject extends SupportedProject { + readonly credential?: { + readonly credentialFingerprint: string; + readonly viewer: string | null; + readonly error: PullRequestProviderError | null; + }; + readonly runRead: ( + read: Effect.Effect, + ) => Effect.Effect; +} + +/** Pin each workspace before grouping reads; a failed account must not borrow another's access. */ +const prepareReadProject = ( + project: SupportedProject, + rateLimits: SourceControlRateLimit.SourceControlRateLimit["Service"], +): Effect.Effect => { + const verify = project.api.withVerifiedCredential; + if (verify === undefined) return Effect.succeed({ ...project, runRead: (read) => read }); + const checkPause = rateLimits.check({ provider: project.api.kind, host: project.host }).pipe( + Effect.mapError( + (error) => + new PullRequestProviderError({ + provider: project.api.kind, + operation: "getViewer", + reason: "rate-limited", + detail: error.detail, + retryAt: error.retryAt, + cause: error, + }), + ), + ); + return checkPause.pipe( + Effect.andThen( + verify({ cwd: project.project.workspaceRoot, host: project.host }, (identity) => + Effect.gen(function* (): Effect.fn.Return { + yield* checkPause; + const pinned = yield* PinnedGitHubCredential; + const scope = yield* SourceControlRateLimit.CredentialScope; + return { + ...project, + credential: { ...identity, error: null }, + runRead: (read) => + read.pipe( + Effect.provideService(PinnedGitHubCredential, pinned), + Effect.provideService(SourceControlRateLimit.CredentialScope, scope), + ), + }; + }), + ), + ), + Effect.catch((error) => + Effect.succeed({ + ...project, + credential: { + credentialFingerprint: `unavailable:${project.project.workspaceRoot}`, + viewer: null, + error, + }, + runRead: () => Effect.fail(error), + }), + ), + ); +}; + /** * What the workspace has, split by whether this build can read it. Hosts with no * implementation are counted rather than dropped, so their projects are explained in the @@ -617,6 +684,7 @@ const observeRead = Effect.fnUntraced(function* (read: Effect.Effect(64); const pullRequestRefreshes = yield* SubscriptionRef.make(0); @@ -1123,8 +1191,15 @@ export const make = Effect.gen(function* () { // authored/reviewing searches for that same repository are therefore real empty answers, not // a reason to issue the two-command per-repository fallback again. const searchVisibleAt = new Map(); - const searchVisibilityKey = (host: string, repository: string) => - `${host}\n${repository.trim().toLowerCase()}`; + const searchVisibilityKey = (host: string, repository: string, credential = "") => + readGroupKey([host, repository.trim().toLowerCase(), credential]); + + // Presentation only: a continuation must not reauthenticate unrelated hosts just to keep + // their switcher state. These values never authorize a read or supply a row's viewer. + const listingHostStates = new Map< + string, + { readonly at: number; readonly configured: boolean; readonly detail: string | null } + >(); const listUncached: PullRequestService["Service"]["list"] = (input) => Effect.gen(function* () { @@ -1134,34 +1209,99 @@ export const make = Effect.gen(function* () { // and reading part of the listing under that assumption would quietly lose rows. const continuation = yield* decodeCursors(input.cursors); const { - supported: projects, + supported: workspaceProjects, unimplemented, viewerRoots, } = yield* listWorkspaceProjects(input); + const projects = yield* Effect.forEach( + continuation === null + ? workspaceProjects + : workspaceProjects.filter(({ cursorKey }) => continuation.has(cursorKey)), + (project) => prepareReadProject(project, rateLimits), + { + concurrency: REPOSITORY_CONCURRENCY, + }, + ); const projectCounts = new Map(); - for (const { host } of projects) { + for (const { host } of workspaceProjects) { projectCounts.set(host, (projectCounts.get(host) ?? 0) + 1); } - const viewerResults = yield* resolveViewers(projects, viewerRoots); + const unpinnedViewers = yield* resolveViewers( + workspaceProjects.filter((project) => project.api.withVerifiedCredential === undefined), + viewerRoots, + ); + const byHost = new Map(unpinnedViewers.map((result) => [result.host, result])); + for (const project of projects) { + const credential = project.credential; + if (credential === undefined || byHost.get(project.host)?.viewer != null) continue; + byHost.set(project.host, { + host: project.host, + kind: project.api.kind, + viewer: credential.viewer, + error: credential.error, + }); + } + const viewerResults = [...byHost.values()]; const viewers: Record = {}; for (const result of viewerResults) { if (result.viewer !== null) viewers[result.host] = result.viewer; } + const viewerFor = (project: ReadProject) => + project.credential === undefined ? viewers[project.host] : project.credential.viewer; + for (const project of projects) { + const viewer = viewerFor(project); + if (viewer != null) + viewers[`${encodeURIComponent(project.project.id)} ${project.host}`] = viewer; + } + const observedAt = yield* Clock.currentTimeMillis; + for (const result of viewerResults) { + // Failure of a continued subset cannot declare the whole host signed out. + if ( + continuation !== null && + result.viewer === null && + projects.some( + (project) => project.host === result.host && project.credential !== undefined, + ) && + projects.filter((project) => project.host === result.host).length < + (projectCounts.get(result.host) ?? 0) + ) + continue; + listingHostStates.delete(result.host); + listingHostStates.set(result.host, { + at: observedAt, + configured: result.viewer !== null, + detail: result.error === null ? null : providerDetail(result.error), + }); + if (listingHostStates.size > VIEWER_CACHE_CAPACITY) { + listingHostStates.delete(listingHostStates.keys().next().value!); + } + } // One summary per host, which is what the viewer lookup already answers for: two GitHub // hosts sign in separately, so collapsing them by kind would report one as the other. const providers: ReadonlyArray = [ - ...viewerResults.map((result) => ({ - host: result.host, - kind: result.kind, - searchesOnHost: - projects.find((project) => project.host === result.host)?.api.capabilities.search ?? - false, - projectCount: projectCounts.get(result.host) ?? 1, - configured: result.viewer !== null, - detail: result.error === null ? null : providerDetail(result.error), - })), + ...[...projectCounts].map(([host, projectCount]) => { + const project = workspaceProjects.find((project) => project.host === host)!; + const held = listingHostStates.get(host); + const state = + held !== undefined && observedAt - held.at <= Duration.toMillis(VIEWER_CACHE_TTL) + ? held + : undefined; + return { + host, + kind: project.api.kind, + searchesOnHost: project.api.capabilities.search, + projectCount, + configured: state?.configured ?? false, + detail: + state !== undefined + ? state.detail + : continuation !== null + ? "Sign-in was not checked on this continuation page." + : null, + }; + }), ...[...unimplemented].map(([host, { kind, projectCount }]) => ({ host, kind, @@ -1176,15 +1316,12 @@ export const make = Effect.gen(function* () { // other one is already on the page, and reading it again is the whole cost this is here to // avoid. The host summaries above stay over the whole workspace, because the switcher they // fill is about the workspace rather than about this slice. - const selected = - continuation === null - ? projects - : projects.filter(({ cursorKey }) => continuation.has(cursorKey)); - const readable = selected.filter(({ host }) => viewers[host] !== undefined); + const selected = projects; + const readable = selected.filter((project) => viewerFor(project) != null); // A host that could not be read still has projects, and they are absent from the list. // Reporting them keeps "N repositories were unavailable" honest instead of dropping them. const unreadable = selected - .filter(({ host }) => viewers[host] === undefined) + .filter((project) => viewerFor(project) == null) .map(({ project, repository }) => ({ projectId: project.id, projectTitle: project.title, @@ -1199,11 +1336,14 @@ export const make = Effect.gen(function* () { // Only the hosts this request was actually going to read: a continuation that named // nothing has asked for nothing, and a host it never mentioned being signed out is no // reason to refuse it. - const errors = viewerResults.flatMap((result) => - result.error === null || !selected.some(({ host }) => host === result.host) - ? [] - : [result.error], - ); + const errors = [ + ...selected.flatMap(({ credential }) => (credential?.error ? [credential.error] : [])), + ...viewerResults.flatMap((result) => + result.error === null || !selected.some(({ host }) => host === result.host) + ? [] + : [result.error], + ), + ]; const blocking = errors.find(isProviderUnusable) ?? errors[0]; if (blocking) { return yield* toPullRequestError("list")(blocking); @@ -1227,9 +1367,9 @@ export const make = Effect.gen(function* () { * One repository asked on its own. What every host without a search across repositories * does, and what a batched read falls back to for a repository it could not answer for. */ - const readRepository = (project: SupportedProject): Effect.Effect => { + const readRepository = (project: ReadProject): Effect.Effect => { { - const viewer = viewers[project.host]!; + const viewer = viewerFor(project)!; const key = project.cursorKey; const cursor = cursorOf(project); return project.api @@ -1254,6 +1394,7 @@ export const make = Effect.gen(function* () { }), }) .pipe( + project.runRead, observeRead, Effect.map(({ value: page, observedAt }): RepositoryBatch => { // The boundary instant was asked for inclusively, so the rows already sent at it @@ -1300,7 +1441,7 @@ export const make = Effect.gen(function* () { }; /** - * One host's repositories in one read. The slice is the newest `limit` rows across all of + * One credential's repositories in one read. The slice is the newest `limit` rows across all of * them, so it is split back up by repository here: the page still reports per project, and * each repository still carries on from a cursor of its own. * @@ -1309,14 +1450,14 @@ export const make = Effect.gen(function* () { * repositories as unreadable before anyone has asked it about them one at a time. */ const readTogether = ( - chunk: ReadonlyArray, + chunk: ReadonlyArray, ): Effect.Effect> => { const first = chunk[0]!; const readAcross = first.api.listChangeRequestsAcross; const separately = () => Effect.forEach(chunk, readRepository, { concurrency: REPOSITORY_CONCURRENCY }); if (readAcross === undefined) return separately(); - const viewer = viewers[first.host]!; + const viewer = viewerFor(first)!; const cursor = cursorOf(first); return readAcross({ cwd: first.project.workspaceRoot, @@ -1332,6 +1473,7 @@ export const make = Effect.gen(function* () { ? {} : { cursor: { updatedBefore: cursor.updatedBefore, delivered: cursor.delivered } }), }).pipe( + first.runRead, observeRead, Effect.flatMap(({ value: page, observedAt }) => Effect.flatMap(Clock.currentTimeMillis, (now) => { @@ -1346,7 +1488,14 @@ export const make = Effect.gen(function* () { const held = rows.get(key); if (held === undefined) rows.set(key, [item]); else held.push(item); - searchVisibleAt.set(searchVisibilityKey(first.host, item.repository), now); + searchVisibleAt.set( + searchVisibilityKey( + first.host, + item.repository, + first.credential?.credentialFingerprint, + ), + now, + ); } // The oldest row of the whole slice, which is how far every repository in it has now // been read — including the ones that contributed nothing to it. @@ -1368,7 +1517,11 @@ export const make = Effect.gen(function* () { // price of one request per repository with nothing in the first slice — which // run together, and only there. const lastVisible = searchVisibleAt.get( - searchVisibilityKey(project.host, project.repository), + searchVisibilityKey( + project.host, + project.repository, + project.credential?.credentialFingerprint, + ), ); const searchIsKnownVisible = !page.truncated && @@ -1414,14 +1567,18 @@ export const make = Effect.gen(function* () { // A host with a search across repositories is asked once for all of them; everyone else is // asked once each. Repositories standing at different points of the same listing are // different questions, so they are grouped by the boundary they carry on from. - const together = new Map>(); - const separate: Array = []; + const together = new Map>(); + const separate: Array = []; for (const project of readable) { if (project.api.listChangeRequestsAcross === undefined) { separate.push(project); continue; } - const key = `${project.host}\n${cursorOf(project)?.updatedBefore ?? ""}`; + const key = readGroupKey([ + project.host, + project.credential?.credentialFingerprint ?? "", + cursorOf(project)?.updatedBefore ?? "", + ]); const group = together.get(key); if (group === undefined) together.set(key, [project]); else group.push(project); @@ -1454,9 +1611,8 @@ export const make = Effect.gen(function* () { }); /** - * Who this project's host says the reader is. Shared with the listing's own lookup — the same - * ten-minute answer per host — so a page that has already listed anything pays nothing for it, - * and a host that cannot say leaves it null rather than failing the read it decorates. + * Who this project's host says the reader is. A routed read uses its verified viewer; + * other reads share the host lookup and leave it null when the host cannot answer. */ const viewerOf = (project: SupportedProject): Effect.Effect => routingCredential.pipe( @@ -1488,6 +1644,7 @@ export const make = Effect.gen(function* () { return { ...identity, host, provider: "github" as const }; }); + /** Verify the expected account before an operation, retaining actionable rate-limit errors. */ const withRoutingCredential: PullRequestService["Service"]["withRoutingCredential"] = ( input, operation, @@ -1515,7 +1672,15 @@ export const make = Effect.gen(function* () { ? operation.pipe(Effect.provideService(routingCredential, identity), Effect.result) : Effect.fail(rejected()), ) - .pipe(Effect.catchTag("PullRequestProviderError", () => Effect.fail(rejected()))); + .pipe( + Effect.catchTag("PullRequestProviderError", (error) => + Effect.fail( + error.reason === "rate-limited" + ? toPullRequestError("routeIdentity")(error) + : rejected(), + ), + ), + ); return yield* Effect.fromResult(result); }); @@ -2434,11 +2599,14 @@ export const make = Effect.gen(function* () { Effect.gen(function* () { if (input.refs.length === 0) return { stats: [] }; const { supported } = yield* listWorkspaceProjects({}); - const byProject = new Map(supported.map((project) => [project.project.id, project])); - const wanted = new Map< - string, - { readonly project: SupportedProject; readonly number: number } - >(); + const requestedIds = new Set(input.refs.map((ref) => ref.projectId)); + const projects = yield* Effect.forEach( + supported.filter((project) => requestedIds.has(project.project.id)), + (project) => prepareReadProject(project, rateLimits), + { concurrency: REPOSITORY_CONCURRENCY }, + ); + const byProject = new Map(projects.map((project) => [project.project.id, project])); + const wanted = new Map(); for (const ref of input.refs) { const project = byProject.get(ref.projectId); // The repository travels through the client, so it is checked against the project's own @@ -2452,10 +2620,14 @@ export const make = Effect.gen(function* () { } wanted.set(`${project.project.id} ${ref.number}`, { project, number: ref.number }); } - const byHost = new Map>(); + const byHost = new Map>(); for (const entry of wanted.values()) { - const held = byHost.get(entry.project.host); - if (held === undefined) byHost.set(entry.project.host, [entry]); + const key = readGroupKey([ + entry.project.host, + entry.project.credential?.credentialFingerprint ?? "", + ]); + const held = byHost.get(key); + if (held === undefined) byHost.set(key, [entry]); else held.push(entry); } const stats = yield* Effect.forEach( @@ -2479,6 +2651,7 @@ export const make = Effect.gen(function* () { number: entry.number, })), }).pipe( + first.project.runRead, Effect.map((read) => read.flatMap((stat): ReadonlyArray => { const project = projectsByRepository.get( @@ -3115,6 +3288,7 @@ export const make = Effect.gen(function* () { listingsEpoch = ++epochCounter; everyFileRevisionEpoch = ++epochCounter; viewersByHost.clear(); + listingHostStates.clear(); yield* Cache.invalidateAll(viewerFlights); } if (options?.notifyReaders) { diff --git a/apps/web/src/components/pullRequest/pullRequestList.logic.test.ts b/apps/web/src/components/pullRequest/pullRequestList.logic.test.ts index 357ea20ad76d..3c4359286ff9 100644 --- a/apps/web/src/components/pullRequest/pullRequestList.logic.test.ts +++ b/apps/web/src/components/pullRequest/pullRequestList.logic.test.ts @@ -6,6 +6,7 @@ import { findScopedProject, mergePullRequestLists, pullRequestEntryKey, + pullRequestEntryViewer, pullRequestEnvironmentSetKey, groupPullRequestsByInvolvement, matchesPullRequestFilters, @@ -1395,6 +1396,97 @@ describe('who "I" am, per server', () => { ).toEqual([1]); }); + it("keeps same-host project viewers distinct within and across environments", () => { + const rows = [ + entry({ number: 1, projectId: "a" as ProjectId, author: byBilal }), + entry({ number: 2, projectId: "b" as ProjectId, author: { ...byBilal, login: "Octocat" } }), + entry({ number: 3, projectId: "b" as ProjectId, author: byBilal }), + ]; + const viewers = { "github.com": "Bilal", "a github.com": "Bilal", "b github.com": "Octocat" }; + const merged = mergePullRequestLists([ + [ENV_1, answer(viewers, rows)], + [ENV_2, answer({ "github.com": "Other", "a github.com": "Other" }, [rows[0]!])], + ])!; + const authored = filterPullRequestsByInvolvement(merged.entries, merged.viewers, "authored"); + expect(authored.map((row) => row.number)).toEqual([1, 2]); + expect( + groupPullRequestsByInvolvement(merged.entries, merged.viewers).find( + (group) => group.key === "authored", + )?.entries, + ).toEqual(authored); + expect( + merged.entries + .filter((row) => + matchesPullRequestFilters( + row, + { author: "me" }, + pullRequestEntryViewer(row, merged.viewers), + ), + ) + .map((row) => row.number), + ).toEqual([1, 2]); + // A single-server snapshot has unprefixed project keys. + expect( + filterPullRequestsByInvolvement(rows, viewers, "authored").map((row) => row.number), + ).toEqual([1, 2]); + }); + + it("does not borrow a host viewer for a project omitted by a new server", () => { + const missing = entry({ number: 1, projectId: "missing" as ProjectId, author: byBilal }); + const viewers = { "github.com": "Bilal", "known github.com": "Bilal" }; + expect(pullRequestEntryViewer(missing, viewers)).toBeNull(); + const merged = mergePullRequestLists([ + [ENV_1, answer(viewers, [missing])], + [ENV_2, answer({ "github.com": "Bilal" }, [missing])], + ])!; + expect(merged.entries.map((row) => pullRequestEntryViewer(row, merged.viewers))).toEqual([ + null, + "bilal", + ]); + expect( + filterPullRequestsByInvolvement(merged.entries, merged.viewers, "authored").map( + (row) => row.environmentId, + ), + ).toEqual([ENV_2]); + }); + + it("does not confuse an older environment key with an unscoped project key", () => { + const row = entry({ number: 1, projectId: "env-2" as ProjectId, author: byBilal }); + const merged = mergePullRequestLists([ + [ENV_1, answer({ "github.com": "Octocat" }, [row])], + [ENV_2, answer({ "github.com": "Bilal" }, [])], + ])!; + expect(pullRequestEntryViewer(merged.entries[0]!, merged.viewers)).toBe("octocat"); + }); + + it("keeps spaces and percent signs in identity ids from colliding across environments", () => { + const merged = mergePullRequestLists([ + [ + "a" as EnvironmentId, + answer({ "github.com": "Bilal", "b%20c github.com": "Bilal" }, [ + entry({ number: 1, projectId: "b c" as ProjectId, author: byBilal }), + ]), + ], + [ + "a b" as EnvironmentId, + answer({ "github.com": "Octocat", "c github.com": "Octocat" }, [ + entry({ number: 2, projectId: "c" as ProjectId, author: byBilal }), + ]), + ], + [ + "a%20b" as EnvironmentId, + answer({ "github.com": "Bilal", "c github.com": "Bilal" }, [ + entry({ number: 3, projectId: "c" as ProjectId, author: byBilal }), + ]), + ], + ])!; + expect( + filterPullRequestsByInvolvement(merged.entries, merged.viewers, "authored").map( + (row) => row.number, + ), + ).toEqual([1, 3]); + }); + it("still reads a single server's host-keyed viewers, which is what a snapshot carries", () => { expect( filterPullRequestsByInvolvement( diff --git a/apps/web/src/components/pullRequest/pullRequestList.logic.ts b/apps/web/src/components/pullRequest/pullRequestList.logic.ts index f4cb3a31f46f..ca795a68ce8c 100644 --- a/apps/web/src/components/pullRequest/pullRequestList.logic.ts +++ b/apps/web/src/components/pullRequest/pullRequestList.logic.ts @@ -69,8 +69,9 @@ export type PullRequestViewers = PullRequestListResult["viewers"]; /** A row plus the environment that read it, where the caller has one to give. */ type ScopedEntry = PullRequestListEntry & { readonly environmentId?: string }; +/** Scope a row's account to both its environment and workspace project. */ const pullRequestViewerKey = (entry: ScopedEntry): string => - `${entry.environmentId ?? ""} ${entry.host}`; + `${encodeURIComponent(entry.environmentId ?? "")} ${encodeURIComponent(entry.projectId)} ${entry.host}`; const GROUP_LABELS: Record = { reviewRequested: "Review requested", @@ -144,9 +145,27 @@ export function pullRequestEntryViewer( entry: ScopedEntry, viewers: PullRequestViewers, ): string | null { - // The environment's own answer first; a plain host key is what a single-environment listing - // still writes, and what the snapshot from one carries. - return normalize(viewers[pullRequestViewerKey(entry)] ?? viewers[entry.host]); + const unscoped = Object.hasOwn(viewers, entry.host); + const projectViewer = + viewers[pullRequestViewerKey(entry)] ?? + (unscoped ? viewers[`${encodeURIComponent(entry.projectId)} ${entry.host}`] : undefined); + if (projectViewer !== undefined) return normalize(projectViewer); + // A new server deliberately omits unreadable projects. Do not borrow a healthy project's + // host viewer. Only answers without project identities use the legacy host fallback. + const prefix = !unscoped ? `${encodeURIComponent(entry.environmentId ?? "")} ` : ""; + const hasProjectViewers = Object.keys(viewers).some( + (key) => + key.startsWith(prefix) && + key.endsWith(` ${entry.host}`) && + key.slice(prefix.length).split(" ").length === 2, + ); + return hasProjectViewers + ? null + : normalize( + viewers[`${encodeURIComponent(entry.environmentId ?? "")} ${entry.host}`] ?? + viewers[`${entry.environmentId ?? ""} ${entry.host}`] ?? + viewers[entry.host], + ); } /** @@ -713,7 +732,7 @@ export function mergePullRequestLists( let truncated = false; for (const [environmentId, answer] of answers) { for (const [host, login] of Object.entries(answer.viewers)) { - viewers[`${environmentId} ${host}`] = login; + viewers[`${encodeURIComponent(environmentId)} ${host}`] = login; } for (const provider of answer.providers) { const held = providers.get(provider.host); diff --git a/packages/contracts/src/pullRequest.ts b/packages/contracts/src/pullRequest.ts index 7a249cefb373..7c4f1ddd59f2 100644 --- a/packages/contracts/src/pullRequest.ts +++ b/packages/contracts/src/pullRequest.ts @@ -624,10 +624,9 @@ export type PullRequestListProjectError = typeof PullRequestListProjectError.Typ export const PullRequestListResult = Schema.Struct({ /** - * The signed-in account per host, which is what involvement filtering compares. Keyed by - * host rather than by provider kind: two GitHub hosts are two accounts. A host that could - * not be read is absent rather than present-and-undefined, because an open-keyed record - * cannot carry an optional value through the JSON codec. + * Verified accounts keyed by `${encodeURIComponent(projectId)} ${host}`, with host-only entries for older clients. + * Projects on one host may use different accounts. Unreadable identities are omitted. + * Clients merging environments prefix each key with the encoded environment id and a space. */ viewers: Schema.Record(TrimmedNonEmptyString, TrimmedNonEmptyString), providers: Schema.Array(PullRequestProviderSummary),