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 @@ -1026,6 +1026,42 @@ describe("ProviderRuntimeIngestion", () => {
}),
);

it("settles an active turn when the provider session exits without completion", async () => {
const harness = await createHarness();
const threadId = asThreadId("thread-1");
const turnId = asTurnId("turn-provider-exit");

harness.emit({
type: "turn.started",
eventId: asEventId("evt-turn-started-before-provider-exit"),
provider: ProviderDriverKind.make("opencode"),
threadId,
turnId,
createdAt: "2026-01-01T00:00:00.000Z",
payload: {},
});
await waitForThread(harness.readModel, (entry) => entry.session?.activeTurnId === turnId);

harness.emit({
type: "session.exited",
eventId: asEventId("evt-provider-exited-without-completion"),
provider: ProviderDriverKind.make("opencode"),
threadId,
createdAt: "2026-01-01T00:00:01.000Z",
payload: {
reason: "Provider child process exited.",
recoverable: false,
exitKind: "error",
},
});

const thread = await waitForThread(
harness.readModel,
(entry) => entry.session?.status === "stopped" && entry.session.activeTurnId === null,
);
expect(thread.session?.activeTurnId).toBeNull();
});

it("does not clear active turn when session/thread started arrives mid-turn", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
71 changes: 71 additions & 0 deletions apps/server/src/provider/Layers/OpenCodeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import {
ProviderDriverKind,
ProviderInstanceId,
ThreadId,
TurnId,
} from "@t3tools/contracts";
import { createModelSelection } from "@t3tools/shared/model";
import { ServerConfig } from "../../config.ts";
Expand Down Expand Up @@ -1066,6 +1067,76 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => {
}),
);

for (const liveStatus of ["idle", "busy"] as const) {
it.effect(`resyncs a resumed open turn from confirmed ${liveStatus} provider status`, () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
const threadId = asThreadId(`thread-opencode-resume-${liveStatus}`);
const turnId = TurnId.make(`turn-opencode-resume-${liveStatus}`);
runtimeMock.state.sessionStatusImplementation = async () => ({
data: liveStatus === "busy" ? { ses_persisted: { type: "busy" } } : {},
});
const stateChanged = yield* adapter.streamEvents.pipe(
Stream.filter(
(event) => event.threadId === threadId && event.type === "session.state.changed",
),
Stream.runHead,
Effect.forkChild,
);

const session = yield* adapter.startSession({
provider: ProviderDriverKind.make("opencode"),
threadId,
runtimeMode: "full-access",
resumeCursor: { schemaVersion: 1, sessionId: "ses_persisted" },
resumedActiveTurnId: turnId,
});
const event = Option.getOrThrow(yield* Fiber.join(stateChanged));

NodeAssert.equal(event.type, "session.state.changed");
if (event.type === "session.state.changed") {
NodeAssert.equal(event.payload.state, liveStatus === "busy" ? "running" : "ready");
NodeAssert.equal(event.turnId, turnId);
}
NodeAssert.equal(session.status, liveStatus === "busy" ? "running" : "ready");
NodeAssert.equal(session.activeTurnId, liveStatus === "busy" ? turnId : undefined);

yield* adapter.stopSession(threadId);
}),
);
}

it.effect("leaves a resumed open turn untouched when live status cannot be qualified", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
const threadId = asThreadId("thread-opencode-resume-unknown");
runtimeMock.state.sessionStatusImplementation = async () => ({ data: null });
const lifecycle = yield* adapter.streamEvents.pipe(
Stream.filter((event) => event.threadId === threadId),
Stream.take(2),
Stream.runCollect,
Effect.forkChild,
);

yield* adapter.startSession({
provider: ProviderDriverKind.make("opencode"),
threadId,
runtimeMode: "full-access",
resumeCursor: { schemaVersion: 1, sessionId: "ses_persisted" },
resumedActiveTurnId: TurnId.make("turn-opencode-resume-unknown"),
});
const events = Array.from(yield* Fiber.join(lifecycle));

NodeAssert.deepEqual(
events.map((event) => event.type),
["session.started", "thread.started"],
);
NodeAssert.equal(runtimeMock.state.sessionStatusCalls, 1);

yield* adapter.stopSession(threadId);
}),
);

