From 500391c9d41df47fc7228c460f7bcca93a366d14 Mon Sep 17 00:00:00 2001 From: eimexdev Date: Sat, 3 Oct 2026 18:52:44 -0700 Subject: [PATCH 1/3] fix(server): read paginated review replies when watching PRs --- .../PullRequestWatchReactor.ts | 58 +++++++- .../src/orchestration-v2/runtimeLayer.test.ts | 138 +++++++++++++++--- 2 files changed, 166 insertions(+), 30 deletions(-) diff --git a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts index 4eb6ec80ce17..f3c066989b0b 100644 --- a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts +++ b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts @@ -2,6 +2,8 @@ import { CommandId, MessageId, type OrchestrationV2Notification, + type PullRequestActivity, + type PullRequestRef, type ThreadPullRequestLink, type ThreadPullRequestWatch, } from "@t3tools/contracts"; @@ -125,6 +127,52 @@ export const make = Effect.gen(function* () { }, }).pipe(Effect.catch(() => record(target, null))); + const readRemarks = Effect.fn("PullRequestWatchReactor.readRemarks")( + function* (reference: PullRequestRef, activity: PullRequestActivity) { + // Truncation without a continuation means the initial thread read was incomplete. + if ( + activity.commentsTruncated && + !activity.reviewThreads.some((thread) => thread.nextCommentsCursor !== undefined) + ) + return null; + + const comments = new Map(activity.comments.map((comment) => [comment.id, comment])); + for (const thread of activity.reviewThreads) { + const ids = new Set(thread.comments.map((comment) => comment.id)); + const cursors = new Set(); + let cursor = thread.nextCommentsCursor ?? null; + while (cursor !== null) { + if (cursors.has(cursor)) return null; + cursors.add(cursor); + const page = yield* pullRequests.threadComments({ + ...reference, + threadId: thread.id, + cursor, + }); + for (const comment of page.comments) { + ids.add(comment.id); + comments.set(comment.id, { + ...comment, + kind: "review-comment", + path: thread.path, + reviewState: null, + }); + } + cursor = page.nextCursor; + } + if (ids.size < (thread.commentCount ?? 0)) return null; + } + return [...comments.values()].sort((left, right) => + left.createdAt.localeCompare(right.createdAt), + ); + }, + Effect.catch((error) => + Effect.logWarning("pull request watch comment pagination failed", { error }).pipe( + Effect.as(null), + ), + ), + ); + const check = Effect.fn("PullRequestWatchReactor.check")(function* (target: WatchTarget) { const { thread, link, watch } = target; const pullRequest = identityOf(link); @@ -160,13 +208,9 @@ export const make = Effect.gen(function* () { const [detail, activity] = read.value; if (detail.state !== "open") return yield* record(target, null); - // A degraded read (GitHub's review thread query failed) is truncated with no long thread to - // explain it, and would skip review comments, so remarks wait for a later pass. Replies past - // the first ten of a long review thread are not read. - const degraded = - activity.commentsTruncated && - !activity.reviewThreads.some((reviewThread) => reviewThread.nextCommentsCursor !== undefined); - const report = evaluatePullRequestWatch(watch, detail, degraded ? null : activity.comments); + // Never advance the remark watermark past comments an incomplete read could have missed. + const remarks = yield* readRemarks(reference, activity); + const report = evaluatePullRequestWatch(watch, detail, remarks); if (report.changes.length > 0) { return yield* record( target, diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 2c5235c767d7..c119503efdc0 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -18,6 +18,8 @@ import { type OrchestrationV2Run, ProjectId, type PullRequestDetail, + type PullRequestComment, + PullRequestOperationError, ProviderDriverKind, ProviderInstanceId, ProviderThreadId, @@ -2307,11 +2309,17 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }), ); - it.effect("wakes a watched thread once for failed checks and a review comment", () => + it.effect.each([ + "single page", + "paginated", + "page failure", + "repeated cursor", + "missing comments", + ])("wakes a watched thread once: %s", (mode) => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; - const threadId = ThreadId.make("runtime-pull-request-watch-wake"); - const projectId = ProjectId.make("pr-watch-wake-project"); + const threadId = ThreadId.make(`runtime-pull-request-watch-wake-${mode}`); + const projectId = ProjectId.make(`pr-watch-wake-project-${mode}`); yield* seedProject({ projectId, title: "Watch wake", @@ -2323,7 +2331,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { type: "thread.create", createdBy: "user", creationSource: "web", - commandId: CommandId.make("pr-watch-wake-create"), + commandId: CommandId.make(`pr-watch-wake-create-${mode}`), threadId, projectId, title: "Watch wake", @@ -2337,7 +2345,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { const url = "https://github.com/pingdotgg/t3code/pull/7"; yield* orchestrator.dispatch({ type: "thread.pull-request.link", - commandId: CommandId.make("pr-watch-wake-link"), + commandId: CommandId.make(`pr-watch-wake-link-${mode}`), threadId, ...key, url, @@ -2345,7 +2353,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }); yield* orchestrator.dispatch({ type: "thread.pull-request.watch", - commandId: CommandId.make("pr-watch-wake-start"), + commandId: CommandId.make(`pr-watch-wake-start-${mode}`), threadId, ...key, watching: true, @@ -2398,6 +2406,37 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { mergeCapabilities: { merge: true, squash: true, rebase: true }, viewer: "agent-user", }; + const initialWatch = (yield* orchestrator.getThreadShell(threadId))?.pullRequests?.[0]?.watch; + assert.isDefined(initialWatch); + const remark: PullRequestComment = { + id: "review-1", + kind: "review-comment", + author: { login: "reviewer", name: null, avatarUrl: null }, + body: "One more thing.", + createdAt: "2999-01-01T00:00:03.000Z", + url: null, + path: "src/index.ts", + reviewState: null, + }; + const firstTen = Array.from({ length: 10 }, (_, index) => ({ + ...remark, + id: `old-${index}`, + createdAt: "1900-01-01T00:00:00.000Z", + })); + const eleventh = { + ...remark, + id: "reply-11", + body: "Eleventh reply.", + createdAt: "2999-01-01T00:00:01.000Z", + }; + const twelfth = { + ...remark, + id: "reply-12", + body: "Twelfth reply.", + createdAt: "2999-01-01T00:00:02.000Z", + }; + const incomplete = mode !== "single page" && mode !== "paginated"; + let recovering = false; const reactor = yield* PullRequestWatchReactor.make.pipe( Effect.provide( Layer.mergeAll( @@ -2406,27 +2445,70 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { detail: () => Effect.succeed(detail), activity: () => Effect.succeed({ - comments: [ - { - id: "review-1", - kind: "review-comment", - author: { login: "reviewer", name: null, avatarUrl: null }, - body: "One more thing.", - createdAt: "2999-01-01T00:00:00.000Z", - url: null, - path: "src/index.ts", - reviewState: null, - }, - ], - commentCount: 1, - commentsTruncated: false, - reviewThreads: [], + comments: + mode === "single page" + ? [remark] + : [...firstTen, { ...remark, kind: "issue-comment" }], + commentCount: mode === "single page" ? 1 : 13, + commentsTruncated: mode !== "single page", + reviewThreads: + mode === "single page" + ? [] + : [ + { + id: "review-thread", + path: "src/index.ts", + line: 1, + side: "right", + isResolved: false, + isOutdated: false, + comments: firstTen, + commentCount: 12, + nextCommentsCursor: "after-10", + }, + ], commits: [], }), + threadComments: (input) => { + assert.equal(input.threadId, "review-thread"); + if (input.cursor === "after-10") { + return Effect.succeed({ comments: [eleventh], nextCursor: "after-11" }); + } + assert.equal(input.cursor, "after-11"); + if (!recovering && mode === "page failure") { + return Effect.fail( + new PullRequestOperationError({ + operation: "threadComments", + detail: "Page unavailable", + }), + ); + } + if (!recovering && mode === "repeated cursor") { + return Effect.succeed({ comments: [eleventh], nextCursor: "after-11" }); + } + if (!recovering && mode === "missing comments") { + return Effect.succeed({ comments: [], nextCursor: null }); + } + // Overlapping pages must not report the same reply twice. + return Effect.succeed({ comments: [eleventh, twelfth], nextCursor: null }); + }, }), ), ), ); + if (incomplete) { + yield* reactor.sweep; + const held = (yield* orchestrator.getThreadShell(threadId))?.pullRequests?.[0]?.watch; + assert.equal(held?.remarksThrough, initialWatch?.remarksThrough); + assert.deepEqual(held?.remarkIds, initialWatch?.remarkIds); + assert.deepEqual(held?.failedChecks, ["lint"]); + const { messages } = yield* orchestrator.getThreadRecords(threadId, ["messages"]); + assert.deepEqual( + messages.map((message) => message.notification?.summary), + ["#7: checks failed"], + ); + recovering = true; + } yield* reactor.sweep; yield* reactor.sweep; @@ -2435,12 +2517,22 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { messages.flatMap((message) => message.notification === undefined ? [] : [message.notification.summary], ), - ["#7: checks failed, new comments"], + incomplete + ? ["#7: checks failed", "#7: new comments"] + : ["#7: checks failed, new comments"], ); + if (mode !== "single page") { + const wake = messages.at(-1); + assert.include(wake?.text ?? "", "3 new comments"); + assert.include(wake?.text ?? "", "Eleventh reply."); + assert.include(wake?.text ?? "", "Twelfth reply."); + } const watch = (yield* orchestrator.getThreadShell(threadId))?.pullRequests?.[0]?.watch; + assert.equal(watch?.remarksThrough, remark.createdAt); + assert.deepEqual(watch?.remarkIds, [remark.id]); assert.deepEqual( { headSha: watch?.headSha, failedChecks: watch?.failedChecks, wakes: watch?.wakes }, - { headSha: "abc1234def", failedChecks: ["lint"], wakes: 0 }, + { headSha: "abc1234def", failedChecks: ["lint"], wakes: incomplete ? 1 : 0 }, ); }), ); From ed9a7b26a04f8286e0ea98cf546573a58cc35a5c Mon Sep 17 00:00:00 2001 From: eimexdev Date: Sat, 3 Oct 2026 19:22:22 -0700 Subject: [PATCH 2/3] fix(server): hold watch remarks when review threads are incomplete --- .../src/orchestration-v2/PullRequestWatchReactor.ts | 9 +++------ .../src/orchestration-v2/runtimeLayer.test.ts | 4 ++++ .../src/pullRequest/GitHubPullRequestCli.test.ts | 11 ++++++++++- apps/server/src/pullRequest/GitHubPullRequestCli.ts | 1 + .../pullRequest/GitHubPullRequestProvider.test.ts | 2 ++ .../src/pullRequest/GitHubPullRequestProvider.ts | 2 ++ apps/server/src/pullRequest/PullRequestProvider.ts | 1 + .../src/pullRequest/PullRequestService.test.ts | 13 +++++++++---- apps/server/src/pullRequest/PullRequestService.ts | 3 +++ .../server/src/pullRequest/gitHubPullRequestJson.ts | 1 + packages/contracts/src/pullRequest.ts | 2 ++ 11 files changed, 38 insertions(+), 11 deletions(-) diff --git a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts index f3c066989b0b..9e5e9b0631b0 100644 --- a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts +++ b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts @@ -129,12 +129,9 @@ export const make = Effect.gen(function* () { const readRemarks = Effect.fn("PullRequestWatchReactor.readRemarks")( function* (reference: PullRequestRef, activity: PullRequestActivity) { - // Truncation without a continuation means the initial thread read was incomplete. - if ( - activity.commentsTruncated && - !activity.reviewThreads.some((thread) => thread.nextCommentsCursor !== undefined) - ) - return null; + // Comment cursors cannot account for missing threads. Only finish a truncated read + // when the host confirms that every thread was listed. + if (activity.commentsTruncated && activity.reviewThreadsTruncated !== false) return null; const comments = new Map(activity.comments.map((comment) => [comment.id, comment])); for (const thread of activity.reviewThreads) { diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index c119503efdc0..0e4eb5452b82 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -2315,6 +2315,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { "page failure", "repeated cursor", "missing comments", + "thread list truncated", ])("wakes a watched thread once: %s", (mode) => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; @@ -2451,6 +2452,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { : [...firstTen, { ...remark, kind: "issue-comment" }], commentCount: mode === "single page" ? 1 : 13, commentsTruncated: mode !== "single page", + reviewThreadsTruncated: mode === "thread list truncated" && !recovering, reviewThreads: mode === "single page" ? [] @@ -2507,6 +2509,8 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { messages.map((message) => message.notification?.summary), ["#7: checks failed"], ); + // Keep the two notifications ordered independently of their random message IDs. + yield* TestClock.adjust("1 millis"); recovering = true; } yield* reactor.sweep; diff --git a/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts b/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts index 291105e25cd9..7db96681297f 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestCli.test.ts @@ -3612,7 +3612,14 @@ layer("GitHubPullRequestCli.layer", (it) => { Effect.gen(function* () { // A host that never runs out of pages: the walk has to end itself. mockedExecute.mockReturnValue( - Effect.succeed(output(reviewThreadsPage([thread("PRRT_1", "c1")], "Y3Vyc29yOjE"))), + Effect.succeed( + output( + reviewThreadsPage( + [{ ...thread("PRRT_1", "c1"), comments: threadComments(["c1"], "Y3Vyc29yOjI", 3) }], + "Y3Vyc29yOjE", + ), + ), + ), ); const cli = yield* GitHubPullRequestCli.GitHubPullRequestCli; @@ -3625,6 +3632,7 @@ layer("GitHubPullRequestCli.layer", (it) => { assert.strictEqual(mockedExecute.mock.calls.length, 10); assert.isTrue(conversation.truncated); + assert.isTrue(conversation.reviewThreadsTruncated); }), ); @@ -3656,6 +3664,7 @@ layer("GitHubPullRequestCli.layer", (it) => { nextCommentsCursor: "Y3Vyc29yOjI", }); assert.isTrue(conversation.truncated); + assert.isFalse(conversation.reviewThreadsTruncated); }), ); diff --git a/apps/server/src/pullRequest/GitHubPullRequestCli.ts b/apps/server/src/pullRequest/GitHubPullRequestCli.ts index f40cff8f50bf..f4f95e7afbf9 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestCli.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestCli.ts @@ -2409,6 +2409,7 @@ export const make = Effect.gen(function* () { // where a bound kept some of the words on GitHub. commentCount: entries.reduce((total, entry) => total + entry.commentCount, 0), truncated: cursor !== null || entries.some((entry) => entry.nextCommentCursor !== null), + reviewThreadsTruncated: cursor !== null, reactions, reactionsById, reviewers, diff --git a/apps/server/src/pullRequest/GitHubPullRequestProvider.test.ts b/apps/server/src/pullRequest/GitHubPullRequestProvider.test.ts index 4b2d618a0768..2b65ed729565 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestProvider.test.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestProvider.test.ts @@ -882,6 +882,7 @@ describe("getChangeRequest commits", () => { reviewThreads: [], commentCount: 0, truncated: false, + reviewThreadsTruncated: false, reactions: [], reactionsById: new Map>(), reviewers: [], @@ -967,6 +968,7 @@ describe("getChangeRequestActivity dismissed reviews", () => { reviewThreads: [], commentCount: 0, truncated: false, + reviewThreadsTruncated: false, reactions: [], reactionsById: new Map(), reviewers: [], diff --git a/apps/server/src/pullRequest/GitHubPullRequestProvider.ts b/apps/server/src/pullRequest/GitHubPullRequestProvider.ts index 5e70c2ee1ed4..c45caa0b7bf6 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestProvider.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestProvider.ts @@ -411,6 +411,7 @@ export const make = Effect.gen(function* () { reviewThreads: [], commentCount: 0, truncated: true, + reviewThreadsTruncated: true, reviewers: [], avatarsByLogin: new Map(), botLogins: new Set(), @@ -479,6 +480,7 @@ export const make = Effect.gen(function* () { // are always whole and only the thread walk can stop short of the host. commentCount: pullRequest.comments.length + reviewThreads.commentCount, commentsTruncated: reviewThreads.truncated, + reviewThreadsTruncated: reviewThreads.reviewThreadsTruncated, reviewThreads: reviewThreads.reviewThreads.map((thread) => ({ ...thread, comments: thread.comments.map((comment) => ({ diff --git a/apps/server/src/pullRequest/PullRequestProvider.ts b/apps/server/src/pullRequest/PullRequestProvider.ts index cc370637c601..6cd9f0330646 100644 --- a/apps/server/src/pullRequest/PullRequestProvider.ts +++ b/apps/server/src/pullRequest/PullRequestProvider.ts @@ -248,6 +248,7 @@ export interface ProviderChangeRequestActivity { */ readonly commentCount: number; readonly commentsTruncated: boolean; + readonly reviewThreadsTruncated?: boolean; readonly reviewThreads: ReadonlyArray; readonly commits: ReadonlyArray; /** The change request's own reactions, from a host that has them. */ diff --git a/apps/server/src/pullRequest/PullRequestService.test.ts b/apps/server/src/pullRequest/PullRequestService.test.ts index 7415ae45bfae..1344c3d27cee 100644 --- a/apps/server/src/pullRequest/PullRequestService.test.ts +++ b/apps/server/src/pullRequest/PullRequestService.test.ts @@ -4318,7 +4318,8 @@ it.effect( return Effect.succeed({ comments: [], commentCount: 0, - commentsTruncated: false, + commentsTruncated: true, + reviewThreadsTruncated: true, reviewThreads: [], commits: [], }); @@ -4342,10 +4343,14 @@ it.effect( }, ]); - yield* Effect.all([service.activity(reference), service.activity(reference)], { - concurrency: 2, - }); + const activities = yield* Effect.all( + [service.activity(reference), service.activity(reference)], + { + concurrency: 2, + }, + ); assert.strictEqual(activityCalls, 1); + assert.isTrue(activities.every((activity) => activity.reviewThreadsTruncated === true)); yield* service.invalidate({ reference }); yield* service.activity(reference); diff --git a/apps/server/src/pullRequest/PullRequestService.ts b/apps/server/src/pullRequest/PullRequestService.ts index 293318f91c4b..f444c9dc543b 100644 --- a/apps/server/src/pullRequest/PullRequestService.ts +++ b/apps/server/src/pullRequest/PullRequestService.ts @@ -1778,6 +1778,9 @@ export const make = Effect.gen(function* () { comments: activity.comments, commentCount: activity.commentCount, commentsTruncated: activity.commentsTruncated, + ...(activity.reviewThreadsTruncated === undefined + ? {} + : { reviewThreadsTruncated: activity.reviewThreadsTruncated }), reviewThreads: activity.reviewThreads, commits: activity.commits, ...(activity.reactions === undefined ? {} : { reactions: activity.reactions }), diff --git a/apps/server/src/pullRequest/gitHubPullRequestJson.ts b/apps/server/src/pullRequest/gitHubPullRequestJson.ts index d1600754b50b..ff352f743bb4 100644 --- a/apps/server/src/pullRequest/gitHubPullRequestJson.ts +++ b/apps/server/src/pullRequest/gitHubPullRequestJson.ts @@ -2093,6 +2093,7 @@ export interface GitHubReviewThreadComments { /** The host's own count of the conversation, which a bounded read can fall short of. */ readonly commentCount: number; readonly truncated: boolean; + readonly reviewThreadsTruncated: boolean; /** The pull request's own reactions, which sit on its description. */ readonly reactions: ReadonlyArray; /** Reactions by node id, for the comments and reviews the `gh` JSON read carries no reaction on. */ diff --git a/packages/contracts/src/pullRequest.ts b/packages/contracts/src/pullRequest.ts index c8feaad9a41f..369a2d7e5202 100644 --- a/packages/contracts/src/pullRequest.ts +++ b/packages/contracts/src/pullRequest.ts @@ -916,6 +916,8 @@ export const PullRequestActivity = Schema.Struct({ * however long it is. */ commentsTruncated: Schema.Boolean, + /** Whether the thread listing itself is incomplete, apart from pages within a thread. */ + reviewThreadsTruncated: Schema.optional(Schema.Boolean), reviewThreads: Schema.Array(PullRequestReviewThread), commits: Schema.Array(PullRequestCommit), /** From c1daaef03638d1b1560d7f3cabaa651573d5730f Mon Sep 17 00:00:00 2001 From: maria-rcks <254055478+maria-rcks@users.noreply.github.com> Date: Sun, 4 Oct 2026 02:56:36 +0000 Subject: [PATCH 3/3] fix(server): page a watched review thread again only when its count moves The watcher swept every minute and re-read every page of every long review thread each time. Keep the pages past the first per watch and reuse them until the host's comment count for that thread changes. --- .../PullRequestWatchReactor.ts | 70 ++++++++++++------- .../src/orchestration-v2/runtimeLayer.test.ts | 5 ++ 2 files changed, 50 insertions(+), 25 deletions(-) diff --git a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts index 9e5e9b0631b0..6e0046c5e59a 100644 --- a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts +++ b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts @@ -3,7 +3,9 @@ import { MessageId, type OrchestrationV2Notification, type PullRequestActivity, + type PullRequestComment, type PullRequestRef, + type PullRequestThreadCommentsResult, type ThreadPullRequestLink, type ThreadPullRequestWatch, } from "@t3tools/contracts"; @@ -83,6 +85,12 @@ export const make = Effect.gen(function* () { // Passes in a row that failed, per watch. Kept in memory: a restart only delays the stop. const readFailures = new Map(); + // Replies past each long thread's first page, per watch, so a pass pages a thread again only + // when the host's count of it moves. Kept in memory: a restart pages each thread once more. + const threadTails = new Map< + string, + Map }> + >(); // Host-level identity, with the repository as linked, the way pull request sync reads it. const identityOf = (link: ThreadPullRequestLink) => ({ @@ -128,40 +136,51 @@ export const make = Effect.gen(function* () { }).pipe(Effect.catch(() => record(target, null))); const readRemarks = Effect.fn("PullRequestWatchReactor.readRemarks")( - function* (reference: PullRequestRef, activity: PullRequestActivity) { + function* (target: WatchTarget, reference: PullRequestRef, activity: PullRequestActivity) { // Comment cursors cannot account for missing threads. Only finish a truncated read // when the host confirms that every thread was listed. if (activity.commentsTruncated && activity.reviewThreadsTruncated !== false) return null; - const comments = new Map(activity.comments.map((comment) => [comment.id, comment])); + const key = failureKey(target); + let tails = threadTails.get(key); + if (tails === undefined) { + tails = new Map(); + threadTails.set(key, tails); + } + const remarks = [...activity.comments]; for (const thread of activity.reviewThreads) { - const ids = new Set(thread.comments.map((comment) => comment.id)); - const cursors = new Set(); let cursor = thread.nextCommentsCursor ?? null; - while (cursor !== null) { - if (cursors.has(cursor)) return null; - cursors.add(cursor); - const page = yield* pullRequests.threadComments({ - ...reference, - threadId: thread.id, - cursor, - }); - for (const comment of page.comments) { - ids.add(comment.id); - comments.set(comment.id, { - ...comment, - kind: "review-comment", - path: thread.path, - reviewState: null, + if (cursor === null) continue; + const count = thread.commentCount ?? 0; + let tail = tails.get(thread.id); + if (tail?.count !== count) { + const comments = new Map(); + const cursors = new Set(); + while (cursor !== null) { + if (cursors.has(cursor)) return null; + cursors.add(cursor); + const page: PullRequestThreadCommentsResult = yield* pullRequests.threadComments({ + ...reference, + threadId: thread.id, + cursor, }); + for (const comment of page.comments) { + comments.set(comment.id, { + ...comment, + kind: "review-comment", + path: thread.path, + reviewState: null, + }); + } + cursor = page.nextCursor; } - cursor = page.nextCursor; + if (thread.comments.length + comments.size < count) return null; + tail = { count, comments: [...comments.values()] }; + tails.set(thread.id, tail); } - if (ids.size < (thread.commentCount ?? 0)) return null; + remarks.push(...tail.comments); } - return [...comments.values()].sort((left, right) => - left.createdAt.localeCompare(right.createdAt), - ); + return remarks.sort((left, right) => left.createdAt.localeCompare(right.createdAt)); }, Effect.catch((error) => Effect.logWarning("pull request watch comment pagination failed", { error }).pipe( @@ -206,7 +225,7 @@ export const make = Effect.gen(function* () { if (detail.state !== "open") return yield* record(target, null); // Never advance the remark watermark past comments an incomplete read could have missed. - const remarks = yield* readRemarks(reference, activity); + const remarks = yield* readRemarks(target, reference, activity); const report = evaluatePullRequestWatch(watch, detail, remarks); if (report.changes.length > 0) { return yield* record( @@ -233,6 +252,7 @@ export const make = Effect.gen(function* () { ); const keys = new Set(targets.map(failureKey)); for (const key of readFailures.keys()) if (!keys.has(key)) readFailures.delete(key); + for (const key of threadTails.keys()) if (!keys.has(key)) threadTails.delete(key); yield* Effect.forEach( targets, (target) => diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 0e4eb5452b82..84d6df90160a 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -2438,6 +2438,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }; const incomplete = mode !== "single page" && mode !== "paginated"; let recovering = false; + let pagesRead = 0; const reactor = yield* PullRequestWatchReactor.make.pipe( Effect.provide( Layer.mergeAll( @@ -2472,6 +2473,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { commits: [], }), threadComments: (input) => { + pagesRead += 1; assert.equal(input.threadId, "review-thread"); if (input.cursor === "after-10") { return Effect.succeed({ comments: [eleventh], nextCursor: "after-11" }); @@ -2514,7 +2516,10 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { recovering = true; } yield* reactor.sweep; + // A thread whose count has not moved is not paged again. + const pagesBefore = pagesRead; yield* reactor.sweep; + assert.equal(pagesRead, pagesBefore); const { messages } = yield* orchestrator.getThreadRecords(threadId, ["messages"]); assert.deepEqual(