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 @@ -7,10 +7,13 @@ import {
ProjectId,
ProviderDriverKind,
ProviderInstanceId,
type ProviderSendTurnInput,
ThreadId,
TurnId,
} from "@t3tools/contracts";
import { assert, it } from "@effect/vitest";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
Expand Down Expand Up @@ -361,3 +364,123 @@ it.effect(
),
),
);

it.effect.each(["opt-in desktop restart", "marked remote update"] as const)(
"continues a newer persisted turn after %s",
(restart) =>
Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
const createdAt = DateTime.formatIso(yield* DateTime.now);
const activeTurnId = TurnId.make("turn-started-after-original-send");
const originalTurnId = TurnId.make("turn-from-original-send");
const sent = yield* Deferred.make<ProviderSendTurnInput>();

yield* Effect.gen(function* () {
const engine = yield* OrchestrationEngine.OrchestrationEngineService;
const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory;
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make("create-restart-project"),
projectId,
title: "Restart continuation",
workspaceRoot: "/tmp/startup-orphan-project",
defaultModelSelection: { instanceId: providerInstanceId, model: "gpt-5" },
createdAt,
});
yield* engine.dispatch({
type: "thread.create",
commandId: CommandId.make("create-restart-thread"),
threadId,
projectId,
title: "Newer running turn",
modelSelection: { instanceId: providerInstanceId, model: "gpt-5" },
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt,
});
yield* engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("persist-newer-running-turn"),
threadId,
session: {
threadId,
status: "running",
providerName: "codex",
providerInstanceId,
runtimeMode: "full-access",
activeTurnId,
lastError: null,
updatedAt: createdAt,
},
createdAt,
});
yield* directory.upsert({
threadId,
provider: ProviderDriverKind.make("codex"),
providerInstanceId,
status: "running",
resumeCursor,
runtimePayload: { activeTurnId: originalTurnId },
});
if (restart === "marked remote update") {
assert.deepStrictEqual(
yield* ServerRuntimeStartup.markRunningProviderSessionsForContinuation,
[threadId],
);
}
}).pipe(Effect.provide(makePersistedRuntimeLayer(config.dbPath)));

yield* Effect.gen(function* () {
const query = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery;
const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory;
const provider = yield* ProviderService.ProviderService;
const before = Option.getOrThrow(yield* query.getThreadDetailById(threadId));
assert.equal(before.session?.activeTurnId, activeTurnId);
assert.propertyVal(
Option.getOrThrow(yield* directory.getBinding(threadId)).runtimePayload,
"activeTurnId",
originalTurnId,
);
yield* ServerRuntimeStartup.reconcileProviderSessions.pipe(
Effect.provideService(ProviderService.ProviderService, {
...provider,
getCapabilities: () =>
Effect.succeed({
sessionModelSwitch: "in-session",
promptlessTurnContinuation: true,
}),
sendTurn: (input) =>
Deferred.succeed(sent, input).pipe(
Effect.as({ threadId, turnId: TurnId.make("continued-turn") }),
),
}),
Effect.provide(
ServerSettings.layerTest({
continueThreadsAfterServerUpdate: restart === "opt-in desktop restart",
}),
),
);
const after = Option.getOrThrow(yield* query.getThreadDetailById(threadId));
assert.equal(after.session?.status, "starting");
assert.equal(after.session?.activeTurnId, null);
assert.equal(after.session?.lastError, null);
assert.deepStrictEqual(yield* Deferred.await(sent), {
threadId,
continuation: true,
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
});
}).pipe(
Effect.provide(
Layer.mergeAll(makePersistedRuntimeLayer(config.dbPath), startupDependencies),
),
);
}).pipe(
Effect.provide(
ServerConfig.layerTest(process.cwd(), { prefix: "t3-restart-newer-turn-" }).pipe(
Layer.provideMerge(NodeServices.layer),
),
),
),
);
130 changes: 111 additions & 19 deletions apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -781,25 +781,123 @@ describe("isRecoverableThreadResumeError", () => {
});

