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
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,7 @@ function createProviderServiceHarness(
interruptTurn: () => unsupported(),
respondToRequest: () => unsupported(),
respondToUserInput: () => unsupported(),
discardTransientThread: () => unsupported(),
stopSession: () => unsupported(),
listSessions,
getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -321,6 +321,7 @@ describe("ProviderCommandReactor", () => {
interruptTurn: interruptTurn as ProviderServiceShape["interruptTurn"],
respondToRequest: respondToRequest as ProviderServiceShape["respondToRequest"],
respondToUserInput: respondToUserInput as ProviderServiceShape["respondToUserInput"],
discardTransientThread: () => unsupported(),
stopSession: stopSession as ProviderServiceShape["stopSession"],
listSessions: () => Effect.succeed(runtimeSessions),
getCapabilities: (_provider) =>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,7 @@ function createProviderServiceHarness() {
interruptTurn: () => unsupported(),
respondToRequest: () => unsupported(),
respondToUserInput: () => unsupported(),
discardTransientThread: () => unsupported(),
stopSession: () => unsupported(),
listSessions: () => Effect.succeed([...runtimeSessions]),
getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import type { ProjectionThread } from "../../persistence/Services/ProjectionThre
import {
logCleanupCauseUnlessInterrupted,
managedWorktreeCleanupTarget,
shouldDiscardTransientProviderThread,
} from "./ThreadDeletionReactor.ts";

const threadId = ThreadId.make("thread-deletion-reactor-test");
Expand Down Expand Up @@ -160,3 +161,10 @@ describe("managedWorktreeCleanupTarget", () => {
).toBeNull();
});
});

describe("shouldDiscardTransientProviderThread", () => {
it("selects only durable first-message retraction deletions", () => {
expect(shouldDiscardTransientProviderThread(deletedEvent())).toBe(true);
expect(shouldDiscardTransientProviderThread(deletedEvent(false))).toBe(false);
});
});
14 changes: 14 additions & 0 deletions apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,10 @@ import { forkParked } from "../../serverActivation.ts";

type ThreadDeletedEvent = Extract<OrchestrationEvent, { type: "thread.deleted" }>;

export function shouldDiscardTransientProviderThread(event: ThreadDeletedEvent): boolean {
return event.payload.retraction?.firstUserMessage === true;
}

export function managedWorktreeCleanupTarget(input: {
readonly event: ThreadDeletedEvent;
readonly thread: ProjectionThread;
Expand Down Expand Up @@ -77,6 +81,13 @@ const make = Effect.gen(function* () {
threadId,
});

const discardTransientProviderThread = (threadId: ThreadDeletedEvent["payload"]["threadId"]) =>
logCleanupCauseUnlessInterrupted({
effect: providerService.discardTransientThread({ threadId }),
message: "thread retraction cleanup skipped transient provider thread discard",
threadId,
});

const closeThreadTerminals = (threadId: ThreadDeletedEvent["payload"]["threadId"]) =>
logCleanupCauseUnlessInterrupted({
effect: terminalManager.close({ threadId, deleteHistory: true }),
Expand Down Expand Up @@ -113,6 +124,9 @@ const make = Effect.gen(function* () {
event: ThreadDeletedEvent,
) {
const { threadId } = event.payload;
if (shouldDiscardTransientProviderThread(event)) {
yield* discardTransientProviderThread(threadId);
}
yield* stopProviderSession(threadId);
yield* closeThreadTerminals(threadId);
yield* removeRetractedManagedWorktree(event);
Expand Down
50 changes: 50 additions & 0 deletions apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -324,6 +324,7 @@ async function startHarness(
}).pipe(Effect.andThen(state.interruptAcknowledgementHangs ? Effect.never : Effect.void)),
respondToRequest: () => unsupported(),
respondToUserInput: () => unsupported(),
discardTransientThread: () => unsupported(),
stopSession: () => unsupported(),
listSessions: () => Effect.succeed([]),
getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }),
Expand Down Expand Up @@ -451,6 +452,55 @@ it("completes a cancelled provider-send path after filesystem convergence", asyn
await stopHarness(harness);
});

it("settles a claimed first-message turn before restoring and completing without rollback", async () => {
const state = makeState("claimed");
state.row = pendingRow("claimed", true);
state.sessionStatus = "starting";
state.activeTurnId = null;
const harness = await startHarness(state);

expect(state.order).toEqual([]);
expect(state.row.status).toBe("requested");

state.sessionStatus = "running";
state.activeTurnId = TURN_ID;
await harness.emitRuntime({
type: "turn.started",
eventId: EventId.make("evt-first-message-turn-started"),
provider: ProviderDriverKind.make("codex"),
createdAt: NOW,
threadId: THREAD_ID,
turnId: TURN_ID,
payload: {},
});
await harness.runtime.runPromise(Effect.yieldNow);
await harness.runtime.runPromise(harness.reactor.drain);
expect(state.order).toEqual(["interrupt"]);

state.sessionStatus = "ready";
state.activeTurnId = null;
await harness.retryTick();
await harness.runtime.runPromise(harness.reactor.drain);

expect(state.filesystemRestored).toBe(true);
expect(state.row.status).toBe("completed");
expect(state.order).toEqual(["interrupt", "restore", "complete"]);
expect(state.interruptedTurnIds).toEqual([TURN_ID]);
expect(state.rollbackTargetTurnIds).toEqual([]);
await stopHarness(harness);
});

it("restores and completes a cancelled first-message send without stopping a provider", async () => {
const state = makeState("cancelled");
state.row = pendingRow("cancelled", true);
const harness = await startHarness(state);

expect(state.filesystemRestored).toBe(true);
expect(state.row.status).toBe("completed");
expect(state.order).toEqual(["restore", "complete"]);
await stopHarness(harness);
});

it("drives claimed convergence from interrupt through a settlement event", async () => {
const state = makeState("claimed");
const harness = await startHarness(state);
Expand Down
15 changes: 15 additions & 0 deletions apps/server/src/orchestration/Layers/TurnRetractionReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -545,6 +545,21 @@ export const makeTurnRetractionReactor = Effect.gen(function* () {
}
}

// A first-message retraction deletes the thread, so its settled provider
// conversation does not need a rollback boundary. Do not stop the session
// while ProviderCommandReactor may still be starting the send: the normal
// start/interrupt settlement above lets that durable worker finish. The
// correlated thread.deleted event owns final provider-session cleanup.
if (row.firstUserMessage) {
yield* restoreFilesystem(row, false);
yield* dispatchCompletion(row, targetTurnId);
clearIssuedInterrupts(row.requestId);
yield* logConvergence(row, "cleanup", "completed", {
action: "restore-and-delete-settled-first-message-thread",
});
return;
}

yield* providerService
.rollbackConversationTo({
threadId: row.threadId,
Expand Down
25 changes: 25 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,8 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape {
}),
);

public readonly deleteThreadImpl = vi.fn((): Promise<void> => Promise.resolve(undefined));

public readonly respondToRequestImpl = vi.fn(
(_requestId: ApprovalRequestId, _decision: ProviderApprovalDecision): Promise<void> =>
Promise.resolve(undefined),
Expand Down Expand Up @@ -141,6 +143,8 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape {
return Effect.promise(() => this.rollbackThreadImpl(numTurns));
}

deleteThread = Effect.promise(() => this.deleteThreadImpl());

respondToRequest(requestId: ApprovalRequestId, decision: ProviderApprovalDecision) {
return Effect.promise(() => this.respondToRequestImpl(requestId, decision));
}
Expand Down Expand Up @@ -309,6 +313,27 @@ const sessionErrorLayer = it.layer(
);

sessionErrorLayer("CodexAdapterLive session errors", (it) => {
it.effect("discards the active provider-owned thread without stopping its session", () =>
Effect.gen(function* () {
const adapter = yield* CodexAdapter;
const threadId = asThreadId("discard-transient-thread");
yield* adapter.startSession({
provider: ProviderDriverKind.make("codex"),
threadId,
runtimeMode: "full-access",
});
const runtime = sessionRuntimeFactory.lastRuntime;
NodeAssert.ok(runtime);
NodeAssert.ok(adapter.discardTransientThread);

yield* adapter.discardTransientThread(threadId);

NodeAssert.equal(runtime.deleteThreadImpl.mock.calls.length, 1);
NodeAssert.equal(runtime.closeImpl.mock.calls.length, 0);
NodeAssert.equal(yield* adapter.hasSession(threadId), true);
}),
);

it.effect("computes the remaining absolute rollback delta and is idempotent", () =>
Effect.gen(function* () {
const adapter = yield* CodexAdapter;
Expand Down
13 changes: 13 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2028,6 +2028,18 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
yield* stopSessionInternal(session);
});

const discardTransientThread: NonNullable<CodexAdapterShape["discardTransientThread"]> = (
threadId,
) =>
requireSession(threadId).pipe(
Effect.flatMap((session) => session.runtime.deleteThread),
Effect.mapError((cause) =>
cause._tag === "ProviderAdapterSessionNotFoundError"
? cause
: mapCodexRuntimeError(threadId, "thread/delete", cause),
),
);

const listSessions: CodexAdapterShape["listSessions"] = () =>
Effect.forEach(
Array.from(sessions.values()).filter((session) => !session.stopped),
Expand Down Expand Up @@ -2066,6 +2078,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
respondToRequest,
respondToUserInput,
stopSession,
discardTransientThread,
listSessions,
hasSession,
stopAll,
Expand Down
24 changes: 24 additions & 0 deletions apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import {
import { codexSessionAppServerArgs } from "./codexLaunchArgs.ts";
import {
buildTurnStartParams,
deleteCodexThread,
hasConfiguredMcpServer,
isRecoverableThreadResumeError,
openCodexThread,
Expand Down Expand Up @@ -469,3 +470,26 @@ describe("openCodexThread", () => {
}),
);
});

describe("deleteCodexThread", () => {
it.effect("uses the provider thread id with thread/delete", () =>
Effect.gen(function* () {
const calls: Array<{ method: string; payload: unknown }> = [];
const client = {
request: (method: "thread/delete", payload: { readonly threadId: string }) => {
calls.push({ method, payload });
return Effect.succeed({});
},
};

yield* deleteCodexThread(client, "provider-thread-transient");

NodeAssert.deepStrictEqual(calls, [
{
method: "thread/delete",
payload: { threadId: "provider-thread-transient" },
},
]);
}),
);
});
20 changes: 20 additions & 0 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,7 @@ export interface CodexSessionRuntimeShape {
readonly rollbackThread: (
numTurns: number,
) => Effect.Effect<CodexThreadSnapshot, CodexSessionRuntimeError>;
readonly deleteThread: Effect.Effect<void, CodexSessionRuntimeError>;
readonly respondToRequest: (
requestId: ApprovalRequestId,
decision: ProviderApprovalDecision,
Expand Down Expand Up @@ -454,6 +455,22 @@ interface CodexThreadOpenClient {
) => Effect.Effect<CodexRpc.ClientRequestResponsesByMethod[M], CodexErrors.CodexAppServerError>;
}

interface CodexThreadDeleteClient {
readonly request: (
method: "thread/delete",
payload: CodexRpc.ClientRequestParamsByMethod["thread/delete"],
) => Effect.Effect<
CodexRpc.ClientRequestResponsesByMethod["thread/delete"],
CodexErrors.CodexAppServerError
>;
}

export const deleteCodexThread = (
client: CodexThreadDeleteClient,
threadId: string,
): Effect.Effect<void, CodexErrors.CodexAppServerError> =>
client.request("thread/delete", { threadId }).pipe(Effect.asVoid);

export const openCodexThread = (input: {
readonly client: CodexThreadOpenClient;
readonly threadId: ThreadId;
Expand Down Expand Up @@ -1853,6 +1870,9 @@ export const makeCodexSessionRuntime = (
});
return parseThreadSnapshot(response);
}),
deleteThread: Effect.flatMap(readProviderThreadId, (providerThreadId) =>
deleteCodexThread(client, providerThreadId),
),
respondToRequest: (requestId, decision) =>
Effect.gen(function* () {
const pending = (yield* Ref.get(pendingApprovalsRef)).get(requestId);
Expand Down
51 changes: 50 additions & 1 deletion apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,10 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) {
}),
);

const discardTransientThread = vi.fn(
(_threadId: ThreadId): Effect.Effect<void, ProviderAdapterError> => Effect.void,
);

const listSessions = vi.fn(
(): Effect.Effect<ReadonlyArray<ProviderSession>> =>
Effect.sync(() => Array.from(sessions.values())),
Expand Down Expand Up @@ -223,6 +227,7 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) {
respondToRequest,
respondToUserInput,
stopSession,
...(provider === CODEX_DRIVER ? { discardTransientThread } : {}),
listSessions,
hasSession,
readThread,
Expand Down Expand Up @@ -259,6 +264,7 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) {
respondToRequest,
respondToUserInput,
stopSession,
discardTransientThread,
listSessions,
hasSession,
readThread,
Expand Down Expand Up @@ -904,8 +910,14 @@ routing.layer("ProviderServiceLive routing", (it) => {
assert.equal(Option.isSome(persisted), true);
if (Option.isSome(persisted)) {
assert.deepEqual(persisted.value.resumeCursor, rolledBackCursor);
const runtimePayload = persisted.value.runtimePayload;
assert.equal(
persisted.value.runtimePayload.lastRuntimeEvent,
runtimePayload !== null &&
typeof runtimePayload === "object" &&
!Array.isArray(runtimePayload) &&
"lastRuntimeEvent" in runtimePayload
? runtimePayload.lastRuntimeEvent
: undefined,
"provider.rollbackConversationTo",
);
}
Expand All @@ -921,6 +933,43 @@ routing.layer("ProviderServiceLive routing", (it) => {
}),
);

it.effect("discards only provider threads whose adapter explicitly supports it", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
const codexThreadId = asThreadId("transient-codex-thread");
const claudeThreadId = asThreadId("transient-claude-thread");
yield* provider.startSession(codexThreadId, {
provider: CODEX_DRIVER,
providerInstanceId: codexInstanceId,
threadId: codexThreadId,
runtimeMode: "full-access",
});
yield* provider.startSession(claudeThreadId, {
provider: CLAUDE_AGENT_DRIVER,
providerInstanceId: claudeAgentInstanceId,
threadId: claudeThreadId,
runtimeMode: "full-access",
});
routing.codex.discardTransientThread.mockClear();
routing.claude.discardTransientThread.mockClear();

yield* provider.discardTransientThread({ threadId: codexThreadId });
yield* provider.discardTransientThread({ threadId: claudeThreadId });

assert.deepEqual(routing.codex.discardTransientThread.mock.calls, [[codexThreadId]]);
assert.equal(routing.claude.discardTransientThread.mock.calls.length, 0);

yield* provider.stopSession({ threadId: codexThreadId });
yield* provider.stopSession({ threadId: claudeThreadId });
routing.codex.startSession.mockClear();
routing.codex.stopSession.mockClear();
routing.codex.discardTransientThread.mockClear();
routing.claude.startSession.mockClear();
routing.claude.stopSession.mockClear();
routing.claude.discardTransientThread.mockClear();
}),
);

it.effect("routes provider operations and rollback conversation", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
Expand Down
Loading
Loading