it.effect("sends follow-up turns to the resumed session id", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
Expand Down
54 changes: 54 additions & 0 deletions apps/server/src/provider/Layers/OpenCodeAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1480,6 +1480,57 @@ export function makeOpenCodeAdapter(
promptAdmission.recoveryFiber = yield* recover.pipe(Effect.forkIn(context.sessionScope));
});

const resyncResumedOpenCodeTurn = Effect.fn("resyncResumedOpenCodeTurn")(function* (
context: OpenCodeSessionContext,
turnId: TurnId,
) {
const response = yield* runOpenCodeSdk("session.status", (signal) =>
context.client.session.status(undefined, { signal }),
).pipe(Effect.timeout("1 second"), Effect.option);
const statuses = Option.isSome(response)
? Option.getOrUndefined(decodeOpenCodeSessionStatusMap(response.value.data))
: undefined;
if (statuses === undefined) {
return;
}
const status = statuses[context.openCodeSessionId];
const running = status?.type === "busy" || status?.type === "retry";
if (status !== undefined && !running && status.type !== "idle") {
return;
}

const updatedAt = yield* nowIso;
if (running) {
context.activeTurnId = turnId;
applyProviderSessionUpdate(
context,
{ status: "running", activeTurnId: turnId },
{ clearLastError: true },
updatedAt,
);
} else {
applyProviderSessionUpdate(
context,
{ status: "ready" },
{ clearActiveTurnId: true, clearLastError: true },
updatedAt,
);
}

yield* emit({
...(yield* buildEventBase({
threadId: context.session.threadId,
turnId,
createdAt: updatedAt,
})),
type: "session.state.changed",
payload: {
state: running ? "running" : "ready",
detail: { source: "resume-status-resync" },
},
});
});

