Skip to content
Closed
Show file tree
Hide file tree
Changes from 3 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
71 changes: 67 additions & 4 deletions apps/server/src/orchestration/Layers/CheckpointReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -147,21 +147,24 @@ async function waitForThread(
readonly threads: ReadonlyArray<{
readonly id: ThreadId;
readonly latestTurn: { readonly turnId: string } | null;
readonly checkpoints: ReadonlyArray<{ readonly checkpointTurnCount: number }>;
readonly checkpoints: ReadonlyArray<{
readonly checkpointTurnCount: number;
readonly completedAt: string;
}>;
readonly activities: ReadonlyArray<{ readonly kind: string }>;
}>;
}>,
predicate: (thread: {
latestTurn: { turnId: string } | null;
checkpoints: ReadonlyArray<{ checkpointTurnCount: number }>;
checkpoints: ReadonlyArray<{ checkpointTurnCount: number; completedAt: string }>;
activities: ReadonlyArray<{ kind: string }>;
}) => boolean,
timeoutMs = 15_000,
) {
const deadline = (await Effect.runPromise(Clock.currentTimeMillis)) + timeoutMs;
const poll = async (): Promise<{
latestTurn: { turnId: string } | null;
checkpoints: ReadonlyArray<{ checkpointTurnCount: number }>;
checkpoints: ReadonlyArray<{ checkpointTurnCount: number; completedAt: string }>;
activities: ReadonlyArray<{ kind: string }>;
}> => {
const snapshot = await readModel();
Expand All @@ -180,7 +183,7 @@ async function waitForThread(

async function waitForEvent(
engine: OrchestrationEngineShape,
predicate: (event: { type: string }) => boolean,
predicate: (event: { type: string; occurredAt: string }) => boolean,
timeoutMs = 15_000,
) {
const deadline = (await Effect.runPromise(Clock.currentTimeMillis)) + timeoutMs;
Expand Down Expand Up @@ -364,6 +367,10 @@ describe("CheckpointReactor", () => {
scope = await Effect.runPromise(Scope.make("sequential"));
await Effect.runPromise(reactor.start().pipe(Scope.provide(scope)));
const drain = () => Effect.runPromise(reactor.drain);
const captureCheckpoint = (input: CheckpointStore.CaptureCheckpointInput) =>
runtime!.runPromise(checkpointStore.captureCheckpoint(input));
const dispatch = (command: Parameters<OrchestrationEngineShape["dispatch"]>[0]) =>
runtime!.runPromise(engine.dispatch(command));

const createdAt = "2026-01-01T00:00:00.000Z";
await Effect.runPromise(
Expand Down Expand Up @@ -449,6 +456,8 @@ describe("CheckpointReactor", () => {
engine,
readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()),
provider,
captureCheckpoint,
dispatch,
cwd,
drain,
};
Expand Down Expand Up @@ -530,6 +539,60 @@ describe("CheckpointReactor", () => {
).toBe("v2\n");
});

it("replaces a legacy mid-turn checkpoint with the filesystem state at completion", async () => {
const harness = await createHarness({ seedFilesystemCheckpoints: false });
const threadId = ThreadId.make("thread-1");
const turnId = asTurnId("turn-1");
const checkpointRef = checkpointRefForThreadTurn(threadId, 1);
const midTurnAt = "2026-01-01T00:00:10.000Z";
const completedAt = "2026-01-01T00:00:30.000Z";

await harness.captureCheckpoint({
cwd: harness.cwd,
checkpointRef: checkpointRefForThreadTurn(threadId, 0),
});
NodeFS.writeFileSync(NodePath.join(harness.cwd, "README.md"), "mid-turn\n", "utf8");
await harness.captureCheckpoint({ cwd: harness.cwd, checkpointRef });
await harness.dispatch({
type: "thread.turn.diff.complete",
commandId: CommandId.make("cmd-legacy-mid-turn-diff"),
threadId,
turnId,
completedAt: midTurnAt,
checkpointRef,
status: "ready",
files: [{ path: "README.md", kind: "modified", additions: 1, deletions: 1 }],
checkpointTurnCount: 1,
createdAt: midTurnAt,
});
await harness.drain();
expect(gitShowFileAtRef(harness.cwd, checkpointRef, "README.md")).toBe("mid-turn\n");

NodeFS.writeFileSync(NodePath.join(harness.cwd, "README.md"), "final\n", "utf8");
harness.provider.emit({
type: "turn.completed",
eventId: EventId.make("evt-turn-completed-after-legacy-checkpoint"),
provider: ProviderDriverKind.make("codex"),
createdAt: completedAt,
threadId,
turnId,
payload: { state: "completed" },
});

await waitForEvent(
harness.engine,
(event) => event.type === "thread.turn-diff-completed" && event.occurredAt === completedAt,
);
const thread = await waitForThread(
harness.readModel,
(entry) =>
entry.checkpoints.length === 1 && entry.checkpoints[0]?.completedAt === completedAt,
);

expect(thread.checkpoints[0]?.checkpointTurnCount).toBe(1);
expect(gitShowFileAtRef(harness.cwd, checkpointRef, "README.md")).toBe("final\n");
});

it("refreshes local git status state on turn completion using the session cwd", async () => {
const gitStatusRefreshCalls: string[] = [];
const harness = await createHarness({
Expand Down
109 changes: 11 additions & 98 deletions apps/server/src/orchestration/Layers/CheckpointReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -369,13 +369,16 @@ const make = Effect.gen(function* () {
return;
}

// Only skip if a real (non-placeholder) checkpoint already exists for this turn.
// ProviderRuntimeIngestion may insert placeholder entries with status "missing"
// before this reactor runs; those must not prevent real git capture.
const existingCheckpoint = thread.checkpoints.find(
Comment thread
MohtashamMurshid marked this conversation as resolved.
(checkpoint) => checkpoint.turnId === turnId,
);
// A replayed completion event must not recapture an already-final checkpoint.
// Older servers could capture on a mid-turn diff update, so a checkpoint
// timestamped before the completion event still needs to be replaced.
if (
thread.checkpoints.some(
(checkpoint) => checkpoint.turnId === turnId && checkpoint.status !== "missing",
)
existingCheckpoint &&
existingCheckpoint.status !== "missing" &&
existingCheckpoint.completedAt === event.createdAt
) {
return;
}
Expand All @@ -391,18 +394,11 @@ const make = Effect.gen(function* () {
return;
}

// If a placeholder checkpoint exists for this turn, reuse its turn count
// instead of incrementing past it.
const existingPlaceholder = thread.checkpoints.find(
(checkpoint) => checkpoint.turnId === turnId && checkpoint.status === "missing",
);
const currentTurnCount = thread.checkpoints.reduce(
(maxTurnCount, checkpoint) => Math.max(maxTurnCount, checkpoint.checkpointTurnCount),
0,
);
const nextTurnCount = existingPlaceholder
? existingPlaceholder.checkpointTurnCount
: currentTurnCount + 1;
const nextTurnCount = existingCheckpoint?.checkpointTurnCount ?? currentTurnCount + 1;

yield* captureAndDispatchCheckpoint({
threadId: thread.id,
Expand All @@ -417,68 +413,6 @@ const make = Effect.gen(function* () {
},
);

// Captures a real git checkpoint when a placeholder checkpoint (status "missing")
// is detected via a domain event. This replaces the placeholder with a real
// git-ref-based checkpoint.
//
// ProviderRuntimeIngestion creates placeholder checkpoints on turn.diff.updated
// events from the Codex runtime. This handler fires when the corresponding
// domain event arrives, allowing the reactor to capture the actual filesystem
// state into a git ref and dispatch a replacement checkpoint.
const captureCheckpointFromPlaceholder = Effect.fn("captureCheckpointFromPlaceholder")(function* (
event: Extract<OrchestrationEvent, { type: "thread.turn-diff-completed" }>,
) {
const { threadId, turnId, checkpointTurnCount, status } = event.payload;

// Only replace placeholders; skip events from our own real captures.
if (status !== "missing") {
return;
}

const thread = yield* resolveThreadDetail(threadId);
if (!thread) {
yield* Effect.logWarning("checkpoint capture from placeholder skipped: thread not found", {
threadId,
});
return;
}

// If a real checkpoint already exists for this turn, skip.
if (
thread.checkpoints.some(
(checkpoint) => checkpoint.turnId === turnId && checkpoint.status !== "missing",
)
) {
yield* Effect.logDebug(
"checkpoint capture from placeholder skipped: real checkpoint already exists",
{ threadId, turnId },
);
return;
}

const projects = yield* resolveThreadProjects(thread.projectId);
const checkpointCwd = yield* resolveCheckpointCwd({
threadId,
thread,
projects,
preferSessionRuntime: true,
});
if (!checkpointCwd) {
return;
}

yield* captureAndDispatchCheckpoint({
threadId,
turnId,
thread,
cwd: checkpointCwd,
turnCount: checkpointTurnCount,
status: "ready",
assistantMessageId: event.payload.assistantMessageId ?? undefined,
createdAt: event.payload.completedAt,
});
});

const ensurePreTurnBaselineFromTurnStart = Effect.fn("ensurePreTurnBaselineFromTurnStart")(
function* (event: Extract<ProviderRuntimeEvent, { type: "turn.started" }>) {
const turnId = toTurnId(event.turnId);
Expand Down Expand Up @@ -838,26 +772,6 @@ const make = Effect.gen(function* () {
);
return;
}

// When ProviderRuntimeIngestion creates a placeholder checkpoint (status "missing")
// from a turn.diff.updated runtime event, capture the real git checkpoint to
// replace it. The providerService.streamEvents PubSub does not reliably deliver
// 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") {
yield* captureCheckpointFromPlaceholder(event).pipe(
Effect.catch((error) =>
Effect.flatMap(nowIso, (createdAt) =>
appendCaptureFailureActivity({
threadId: event.payload.threadId,
turnId: event.payload.turnId,
detail: error.message,
createdAt,
}).pipe(Effect.catch(() => Effect.void)),
),
),
);
}
});

const processRuntimeEvent = Effect.fn("processRuntimeEvent")(function* (
Expand Down Expand Up @@ -918,8 +832,7 @@ const make = Effect.gen(function* () {
if (
event.type !== "thread.turn-start-requested" &&
event.type !== "thread.message-sent" &&
event.type !== "thread.checkpoint-revert-requested" &&
event.type !== "thread.turn-diff-completed"
event.type !== "thread.checkpoint-revert-requested"
) {
return Effect.void;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,6 @@ type ProviderRuntimeTestThread = ProviderRuntimeTestReadModel["threads"][number]
type ProviderRuntimeTestMessage = ProviderRuntimeTestThread["messages"][number];
type ProviderRuntimeTestProposedPlan = ProviderRuntimeTestThread["proposedPlans"][number];
type ProviderRuntimeTestActivity = ProviderRuntimeTestThread["activities"][number];
type ProviderRuntimeTestCheckpoint = ProviderRuntimeTestThread["checkpoints"][number];

async function waitForThread(
readModel: () => Promise<ProviderRuntimeTestReadModel>,
Expand Down Expand Up @@ -2856,7 +2855,7 @@ describe("ProviderRuntimeIngestion", () => {
});
});

it("consumes P1 runtime events into thread metadata, diff checkpoints, and activities", async () => {
it("consumes P1 runtime events without checkpointing mid-turn diff updates", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";

Expand Down Expand Up @@ -2906,28 +2905,28 @@ describe("ProviderRuntimeIngestion", () => {
});

harness.emit({
type: "runtime.warning",
eventId: asEventId("evt-runtime-warning"),
type: "turn.diff.updated",
eventId: asEventId("evt-turn-diff-updated"),
provider: ProviderDriverKind.make("codex"),
createdAt: now,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-p1"),
itemId: asItemId("item-p1-assistant"),
payload: {
message: "Provider got slow",
detail: { latencyMs: 1500 },
unifiedDiff: "diff --git a/file.txt b/file.txt\n+hello\n",
},
});

harness.emit({
type: "turn.diff.updated",
eventId: asEventId("evt-turn-diff-updated"),
type: "runtime.warning",
eventId: asEventId("evt-runtime-warning"),
provider: ProviderDriverKind.make("codex"),
createdAt: now,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-p1"),
itemId: asItemId("item-p1-assistant"),
payload: {
unifiedDiff: "diff --git a/file.txt b/file.txt\n+hello\n",
message: "Provider got slow",
detail: { latencyMs: 1500 },
},
});

Expand All @@ -2943,9 +2942,6 @@ describe("ProviderRuntimeIngestion", () => {
) &&
entry.activities.some(
(activity: ProviderRuntimeTestActivity) => activity.kind === "runtime.warning",
) &&
entry.checkpoints.some(
(checkpoint: ProviderRuntimeTestCheckpoint) => checkpoint.turnId === "turn-p1",
),
);

Expand Down Expand Up @@ -2983,12 +2979,7 @@ describe("ProviderRuntimeIngestion", () => {
expect(warning?.kind).toBe("runtime.warning");
expect(warningPayload?.message).toBe("Provider got slow");

const checkpoint = thread.checkpoints.find(
(entry: ProviderRuntimeTestCheckpoint) => entry.turnId === "turn-p1",
);
expect(checkpoint?.status).toBe("missing");
expect(checkpoint?.assistantMessageId).toBe("assistant:item-p1-assistant");
expect(checkpoint?.checkpointRef).toBe("provider-diff:evt-turn-diff-updated");
expect(thread.checkpoints).toEqual([]);
});

it("mirrors a provider title only while the thread still has the default title", async () => {
Expand Down
Loading
Loading