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
97 changes: 94 additions & 3 deletions apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -781,6 +781,97 @@ describe("isRecoverableThreadResumeError", () => {
});

describe("openCodexThread", () => {
it.effect("unarchives an archived session and resumes the same thread", () =>
Effect.gen(function* () {
const calls: Array<{ method: string; payload: unknown }> = [];
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const opened = yield* openCodexThread({
client: {
request: () => Effect.die("An archived session must not start fresh"),
raw: {
request: (method, payload) =>
Effect.suspend(() => {
calls.push({ method, payload });
if (calls.length === 1) {
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32600,
errorMessage:
"session saved-thread is archived. Run `codex unarchive saved-thread` to unarchive it first.",
}),
);
}
return Effect.succeed(makeThreadOpenResponse("saved-thread"));
}),
},
},
threadId: ThreadId.make("thread-1"),
runtimeMode: "auto",
cwd: "/tmp/project",
requestedModel: "gpt-5.3-codex",
serviceTier: "fast",
resumeThreadId: "saved-thread",
});

NodeAssert.equal(opened.thread.id, "saved-thread");
NodeAssert.deepStrictEqual(
calls.map((call) => call.method),
["thread/resume", "thread/unarchive", "thread/resume"],
);
NodeAssert.deepStrictEqual(calls[1]?.payload, { threadId: "saved-thread" });
NodeAssert.deepStrictEqual(calls[2]?.payload, calls[0]?.payload);
}),
);

it.effect(
"propagates unarchive and retry failures without starting fresh or retrying again",
() =>
Effect.gen(function* () {
for (const failAt of [2, 3]) {
for (const errorMessage of ["thread not found", "session saved-thread is archived"]) {
const failure = new CodexErrors.CodexAppServerRequestError({
code: -32600,
errorMessage,
});
const calls: string[] = [];
const error = yield* openCodexThread({
client: {
request: () => Effect.die("An archived session must not start fresh"),
raw: {
request: (method) =>
Effect.suspend(() => {
calls.push(method);
if (calls.length === failAt) return Effect.fail(failure);
if (calls.length === 1) {
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32600,
errorMessage:
"Run `codex unarchive saved-thread` to unarchive it first.",
}),
);
}
return Effect.succeed({});
}),
},
},
threadId: ThreadId.make("thread-1"),
runtimeMode: "full-access",
cwd: "/tmp/project",
requestedModel: undefined,
serviceTier: undefined,
resumeThreadId: "saved-thread",
}).pipe(Effect.flip);

NodeAssert.strictEqual(error, failure);
NodeAssert.deepStrictEqual(
calls,
["thread/resume", "thread/unarchive", "thread/resume"].slice(0, failAt),
);
}
}
}),
);

it.effect("resumes metadata when historical turns contain unknown error values", () =>
Effect.gen(function* () {
const response = makeThreadOpenResponse("saved-thread");
Expand Down Expand Up @@ -875,13 +966,13 @@ describe("openCodexThread", () => {

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 calls: Array<{ method: string; payload: unknown }> = [];
const started = makeThreadOpenResponse("fresh-thread");
const client = {
raw: {
request: (
method: "thread/resume",
payload: CodexRpc.ClientRequestParamsByMethod["thread/resume"],
method: "thread/resume" | "thread/unarchive",
payload: CodexRpc.ClientRequestParamsByMethod["thread/resume" | "thread/unarchive"],
) => {
calls.push({ method, payload });
return Effect.fail(
Expand Down
27 changes: 19 additions & 8 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -686,8 +686,8 @@ const decodeCodexThreadResumeMetadata = Schema.decodeUnknownEffect(CodexThreadRe
interface CodexThreadOpenClient {
readonly raw: {
readonly request: (
method: "thread/resume",
payload: CodexRpc.ClientRequestParamsByMethod["thread/resume"] & {
method: "thread/resume" | "thread/unarchive",
payload: CodexRpc.ClientRequestParamsByMethod["thread/resume" | "thread/unarchive"] & {
readonly excludeTurns?: boolean;
},
) => Effect.Effect<unknown, CodexErrors.CodexAppServerError>;
Expand Down Expand Up @@ -725,7 +725,7 @@ export const openCodexThread = (input: {
// 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
const resume = input.client.raw
.request("thread/resume", {
threadId: resumeThreadId,
...startParams,
Expand All @@ -743,16 +743,27 @@ export const openCodexThread = (input: {
),
),
),
Effect.catchIf(isRecoverableThreadResumeError, (error) =>
Effect.logWarning("codex app-server thread resume fell back to fresh start", {
);

return resume.pipe(
Effect.catch((error) => {
if (/is archived|codex unarchive/i.test(error.message)) {
return input.client.raw
.request("thread/unarchive", { threadId: resumeThreadId })
.pipe(Effect.andThen(resume));
}
if (isRecoverableThreadResumeError(error)) {
return Effect.logWarning("codex app-server thread resume fell back to fresh start", {
threadId: input.threadId,
requestedRuntimeMode: input.runtimeMode,
resumeThreadId,
recoverable: true,
cause: error,
}).pipe(Effect.andThen(input.client.request("thread/start", startParams))),
),
);
}).pipe(Effect.andThen(input.client.request("thread/start", startParams)));
}
return Effect.fail(error);
}),
);
};

function readNotificationThreadId(notification: CodexServerNotification): string | undefined {
Expand Down
Loading