Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
f5a4410
fix(pull-requests): refresh data after thread turns
maria-rcks Sep 3, 2026
54077db
fix(pull-requests): refresh open diff after turns
maria-rcks Sep 3, 2026
d2126f2
fix(pull-requests): keep turn refreshes scoped
maria-rcks Sep 3, 2026
bf50eac
fix(pull-requests): isolate panel turn refreshes
maria-rcks Sep 3, 2026
69fe35c
fix(pull-requests): preserve terminal turn refreshes
maria-rcks Sep 3, 2026
c181635
fix(pull-requests): ignore mid-turn refresh placeholders
maria-rcks Sep 3, 2026
8b943a9
fix(pull-requests): refresh sessionless completions
maria-rcks Sep 3, 2026
a90e325
refactor(pull-requests): simplify turn refreshes
maria-rcks Sep 4, 2026
d91144c
chore(pull-requests): merge main into refresh branch
maria-rcks Sep 4, 2026
65b1ef2
fix(pull-requests): refresh startup completions
maria-rcks Sep 4, 2026
3b54ca4
fix(pull-requests): bound turn refresh state
maria-rcks Sep 4, 2026
0ac7540
fix(pull-requests): refresh diff on interval
maria-rcks Sep 4, 2026
e588f96
refactor(pull-requests): simplify turn refresh trigger
maria-rcks Sep 4, 2026
19c9de8
fix(pull-requests): preserve terminal refresh ordering
maria-rcks Sep 4, 2026
b372bed
fix(pull-requests): keep primary turn refresh tracked
maria-rcks Sep 4, 2026
ea106a9
fix(pull-requests): track steered primary turns
maria-rcks Sep 4, 2026
6571dd6
fix(pull-requests): validate replacement turn starts
maria-rcks Sep 4, 2026
2f435f3
fix(pull-requests): require pending replacement turns
maria-rcks Sep 4, 2026
e977b4c
refactor(pull-requests): trim refresh plumbing
maria-rcks Sep 4, 2026
c8a4a30
fix(pull-requests): replay refresh revision
maria-rcks Sep 4, 2026
951ad84
fix(pull-requests): ignore baseline refresh signal
maria-rcks Sep 4, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ import { VcsStatusBroadcaster } from "../src/vcs/VcsStatusBroadcaster.ts";
import { GitWorkflowService } from "../src/git/GitWorkflowService.ts";
import * as VcsProcess from "../src/vcs/VcsProcess.ts";
import * as AgentAwarenessRelay from "../src/relay/AgentAwarenessRelay.ts";
import * as PullRequestService from "../src/pullRequest/PullRequestService.ts";

const decodeCodexSettings = Schema.decodeEffect(CodexSettings);