const interruptOpenCodeTurn = Effect.fn("interruptOpenCodeTurn")(function* (
context: OpenCodeSessionContext,
turnId: TurnId,
Expand Down Expand Up @@ -3046,6 +3097,9 @@ export function makeOpenCodeAdapter(
}
yield* awaitOpenCodeContextReady(context);
if (!started.created) {
if (input.resumedActiveTurnId !== undefined) {
yield* resyncResumedOpenCodeTurn(context, input.resumedActiveTurnId);
}
yield* schedulePendingRequestRecovery(context);
}

Expand Down
32 changes: 28 additions & 4 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ import {
PROVIDER_SEND_TURN_MAX_INPUT_CHARS,
ProviderDriverKind,
ProviderInstanceId,
ProviderSessionStartInput,
ThreadId,
TurnId,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -60,7 +59,10 @@ import {
ProviderWorkspaceMissingError,
type ProviderAdapterError,
} from "../Errors.ts";
import type { ProviderAdapterShape } from "../Services/ProviderAdapter.ts";
import type {
ProviderAdapterSessionStartInput,
ProviderAdapterShape,
} from "../Services/ProviderAdapter.ts";
import * as ProviderAdapterRegistry from "../Services/ProviderAdapterRegistry.ts";
import * as ProviderService from "../Services/ProviderService.ts";
import * as ProviderSessionDirectory from "../Services/ProviderSessionDirectory.ts";
Expand Down Expand Up @@ -143,7 +145,7 @@ function makeFakeCodexAdapter(
const sessions = new Map<ThreadId, ProviderSession>();
const runtimeEventPubSub = Effect.runSync(PubSub.unbounded<ProviderRuntimeEvent>());

const startSession = vi.fn((input: ProviderSessionStartInput) =>
const startSession = vi.fn((input: ProviderAdapterSessionStartInput) =>
Effect.sync(() => {
const now = "2026-01-01T00:00:00.000Z";
const session: ProviderSession = {
Expand Down Expand Up @@ -2831,6 +2833,28 @@ routing.layer("ProviderServiceLive routing", (it) => {
}),
);

it.effect("passes the persisted open turn to a resumed provider session", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
const threadId = asThreadId("thread-resume-open-turn");
yield* provider.startSession(threadId, {
provider: ProviderDriverKind.make("codex"),
providerInstanceId: codexInstanceId,
threadId,
cwd: fixtureCwd("project-resume-open-turn"),
runtimeMode: "full-access",
});
const turn = yield* provider.sendTurn({ threadId, input: "work", attachments: [] });
yield* routing.codex.stopAll();
routing.codex.startSession.mockClear();

yield* provider.sendTurn({ threadId, input: "resume", attachments: [] });

const resumedStartInput = routing.codex.startSession.mock.calls[0]?.[0];
assert.equal(resumedStartInput?.resumedActiveTurnId, turn.turnId);
}),
);

it.effect("recovers stale claudeAgent sessions for sendTurn using persisted cwd", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
Expand Down Expand Up @@ -4806,7 +4830,7 @@ validation.layer("ProviderServiceLive validation", (it) => {
const provider = yield* ProviderService.ProviderService;
const runtimeRepository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository;

validation.codex.startSession.mockImplementationOnce((input: ProviderSessionStartInput) =>
validation.codex.startSession.mockImplementationOnce((input) =>
Effect.sync(() => {
const now = "2026-01-01T00:00:00.000Z";
return {
Expand Down
11 changes: 11 additions & 0 deletions apps/server/src/provider/Layers/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,9 @@ import * as ServerSettings from "../../serverSettings.ts";
import * as ProjectionSnapshotQuery from "../../orchestration/Services/ProjectionSnapshotQuery.ts";
const isModelSelection = Schema.is(ModelSelection);
const encodePromptJson = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));
const decodePersistedActiveTurn = Schema.decodeUnknownOption(
Schema.Struct({ activeTurnId: Schema.optional(TurnId) }),
);

interface SnapShotPromptAccessibilityNode {
readonly role: string;
Expand Down Expand Up @@ -427,6 +430,12 @@ function readPersistedCwd(
return trimmed.length > 0 ? trimmed : undefined;
}

function readPersistedActiveTurnId(
runtimePayload: ProviderSessionDirectory.ProviderRuntimeBinding["runtimePayload"],
): TurnId | undefined {
return Option.getOrUndefined(decodePersistedActiveTurn(runtimePayload))?.activeTurnId;
}

const dieOnMissingBindingInstanceId = (
operation: string,
payload: {
Expand Down Expand Up @@ -1268,6 +1277,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (

const persistedCwd = readPersistedCwd(input.binding.runtimePayload);
const persistedModelSelection = readPersistedModelSelection(input.binding.runtimePayload);
const resumedActiveTurnId = readPersistedActiveTurnId(input.binding.runtimePayload);

yield* prepareMcpSession(input.binding.threadId, bindingInstanceId);
const resumed = yield* adapter
Expand All @@ -1278,6 +1288,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
...(persistedCwd ? { cwd: persistedCwd } : {}),
...(persistedModelSelection ? { modelSelection: persistedModelSelection } : {}),
...(hasResumeCursor ? { resumeCursor: input.binding.resumeCursor } : {}),
...(resumedActiveTurnId ? { resumedActiveTurnId } : {}),
runtimeMode: input.binding.runtimeMode ?? "full-access",
})
.pipe(Effect.onError(() => clearMcpSession(input.binding.threadId)));
Expand Down
8 changes: 7 additions & 1 deletion apps/server/src/provider/Services/ProviderAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,12 @@ import type * as Stream from "effect/Stream";

export type ProviderSessionModelSwitchMode = "in-session" | "unsupported";

export type ProviderAdapterSessionStartInput = ProviderSessionStartInput & {
/** An orchestration turn remembered across provider-session loss. Adapters
may adopt it only after reconciling against live provider state. */
readonly resumedActiveTurnId?: TurnId;
};

/**
* How ProviderService runs manual context compaction for an adapter.
* Native adapters expose a start call and must emit a compacted thread state
Expand Down Expand Up @@ -75,7 +81,7 @@ export interface ProviderAdapterShape<TError> {
* Start a provider-backed session.
*/
readonly startSession: (
input: ProviderSessionStartInput,
input: ProviderAdapterSessionStartInput,
) => Effect.Effect<ProviderSession, TError>;

/**
Expand Down
Loading