Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -888,7 +888,7 @@ it.live("reverts to an earlier checkpoint and trims checkpoint projections + git
);

it.live(
"appends checkpoint.revert.failed activity when revert is requested without an active session",
"appends checkpoint.revert.failed activity when revert has no persisted provider binding",
() =>
withHarness((harness) =>
Effect.gen(function* () {
Expand Down Expand Up @@ -917,7 +917,7 @@ it.live(
assert.equal(
String(
(failureActivity?.payload as { readonly detail?: string } | undefined)?.detail,
).includes("No active provider session"),
).includes("no persisted provider binding"),
true,
);
}),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ const startupDependencies = Layer.mergeAll(
stopSession: () => Effect.die("unused"),
listSessions: () => Effect.succeed([]),
getCapabilities: () => Effect.die("unused"),
assertConversationRollbackSupported: () => Effect.die("unused"),
prepareConversationRollback: () => Effect.die("unused"),
getInstanceInfo: () => Effect.die("unused"),
rollbackConversation: () => Effect.die("unused"),
uploadFeedback: () => Effect.die("unused"),
Expand Down
126 changes: 117 additions & 9 deletions apps/server/src/orchestration/Layers/CheckpointReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,9 +94,9 @@ function createProviderServiceHarness(
const rollbackConversation = vi.fn(
(_input: { readonly threadId: ThreadId; readonly numTurns: number }) => Effect.void,
);
const assertConversationRollbackSupported = vi.fn<
ProviderServiceShape["assertConversationRollbackSupported"]
>(() => Effect.void);
const prepareConversationRollback = vi.fn<ProviderServiceShape["prepareConversationRollback"]>(
() => Effect.void,
);

const unsupported = <A>() =>
Effect.die(new Error("Unsupported provider call in test")) as Effect.Effect<A, never>;
Expand Down Expand Up @@ -124,7 +124,7 @@ function createProviderServiceHarness(
stopSession: () => unsupported(),
listSessions,
getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }),
assertConversationRollbackSupported,
prepareConversationRollback,
getInstanceInfo: (instanceId) =>
Effect.succeed({
instanceId,
Expand All @@ -149,7 +149,10 @@ function createProviderServiceHarness(

return {
service,
assertConversationRollbackSupported,
setSessionActive: (active: boolean) => {
hasSession = active;
},
prepareConversationRollback,
rollbackConversation,
emit,
};
Expand Down Expand Up @@ -1608,12 +1611,12 @@ describe("CheckpointReactor", () => {
const threadId = ThreadId.make("thread-1");
const createdAt = "2026-01-01T00:00:00.000Z";
const checked = yield* Deferred.make<void>();
harness.provider.assertConversationRollbackSupported.mockImplementation(() =>
harness.provider.prepareConversationRollback.mockImplementation(() =>
Deferred.succeed(checked, undefined).pipe(
Effect.andThen(
Effect.fail(
new ProviderValidationError({
operation: "ProviderService.assertConversationRollbackSupported",
operation: "ProviderService.prepareConversationRollback",
issue: "Provider 'antigravity' does not support conversation rewind.",
}),
),
Expand Down Expand Up @@ -1683,6 +1686,111 @@ describe("CheckpointReactor", () => {
}),
);

effectIt.effect.each([
{ provider: "codex", turnCount: 0, recover: true },
{ provider: "codex", turnCount: 1, recover: true },
{ provider: "claudeAgent", turnCount: 0, recover: true },
{ provider: "claudeAgent", turnCount: 1, recover: true },
{ provider: "claudeAgent", turnCount: 1, recover: false },
])(
"prepares an inactive $provider session before reverting to $turnCount (recover=$recover)",
({ provider, turnCount, recover }) =>
Effect.gen(function* () {
const harness = yield* Effect.promise(() =>
createHarness({
hasSession: false,
providerName: ProviderDriverKind.make(provider),
}),
);
const threadId = ThreadId.make("thread-1");
const createdAt = "2026-01-01T00:00:00.000Z";
yield* harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-recovery-ready-session"),
threadId,
session: {
threadId,
status: "ready",
providerName: provider,
runtimeMode: "approval-required",
activeTurnId: null,
lastError: null,
updatedAt: createdAt,
},
createdAt,
});
for (const checkpointTurnCount of [1, 2]) {
yield* harness.engine.dispatch({
type: "thread.turn.diff.complete",
commandId: CommandId.make(`cmd-recovery-diff-${checkpointTurnCount}`),
threadId,
turnId: asTurnId(`turn-${checkpointTurnCount}`),
completedAt: createdAt,
checkpointRef: checkpointRefForThreadTurn(threadId, checkpointTurnCount),
status: "ready",
files: [],
checkpointTurnCount,
createdAt,
});
}
const before = (yield* Effect.promise(() => harness.readModel())).threads[0];
const prepared = yield* Deferred.make<void>();
harness.provider.prepareConversationRollback.mockImplementation(() =>
Deferred.succeed(prepared, undefined).pipe(
Effect.andThen(
recover
? Effect.sync(() => harness.provider.setSessionActive(true))
: Effect.fail(
new ProviderValidationError({
operation: "ProviderService.prepareConversationRollback",
issue: "Persisted session could not be resumed.",
}),
),
),
),
);
yield* harness.engine.dispatch({
type: "thread.checkpoint.revert",
commandId: CommandId.make("cmd-recovery-revert"),
threadId,
turnCount,
createdAt,
});
yield* Deferred.await(prepared);
yield* Effect.promise(() => harness.drain());
const after = (yield* Effect.promise(() => harness.readModel())).threads[0];
if (recover) {
expect(after?.checkpoints).toHaveLength(turnCount);
expect(
after?.activities.some((activity) => activity.kind === "checkpoint.revert.failed"),
).toBe(false);
expect(harness.provider.rollbackConversation).toHaveBeenCalledWith({
threadId,
numTurns: 2 - turnCount,
});
expect(NodeFS.readFileSync(NodePath.join(harness.cwd, "README.md"), "utf8")).toBe(
turnCount === 0 ? "v1\n" : "v2\n",
);
expect(gitRefExists(harness.cwd, checkpointRefForThreadTurn(threadId, 2))).toBe(false);
} else {
expect(after?.checkpoints).toEqual(before?.checkpoints);
expect(after?.messages).toEqual(before?.messages);
expect(after?.latestTurn).toEqual(before?.latestTurn);
expect(after?.activities).toContainEqual(
expect.objectContaining({
kind: "checkpoint.revert.failed",
payload: expect.objectContaining({
detail: expect.stringContaining("Persisted session could not be resumed."),
}),
}),
);
expect(harness.provider.rollbackConversation).not.toHaveBeenCalled();
expect(NodeFS.readFileSync(NodePath.join(harness.cwd, "README.md"), "utf8")).toBe("v3\n");
expect(gitRefExists(harness.cwd, checkpointRefForThreadTurn(threadId, 2))).toBe(true);
}
}),
);

it("executes provider revert and emits thread.reverted for checkpoint revert requests", async () => {
const harness = await createHarness();
const createdAt = "2026-01-01T00:00:00.000Z";
Expand Down Expand Up @@ -1916,7 +2024,7 @@ describe("CheckpointReactor", () => {
});
});

it("appends an error activity when revert is requested without an active session", async () => {
it("appends an error activity when rollback preparation leaves no session cwd", async () => {
const harness = await createHarness({ hasSession: false });
const createdAt = "2026-01-01T00:00:00.000Z";

Expand All @@ -1925,7 +2033,7 @@ describe("CheckpointReactor", () => {
type: "thread.checkpoint.revert",
commandId: CommandId.make("cmd-revert-no-session"),
threadId: ThreadId.make("thread-1"),
turnCount: 1,
turnCount: 0,
createdAt,
}),
);
Expand Down
42 changes: 21 additions & 21 deletions apps/server/src/orchestration/Layers/CheckpointReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -699,26 +699,6 @@ const make = Effect.gen(function* () {
return;
}

const sessionRuntime = yield* resolveSessionRuntimeForThread(event.payload.threadId);
if (Option.isNone(sessionRuntime)) {
yield* appendRevertFailureActivity({
threadId: event.payload.threadId,
turnCount: event.payload.turnCount,
detail: "No active provider session with workspace cwd is bound to this thread.",
createdAt: now,
}).pipe(Effect.catch(() => Effect.void));
return;
}
if (!(yield* checkpointStore.isGitRepository(sessionRuntime.value.cwd))) {
yield* appendRevertFailureActivity({
threadId: event.payload.threadId,
turnCount: event.payload.turnCount,
detail: "Checkpoints are unavailable because this project is not a git repository.",
createdAt: now,
}).pipe(Effect.catch(() => Effect.void));
return;
}

const currentTurnCount = thread.checkpoints.reduce(
(maxTurnCount, checkpoint) => Math.max(maxTurnCount, checkpoint.checkpointTurnCount),
0,
Expand Down Expand Up @@ -751,7 +731,27 @@ const make = Effect.gen(function* () {
return;
}

yield* providerService.assertConversationRollbackSupported(event.payload.threadId);
yield* providerService.prepareConversationRollback(event.payload.threadId);

const sessionRuntime = yield* resolveSessionRuntimeForThread(event.payload.threadId);
if (Option.isNone(sessionRuntime)) {
yield* appendRevertFailureActivity({
threadId: event.payload.threadId,
turnCount: event.payload.turnCount,
detail: "No active provider session with workspace cwd is bound to this thread.",
createdAt: now,
}).pipe(Effect.catch(() => Effect.void));
return;
}
if (!(yield* checkpointStore.isGitRepository(sessionRuntime.value.cwd))) {
yield* appendRevertFailureActivity({
threadId: event.payload.threadId,
turnCount: event.payload.turnCount,
detail: "Checkpoints are unavailable because this project is not a git repository.",
createdAt: now,
}).pipe(Effect.catch(() => Effect.void));
return;
}

const restored = yield* checkpointStore.restoreCheckpoint({
cwd: sessionRuntime.value.cwd,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -363,7 +363,7 @@ describe("ProviderCommandReactor", () => {
Effect.succeed({
sessionModelSwitch: input?.sessionModelSwitch ?? "in-session",
}),
assertConversationRollbackSupported: () => unsupported(),
prepareConversationRollback: () => unsupported(),
getInstanceInfo: (instanceId) => {
const raw = String(instanceId);
const driverKind = ProviderDriverKind.make(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ function createProviderServiceHarness() {
stopSession: () => unsupported(),
listSessions: () => Effect.succeed([...runtimeSessions]),
getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }),
assertConversationRollbackSupported: () => unsupported(),
prepareConversationRollback: () => unsupported(),
getInstanceInfo: (instanceId) => {
const driverKind = ProviderDriverKind.make(String(instanceId));
return Effect.succeed({
Expand Down
Loading
Loading