Expand Down Expand Up @@ -350,6 +351,11 @@ export const makeOrchestrationIntegrationHarness = (
);
const checkpointReactorLayer = CheckpointReactorLive.pipe(
Layer.provideMerge(runtimeServicesLayer),
Layer.provideMerge(
Layer.mock(PullRequestService.PullRequestService)({
refreshAfterTurn: () => Effect.void,
}),
),
Layer.provideMerge(
Layer.succeed(VcsStatusBroadcaster, {
getStatus: () => Effect.die("getStatus should not be called in this test"),
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ export const RPC_REQUIRED_SCOPES = {
// Read scope like the reads it un-caches: refreshing is part of reading, and a read-only
// client pressing refresh must not be told it may not look again.
[WS_METHODS.pullRequestsInvalidate]: AuthOrchestrationReadScope,
[WS_METHODS.pullRequestsSubscribeRefreshes]: AuthOrchestrationReadScope,
// The candidate list is a read like the detail beside it; asking somebody for a review is a
// write like every other one.
[WS_METHODS.pullRequestsReviewerCandidates]: AuthOrchestrationReadScope,
Expand Down
94 changes: 94 additions & 0 deletions apps/server/src/orchestration/Layers/CheckpointReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ import { ProviderValidationError } from "../../provider/Errors.ts";
import { ServerConfig } from "../../config.ts";
import * as WorkspaceEntries from "../../workspace/WorkspaceEntries.ts";
import * as WorkspacePaths from "../../workspace/WorkspacePaths.ts";
import { PullRequestService } from "../../pullRequest/PullRequestService.ts";

const asProjectId = (value: string): ProjectId => ProjectId.make(value);
const asTurnId = (value: string): TurnId => TurnId.make(value);
Expand Down Expand Up @@ -294,6 +295,7 @@ describe("CheckpointReactor", () => {
readonly providerSessionCwd?: string;
readonly providerName?: ProviderDriverKind;
readonly gitStatusRefreshCalls?: Array<string>;
readonly pullRequestRefreshCalls?: Array<ProjectId>;
}) {
const cwd = createGitRepository();
tempDirs.push(cwd);
Expand Down Expand Up @@ -348,6 +350,14 @@ describe("CheckpointReactor", () => {
Layer.provideMerge(projectionSnapshotLayer),
Layer.provideMerge(RuntimeReceiptBusLive),
Layer.provideMerge(Layer.succeed(ProviderService, provider.service)),
Layer.provideMerge(
Layer.mock(PullRequestService)({
refreshAfterTurn: (projectId) =>
Effect.sync(() => {
options?.pullRequestRefreshCalls?.push(projectId);
}),
}),
),
Layer.provideMerge(vcsStatusBroadcasterLayer),
Layer.provideMerge(CheckpointStore.layer.pipe(Layer.provide(VcsDriverRegistry.layer))),
Layer.provideMerge(
Expand Down Expand Up @@ -538,6 +548,90 @@ describe("CheckpointReactor", () => {
).toBe("v2\n");
});

effectIt.effect("refreshes pull request data once when a running turn terminates", () =>
Effect.gen(function* () {
const pullRequestRefreshCalls: ProjectId[] = [];
const harness = yield* Effect.promise(() =>
createHarness({
seedFilesystemCheckpoints: false,
pullRequestRefreshCalls,
}),
);
const threadId = ThreadId.make("thread-1");
const turnId = asTurnId("turn-refresh-prs");

yield* harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-pr-refresh-running"),
threadId,
session: {
threadId,
status: "running",
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: turnId,
lastError: null,
updatedAt: "2026-01-01T00:00:01.000Z",
},
createdAt: "2026-01-01T00:00:01.000Z",
});
yield* Effect.promise(() => harness.drain());
expect(pullRequestRefreshCalls).toEqual([]);

yield* harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-pr-refresh-ready"),
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: null,
lastError: null,
updatedAt: "2026-01-01T00:00:02.000Z",
},
createdAt: "2026-01-01T00:00:02.000Z",
});
yield* Effect.promise(() => harness.drain());

yield* harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-pr-refresh-ready-again"),
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: null,
lastError: null,
updatedAt: "2026-01-01T00:00:03.000Z",
},
createdAt: "2026-01-01T00:00:03.000Z",
});
yield* Effect.promise(() => harness.drain());

expect(pullRequestRefreshCalls).toEqual([asProjectId("project-1")]);

yield* harness.engine.dispatch({
type: "thread.turn.diff.complete",
commandId: CommandId.make("cmd-pr-refresh-missed-start"),
threadId,
turnId: asTurnId("turn-with-missed-start"),
completedAt: "2026-01-01T00:00:04.000Z",
checkpointRef: checkpointRefForThreadTurn(threadId, 2),
checkpointTurnCount: 2,
status: "ready",
files: [],
createdAt: "2026-01-01T00:00:04.000Z",
});
yield* Effect.promise(() => harness.drain());

expect(pullRequestRefreshCalls).toEqual([asProjectId("project-1"), asProjectId("project-1")]);
}),
);