describe("openCodexThread", () => {
it.effect("resumes metadata when historical turns contain unknown error values", () =>
Effect.gen(function* () {
const response = makeThreadOpenResponse("saved-thread");
const calls: unknown[] = [];
const opened = yield* openCodexThread({
client: {
request: () => Effect.die("A valid resumed thread must not start fresh"),
raw: {
request: (method, payload) => {
calls.push({ method, payload });
return Effect.succeed({
...response,
thread: {
...response.thread,
turns: [
{
id: "old-turn",
status: "failed",
items: [],
error: {
message: "Historical provider error",
codexErrorInfo: "misalignment_policy_violation",
},
},
],
},
});
},
},
},
threadId: ThreadId.make("thread-1"),
runtimeMode: "auto",
cwd: "/tmp/project",
requestedModel: "gpt-5.3-codex",
serviceTier: "fast",
resumeThreadId: "saved-thread",
});

NodeAssert.deepStrictEqual(opened, {
cwd: response.cwd,
model: response.model,
thread: { id: "saved-thread" },
});
NodeAssert.deepStrictEqual(calls, [
{
method: "thread/resume",
payload: {
threadId: "saved-thread",
cwd: "/tmp/project",
model: "gpt-5.3-codex",
serviceTier: "fast",
approvalPolicy: "on-request",
sandbox: "workspace-write",
approvalsReviewer: "auto_review",
excludeTurns: true,
},
},
]);
}),
);

it.effect("rejects malformed required resume metadata without starting a fresh thread", () =>
Effect.gen(function* () {
for (const invalidMetadata of [
{ cwd: null },
{ model: 42 },
{ thread: { id: null } },
{ thread: {} },
]) {
const error = yield* openCodexThread({
client: {
request: () => Effect.die("Invalid resume metadata must not start a fresh thread"),
raw: {
request: () =>
Effect.succeed({ ...makeThreadOpenResponse("saved-thread"), ...invalidMetadata }),
},
},
threadId: ThreadId.make("thread-1"),
runtimeMode: "full-access",
cwd: "/tmp/project",
requestedModel: "gpt-5.3-codex",
serviceTier: undefined,
resumeThreadId: "saved-thread",
}).pipe(Effect.flip);

NodeAssert.ok(isCodexAppServerRequestError(error));
NodeAssert.equal(error.operation, "decode-payload");
NodeAssert.equal(error.method, "thread/resume");
}
}),
);

it.effect("falls back to thread/start when resume fails recoverably", () =>
Effect.gen(function* () {
const calls: Array<{ method: "thread/start" | "thread/resume"; payload: unknown }> = [];
const started = makeThreadOpenResponse("fresh-thread");
const client = {
request: <M extends "thread/start" | "thread/resume">(
method: M,
payload: CodexRpc.ClientRequestParamsByMethod[M],
) => {
calls.push({ method, payload });
if (method === "thread/resume") {
raw: {
request: (
method: "thread/resume",
payload: CodexRpc.ClientRequestParamsByMethod["thread/resume"],
) => {
calls.push({ method, payload });
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "thread not found",
}),
);
}
return Effect.succeed(started as CodexRpc.ClientRequestResponsesByMethod[M]);
},
},
request: (
method: "thread/start",
payload: CodexRpc.ClientRequestParamsByMethod["thread/start"],
) => {
calls.push({ method, payload });
return Effect.succeed(started);
},
};

Expand All @@ -824,21 +922,15 @@ describe("openCodexThread", () => {
it.effect("propagates non-recoverable resume failures", () =>
Effect.gen(function* () {
const client = {
request: <M extends "thread/start" | "thread/resume">(
method: M,
_payload: CodexRpc.ClientRequestParamsByMethod[M],
) => {
if (method === "thread/resume") {
return Effect.fail(
request: () => Effect.die("Non-recoverable resume failures must not start a fresh thread"),
raw: {
request: () =>
Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "timed out waiting for server",
}),
);
}
return Effect.succeed(
makeThreadOpenResponse("fresh-thread") as CodexRpc.ClientRequestResponsesByMethod[M],
);
),
},
};

Expand Down
49 changes: 38 additions & 11 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -676,17 +676,29 @@ export function isRecoverableThreadResumeError(error: unknown): boolean {
return RECOVERABLE_THREAD_RESUME_ERROR_SNIPPETS.some((snippet) => message.includes(snippet));
}

type CodexThreadOpenResponse =
| CodexRpc.ClientRequestResponsesByMethod["thread/start"]
| CodexRpc.ClientRequestResponsesByMethod["thread/resume"];

type CodexThreadOpenMethod = "thread/start" | "thread/resume";
const CodexThreadResumeMetadata = Schema.Struct({
cwd: Schema.String,
model: Schema.String,
thread: Schema.Struct({ id: Schema.String }),
});
const decodeCodexThreadResumeMetadata = Schema.decodeUnknownEffect(CodexThreadResumeMetadata);

interface CodexThreadOpenClient {
readonly request: <M extends CodexThreadOpenMethod>(
method: M,
payload: CodexRpc.ClientRequestParamsByMethod[M],
) => Effect.Effect<CodexRpc.ClientRequestResponsesByMethod[M], CodexErrors.CodexAppServerError>;
readonly raw: {
readonly request: (
method: "thread/resume",
payload: CodexRpc.ClientRequestParamsByMethod["thread/resume"] & {
readonly excludeTurns?: boolean;
},
) => Effect.Effect<unknown, CodexErrors.CodexAppServerError>;
};
readonly request: (
method: "thread/start",
payload: CodexRpc.ClientRequestParamsByMethod["thread/start"],
) => Effect.Effect<
CodexRpc.ClientRequestResponsesByMethod["thread/start"],
CodexErrors.CodexAppServerError
>;
}

export const openCodexThread = (input: {
Expand All @@ -697,7 +709,7 @@ export const openCodexThread = (input: {
readonly requestedModel: string | undefined;
readonly serviceTier: CodexServiceTier | undefined;
readonly resumeThreadId: string | undefined;
}): Effect.Effect<CodexThreadOpenResponse, CodexErrors.CodexAppServerError> => {
}): Effect.Effect<typeof CodexThreadResumeMetadata.Type, CodexErrors.CodexAppServerError> => {
const resumeThreadId = input.resumeThreadId;
const startParams = buildThreadStartParams({
cwd: input.cwd,
Expand All @@ -710,12 +722,27 @@ export const openCodexThread = (input: {
return input.client.request("thread/start", startParams);
}

return input.client
// Older providers may still return history despite excludeTurns. Only the
// session metadata is needed here, so unrelated historical items cannot
// prevent resuming a valid provider thread.
return input.client.raw
.request("thread/resume", {
threadId: resumeThreadId,
...startParams,
excludeTurns: true,
})
.pipe(
Effect.flatMap((response) =>
decodeCodexThreadResumeMetadata(response).pipe(
Effect.mapError((error) =>
CodexErrors.CodexAppServerRequestError.invalidPayload(
"thread/resume",
"decode-payload",
error,
),
),
),
),
Effect.catchIf(isRecoverableThreadResumeError, (error) =>
Effect.logWarning("codex app-server thread resume fell back to fresh start", {
threadId: input.threadId,
Expand Down
Loading
Loading