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
75 changes: 68 additions & 7 deletions apps/server/src/orchestration-v2/PullRequestWatchReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,10 @@ import {
CommandId,
MessageId,
type OrchestrationV2Notification,
type PullRequestActivity,
type PullRequestComment,
type PullRequestRef,
type PullRequestThreadCommentsResult,
type ThreadPullRequestLink,
type ThreadPullRequestWatch,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -81,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<string, number>();
// 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<string, { readonly count: number; readonly comments: ReadonlyArray<PullRequestComment> }>
>();

// Host-level identity, with the repository as linked, the way pull request sync reads it.
const identityOf = (link: ThreadPullRequestLink) => ({
Expand Down Expand Up @@ -125,6 +135,60 @@ export const make = Effect.gen(function* () {
},
}).pipe(Effect.catch(() => record(target, null)));

const readRemarks = Effect.fn("PullRequestWatchReactor.readRemarks")(
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 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) {
let cursor = thread.nextCommentsCursor ?? null;
if (cursor === null) continue;
const count = thread.commentCount ?? 0;
let tail = tails.get(thread.id);
if (tail?.count !== count) {
const comments = new Map<string, PullRequestComment>();
const cursors = new Set<string>();
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;
}
if (thread.comments.length + comments.size < count) return null;
tail = { count, comments: [...comments.values()] };
tails.set(thread.id, tail);
}
remarks.push(...tail.comments);
}
return remarks.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);
Expand Down Expand Up @@ -160,13 +224,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(target, reference, activity);
const report = evaluatePullRequestWatch(watch, detail, remarks);
if (report.changes.length > 0) {
return yield* record(
target,
Expand All @@ -192,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) =>
Expand Down
147 changes: 124 additions & 23 deletions apps/server/src/orchestration-v2/runtimeLayer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@ import {
type OrchestrationV2Run,
ProjectId,
type PullRequestDetail,
type PullRequestComment,
PullRequestOperationError,
ProviderDriverKind,
ProviderInstanceId,
ProviderThreadId,
Expand Down Expand Up @@ -2307,11 +2309,18 @@ 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",
"thread list truncated",
])("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",
Expand All @@ -2323,7 +2332,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",
Expand All @@ -2337,15 +2346,15 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => {
const url = "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/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,
source: "agent",
});
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,
Expand Down Expand Up @@ -2398,6 +2407,38 @@ 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;
let pagesRead = 0;
const reactor = yield* PullRequestWatchReactor.make.pipe(
Effect.provide(
Layer.mergeAll(
Expand All @@ -2406,41 +2447,101 @@ 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",
reviewThreadsTruncated: mode === "thread list truncated" && !recovering,
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) => {
pagesRead += 1;
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"],
);
// Keep the two notifications ordered independently of their random message IDs.
yield* TestClock.adjust("1 millis");
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(
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 },
);
}),
);
Expand Down
11 changes: 10 additions & 1 deletion apps/server/src/pullRequest/GitHubPullRequestCli.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -3625,6 +3632,7 @@ layer("GitHubPullRequestCli.layer", (it) => {

assert.strictEqual(mockedExecute.mock.calls.length, 10);
assert.isTrue(conversation.truncated);
assert.isTrue(conversation.reviewThreadsTruncated);
}),
);

Expand Down Expand Up @@ -3656,6 +3664,7 @@ layer("GitHubPullRequestCli.layer", (it) => {
nextCommentsCursor: "Y3Vyc29yOjI",
});
assert.isTrue(conversation.truncated);
assert.isFalse(conversation.reviewThreadsTruncated);
}),
);

Expand Down
1 change: 1 addition & 0 deletions apps/server/src/pullRequest/GitHubPullRequestCli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/pullRequest/GitHubPullRequestProvider.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -882,6 +882,7 @@ describe("getChangeRequest commits", () => {
reviewThreads: [],
commentCount: 0,
truncated: false,
reviewThreadsTruncated: false,
reactions: [],
reactionsById: new Map<string, ReadonlyArray<PullRequestReaction>>(),
reviewers: [],
Expand Down Expand Up @@ -967,6 +968,7 @@ describe("getChangeRequestActivity dismissed reviews", () => {
reviewThreads: [],
commentCount: 0,
truncated: false,
reviewThreadsTruncated: false,
reactions: [],
reactionsById: new Map(),
reviewers: [],
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/pullRequest/GitHubPullRequestProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -411,6 +411,7 @@ export const make = Effect.gen(function* () {
reviewThreads: [],
commentCount: 0,
truncated: true,
reviewThreadsTruncated: true,
reviewers: [],
avatarsByLogin: new Map<string, string>(),
botLogins: new Set<string>(),
Expand Down Expand Up @@ -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) => ({
Expand Down
Loading
Loading