it("refreshes local git status state on turn completion using the session cwd", async () => {
const gitStatusRefreshCalls: string[] = [];
const harness = await createHarness({
Expand Down
51 changes: 51 additions & 0 deletions apps/server/src/orchestration/Layers/CheckpointReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ import type { OrchestrationDispatchError } from "../Errors.ts";
import { isGitRepository } from "../../git/Utils.ts";
import { VcsStatusBroadcaster } from "../../vcs/VcsStatusBroadcaster.ts";
import * as WorkspaceEntries from "../../workspace/WorkspaceEntries.ts";
import * as PullRequestService from "../../pullRequest/PullRequestService.ts";

const nowIso = Effect.map(DateTime.now, DateTime.formatIso);

Expand Down Expand Up @@ -88,6 +89,24 @@ const make = Effect.gen(function* () {
const receiptBus = yield* RuntimeReceiptBus;
const workspaceEntries = yield* WorkspaceEntries.WorkspaceEntries;
const vcsStatusBroadcaster = yield* VcsStatusBroadcaster;
const pullRequests = yield* PullRequestService.PullRequestService;
const refreshedTurnKeys = new Set<string>();
const REFRESHED_TURN_CAPACITY = 2_048;

const refreshPullRequestsOnce = Effect.fn("refreshPullRequestsOnce")(function* (
projectId: ProjectId,
threadId: ThreadId,
turnId: TurnId,
) {
const key = `${threadId}:${turnId}`;
if (refreshedTurnKeys.has(key)) return;
if (refreshedTurnKeys.size >= REFRESHED_TURN_CAPACITY) {
const oldest = refreshedTurnKeys.values().next().value;
if (oldest !== undefined) refreshedTurnKeys.delete(oldest);
}
refreshedTurnKeys.add(key);
yield* pullRequests.refreshAfterTurn(projectId);
});

const appendRevertFailureActivity = (input: {
readonly threadId: ThreadId;
Expand Down Expand Up @@ -820,6 +839,27 @@ const make = Effect.gen(function* () {
});

const processDomainEvent = Effect.fn("processDomainEvent")(function* (event: OrchestrationEvent) {
if (event.type === "thread.session-set") {
if (
event.payload.session.status === "starting" ||
event.payload.session.status === "running"
) {
return;
}
const thread = yield* resolveThreadDetail(event.payload.threadId);
const turn = thread?.latestTurn;
if (
thread !== undefined &&
turn !== null &&
turn !== undefined &&
turn.state !== "running" &&
turn.completedAt === event.payload.session.updatedAt
) {
yield* refreshPullRequestsOnce(thread.projectId, thread.id, turn.turnId);
}
return;
Comment thread
cursor[bot] marked this conversation as resolved.
Outdated
}

if (event.type === "thread.turn-start-requested" || event.type === "thread.message-sent") {
yield* ensurePreTurnBaselineFromDomainTurnStart(event);
return;
Expand Down Expand Up @@ -847,6 +887,16 @@ const make = Effect.gen(function* () {
// turn.completed runtime events to this reactor (shared subscription), so
// reacting to the domain event is the reliable path.
if (event.type === "thread.turn-diff-completed") {
const thread = yield* resolveThreadDetail(event.payload.threadId);
if (
thread !== undefined &&
thread.session !== null &&
thread.session.status !== "starting" &&
thread.session.status !== "running" &&
thread.latestTurn?.turnId === event.payload.turnId
) {
yield* refreshPullRequestsOnce(thread.projectId, thread.id, event.payload.turnId);
}
Comment thread
cursor[bot] marked this conversation as resolved.
Outdated
yield* captureCheckpointFromPlaceholder(event).pipe(
Effect.catch((error) =>
Effect.flatMap(nowIso, (createdAt) =>
Expand Down Expand Up @@ -920,6 +970,7 @@ const make = Effect.gen(function* () {
if (
event.type !== "thread.turn-start-requested" &&
event.type !== "thread.message-sent" &&
event.type !== "thread.session-set" &&
event.type !== "thread.checkpoint-revert-requested" &&
event.type !== "thread.turn-diff-completed"
) {
Expand Down
67 changes: 67 additions & 0 deletions apps/server/src/pullRequest/PullRequestService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2681,6 +2681,73 @@ it.effect("an explicit invalidation makes the next listing ask the host again",
}),
);

it.effect("refreshes listing and change request caches after a thread turn", () =>
Effect.scoped(
Effect.gen(function* () {
let listCalls = 0;
let unrelatedListCalls = 0;
let detailCalls = 0;
let unrelatedDetailCalls = 0;
const reference = { projectId: "p1" as ProjectId, repository: "acme/web", number: 1 };
const unrelatedReference = {
projectId: "p2" as ProjectId,
repository: "acme/api",
number: 1,
};
const service = yield* makeService({
projects: [
project({ id: "p1", title: "web", workspaceRoot: "/a", repository: "acme/web" }),
project({ id: "p2", title: "api", workspaceRoot: "/b", repository: "acme/api" }),
],
providers: [
fakeProvider("github", {
listChangeRequests: (input) => {
const calls = input.repository === "acme/web" ? ++listCalls : ++unrelatedListCalls;
return Effect.succeed({
items: [changeRequest(calls, "2026-07-02T00:00:00Z")],
truncated: false,
continues: false,
});
},
getChangeRequest: (input) => {
const calls =
input.repository === "acme/web" ? ++detailCalls : ++unrelatedDetailCalls;
return Effect.succeed(hostedChangeRequest(`body ${calls}`));
},
}),
],
});
const observedRefresh = yield* Stream.runHead(service.subscribeRefreshes).pipe(
Effect.forkChild({ startImmediately: true }),
);

const firstList = yield* service.list({ state: "open", projectId: reference.projectId });
yield* service.list({ state: "open", projectId: unrelatedReference.projectId });
const firstDetail = yield* service.detail(reference);
const unrelatedDetail = yield* service.detail(unrelatedReference);
assert.strictEqual(firstList.entries[0]?.number, 1);
assert.strictEqual(firstDetail.body, "body 1");
assert.strictEqual(unrelatedDetail.body, "body 1");

yield* service.refreshAfterTurn(reference.projectId);
const refresh = Option.getOrThrow(yield* Fiber.join(observedRefresh));
const refreshedList = yield* service.list({ state: "open", projectId: reference.projectId });
yield* service.list({ state: "open", projectId: unrelatedReference.projectId });
const refreshedDetail = yield* service.detail(reference);
yield* service.detail(unrelatedReference);

assert.isAbove(refresh.revision, 0);
assert.strictEqual(refresh.projectId, reference.projectId);
assert.strictEqual(refreshedList.entries[0]?.number, 2);
assert.strictEqual(refreshedDetail.body, "body 2");
assert.strictEqual(listCalls, 2);
assert.strictEqual(unrelatedListCalls, 1);
assert.strictEqual(detailCalls, 2);
assert.strictEqual(unrelatedDetailCalls, 1);
}),
),
);

it.effect("a mutation makes the next listing ask the host again, with no client asking", () =>
Effect.gen(function* () {
let hostCalls = 0;
Expand Down
Loading
Loading