Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
109 changes: 94 additions & 15 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,8 @@ describe("ProviderCommandReactor", () => {
readonly requiresNewThreadForModelChange?: boolean;
readonly titleRegenerationCompletionDispatchFailures?: number;
readonly titleRegenerationBeforeStart?: "one" | "two";
readonly interruptTurnEffect?: () => Effect.Effect<void, ProviderAdapterRequestError>;
readonly stopSessionEffect?: () => Effect.Effect<void, ProviderAdapterRequestError>;
readonly startSessionEffect?: (
session: ProviderSession,
) => Effect.Effect<ProviderSession, ProviderAdapterRequestError>;
Expand Down Expand Up @@ -235,23 +237,27 @@ describe("ProviderCommandReactor", () => {
turnId: asTurnId("turn-1"),
}),
);
const interruptTurn = vi.fn((_: unknown) => Effect.void);
const interruptTurn = vi.fn((_: unknown) => input?.interruptTurnEffect?.() ?? Effect.void);
const respondToRequest = vi.fn<ProviderServiceShape["respondToRequest"]>(() => Effect.void);
const respondToUserInput = vi.fn<ProviderServiceShape["respondToUserInput"]>(() => Effect.void);
const stopSession = vi.fn((input: unknown) =>
Effect.sync(() => {
const threadId =
typeof input === "object" && input !== null && "threadId" in input
? (input as { threadId?: ThreadId }).threadId
: undefined;
if (!threadId) {
return;
}
const index = runtimeSessions.findIndex((session) => session.threadId === threadId);
if (index >= 0) {
runtimeSessions.splice(index, 1);
}
}),
const stopSession = vi.fn((stopInput: unknown) =>
(input?.stopSessionEffect?.() ?? Effect.void).pipe(
Effect.tap(() =>
Effect.sync(() => {
const threadId =
typeof stopInput === "object" && stopInput !== null && "threadId" in stopInput
? (stopInput as { threadId?: ThreadId }).threadId
: undefined;
if (!threadId) {
return;
}
const index = runtimeSessions.findIndex((session) => session.threadId === threadId);
if (index >= 0) {
runtimeSessions.splice(index, 1);
}
}),
),
),
);
const renameBranch = vi.fn((input: unknown) =>
Effect.succeed({
Expand Down Expand Up @@ -2486,6 +2492,79 @@ describe("ProviderCommandReactor", () => {
});
});

it("stops a running session and records the failure when provider interrupt fails", async () => {
const harness = await createHarness({
interruptTurnEffect: () =>
Effect.fail(
new ProviderAdapterRequestError({
provider: "codex",
method: "thread.interrupt",
detail: "provider session disappeared",
}),
),
stopSessionEffect: () =>
Effect.fail(
new ProviderAdapterRequestError({
provider: "codex",
method: "session.stop",
detail: "provider process already exited",
}),
),
});
const now = "2026-01-01T00:00:00.000Z";

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-set-interrupt-failure"),
threadId: ThreadId.make("thread-1"),
session: {
threadId: ThreadId.make("thread-1"),
status: "running",
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: asTurnId("turn-1"),
lastError: null,
updatedAt: now,
},
createdAt: now,
}),
);

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.turn.interrupt",
commandId: CommandId.make("cmd-turn-interrupt-provider-failure"),
threadId: ThreadId.make("thread-1"),
turnId: asTurnId("turn-1"),
createdAt: now,
}),
);

await waitFor(async () => {
const thread = (await harness.readModel()).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
return thread?.session?.status === "stopped";
});

const thread = (await harness.readModel()).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(thread?.session).toMatchObject({
status: "stopped",
activeTurnId: null,
lastError: "provider session disappeared",
});
expect(
thread?.activities.find((activity) => activity.kind === "provider.turn.interrupt.failed"),
).toMatchObject({
summary: "Provider turn interrupt failed",
payload: { detail: "provider session disappeared" },
});
expect(harness.stopSession).toHaveBeenCalledWith({ threadId: ThreadId.make("thread-1") });
});

it("starts a fresh session when only projected session state exists", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
52 changes: 49 additions & 3 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1180,8 +1180,8 @@ const make = Effect.gen(function* () {
if (!thread) {
return;
}
const hasSession = thread.session && thread.session.status !== "stopped";
if (!hasSession) {
const session = thread.session;
if (!session || session.status === "stopped") {
return yield* appendProviderFailureActivity({
threadId: event.payload.threadId,
kind: "provider.turn.interrupt.failed",
Expand All @@ -1192,8 +1192,54 @@ const make = Effect.gen(function* () {
});
}

const recoverInterruptFailure = (cause: Cause.Cause<unknown>) => {
if (Cause.hasInterruptsOnly(cause)) {
return Effect.interrupt;
}

const detail = formatFailureDetail(cause);
return Effect.gen(function* () {
yield* providerService.stopSession({ threadId: event.payload.threadId }).pipe(
Effect.catchCause((stopCause) => {
if (Cause.hasInterruptsOnly(stopCause)) {
return Effect.interrupt;
}
return Effect.logWarning(
"provider command reactor failed to stop session after interrupt failure",
{
threadId: event.payload.threadId,
cause: Cause.pretty(stopCause),
originalCause: Cause.pretty(cause),
},
);
}),
);
yield* setThreadSession({
threadId: event.payload.threadId,
session: {
...session,
status: "stopped",
activeTurnId: null,
lastError: detail,
updatedAt: event.payload.createdAt,
},
createdAt: event.payload.createdAt,
});
Comment thread
cursor[bot] marked this conversation as resolved.
yield* appendProviderFailureActivity({
threadId: event.payload.threadId,
kind: "provider.turn.interrupt.failed",
summary: "Provider turn interrupt failed",
detail,
turnId: event.payload.turnId ?? null,
createdAt: event.payload.createdAt,
});
});
};

// Orchestration turn ids are not provider turn ids, so interrupt by session.
yield* providerService.interruptTurn({ threadId: event.payload.threadId });
yield* providerService
.interruptTurn({ threadId: event.payload.threadId })
.pipe(Effect.catchCause(recoverInterruptFailure));
});

const processApprovalResponseRequested = Effect.fn("processApprovalResponseRequested")(function* (
Expand Down
Loading