Skip to content
174 changes: 108 additions & 66 deletions apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
import * as NodeAssert from "node:assert/strict";

import { it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Schema from "effect/Schema";
import { describe, it } from "vite-plus/test";
import { describe } from "vite-plus/test";
import { ThreadId } from "@t3tools/contracts";
import * as CodexErrors from "effect-codex-app-server/errors";
import * as CodexRpc from "effect-codex-app-server/rpc";
Expand All @@ -19,6 +20,23 @@ import {
} from "./CodexSessionRuntime.ts";
const isCodexAppServerRequestError = Schema.is(CodexErrors.CodexAppServerRequestError);

describe("CodexSessionRuntimeIdentifierGenerationError", () => {
it("retains identifier purpose and the random source failure", () => {
const cause = new Error("random source unavailable");
const error = new CodexErrors.CodexAppServerIdentifierGenerationError({
purpose: "provider-event",
cause,
});

NodeAssert.equal(error.purpose, "provider-event");
NodeAssert.strictEqual(error.cause, cause);
NodeAssert.equal(
error.message,
"Failed to generate Codex App Server identifier for provider-event.",
);
});
});

function makeThreadOpenResponse(
threadId: string,
): CodexRpc.ClientRequestResponsesByMethod["thread/start"] {
Expand All @@ -43,6 +61,32 @@ function makeThreadOpenResponse(
}

describe("buildTurnStartParams", () => {
it("keeps invalid turn values only in the schema cause", () => {
const secret = "codex-turn-input-secret-sentinel";
const error = Effect.runSync(
buildTurnStartParams({
threadId: "provider-thread-1",
runtimeMode: "full-access",
attachments: [
{
type: "image",
url: { secret } as unknown as string,
},
],
}).pipe(Effect.flip),
);
const { cause, ...directDiagnostics } = error;

NodeAssert.equal(error.operation, "decode-request-payload");
NodeAssert.equal(error.method, "turn/start");
NodeAssert.ok((error.issueCount ?? 0) > 0);
NodeAssert.ok(error.issueKinds?.includes("Pointer"));
NodeAssert.ok((error.maximumPathDepth ?? 0) > 0);
NodeAssert.ok(Schema.isSchemaError(cause));
NodeAssert.doesNotMatch(error.message, new RegExp(secret));
NodeAssert.doesNotMatch(JSON.stringify(directDiagnostics), new RegExp(secret));
});

it("includes plan collaboration mode when requested", () => {
const params = Effect.runSync(
buildTurnStartParams({
Expand Down Expand Up @@ -223,81 +267,79 @@ describe("isRecoverableThreadResumeError", () => {
});

describe("openCodexThread", () => {
it("falls back to thread/start when resume fails recoverably", async () => {
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") {
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "thread not found",
}),
);
}
return Effect.succeed(started as CodexRpc.ClientRequestResponsesByMethod[M]);
},
};
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") {
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "thread not found",
}),
);
}
return Effect.succeed(started as CodexRpc.ClientRequestResponsesByMethod[M]);
},
};

const opened = await Effect.runPromise(
openCodexThread({
const opened = yield* openCodexThread({
client,
threadId: ThreadId.make("thread-1"),
runtimeMode: "full-access",
cwd: "/tmp/project",
requestedModel: "gpt-5.3-codex",
serviceTier: undefined,
resumeThreadId: "stale-thread",
}),
);
});

NodeAssert.equal(opened.thread.id, "fresh-thread");
NodeAssert.deepStrictEqual(
calls.map((call) => call.method),
["thread/resume", "thread/start"],
);
});
NodeAssert.equal(opened.thread.id, "fresh-thread");
NodeAssert.deepStrictEqual(
calls.map((call) => call.method),
["thread/resume", "thread/start"],
);
}),
);

it("propagates non-recoverable resume failures", async () => {
const client = {
request: <M extends "thread/start" | "thread/resume">(
method: M,
_payload: CodexRpc.ClientRequestParamsByMethod[M],
) => {
if (method === "thread/resume") {
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "timed out waiting for server",
}),
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(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "timed out waiting for server",
}),
);
}
return Effect.succeed(
makeThreadOpenResponse("fresh-thread") as CodexRpc.ClientRequestResponsesByMethod[M],
);
}
return Effect.succeed(
makeThreadOpenResponse("fresh-thread") as CodexRpc.ClientRequestResponsesByMethod[M],
);
},
};
},
};

await NodeAssert.rejects(
Effect.runPromise(
openCodexThread({
client,
threadId: ThreadId.make("thread-1"),
runtimeMode: "full-access",
cwd: "/tmp/project",
requestedModel: "gpt-5.3-codex",
serviceTier: undefined,
resumeThreadId: "stale-thread",
}),
),
(error: unknown) =>
isCodexAppServerRequestError(error) &&
error.errorMessage === "timed out waiting for server",
);
});
const error = yield* openCodexThread({
client,
threadId: ThreadId.make("thread-1"),
runtimeMode: "full-access",
cwd: "/tmp/project",
requestedModel: "gpt-5.3-codex",
serviceTier: undefined,
resumeThreadId: "stale-thread",
}).pipe(Effect.flip);

NodeAssert.ok(isCodexAppServerRequestError(error));
NodeAssert.equal(error.errorMessage, "timed out waiting for server");
}),
);
});
59 changes: 30 additions & 29 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,10 +26,9 @@ import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
import * as Scope from "effect/Scope";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import * as SchemaIssue from "effect/SchemaIssue";
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process";
import * as CodexClient from "effect-codex-app-server/client";
import * as CodexErrors from "effect-codex-app-server/errors";
Expand Down Expand Up @@ -89,7 +88,6 @@ const decodeCodexTurnStartParamsWithCollaborationMode = Schema.decodeUnknownEffe

export type CodexTurnStartParamsWithCollaborationMode =
typeof CodexTurnStartParamsWithCollaborationMode.Type;
const formatSchemaIssue = SchemaIssue.makeFormatterDefault();

export type CodexResumeCursor = typeof CodexResumeCursorSchema.Type;
type CodexServiceTier = NonNullable<EffectCodexSchema.V2ThreadStartParams["serviceTier"]>;
Expand Down Expand Up @@ -390,7 +388,13 @@ export function buildTurnStartParams(input: {
...(input.effort ? { effort: input.effort } : {}),
...(collaborationMode ? { collaborationMode } : {}),
}).pipe(
Effect.mapError((error) => toProtocolParseError("Invalid turn/start request payload", error)),
Effect.mapError((cause) =>
CodexErrors.CodexAppServerProtocolParseError.fromSchemaError(
"decode-request-payload",
cause,
{ method: "turn/start" },
),
),
);
}

Expand Down Expand Up @@ -468,7 +472,7 @@ export const openCodexThread = (input: {
requestedRuntimeMode: input.runtimeMode,
resumeThreadId,
recoverable: true,
cause: error.message,
cause: error,
}).pipe(Effect.andThen(input.client.request("thread/start", startParams))),
),
);
Expand Down Expand Up @@ -658,16 +662,6 @@ function toCodexUserInputAnswers(
).pipe(Effect.map((entries) => Object.fromEntries(entries)));
}

function toProtocolParseError(
detail: string,
cause: Schema.SchemaError,
): CodexErrors.CodexAppServerProtocolParseError {
return new CodexErrors.CodexAppServerProtocolParseError({
detail: `${detail}: ${formatSchemaIssue(cause.issue)}`,
cause,
});
}

function currentProviderThreadId(session: ProviderSession): string | undefined {
return readResumeCursorThreadId(session.resumeCursor);
}
Expand Down Expand Up @@ -760,15 +754,16 @@ export const makeCodexSessionRuntime = (
);
const serverNotifications = yield* Queue.unbounded<CodexServerNotification>();
const nowIso = Effect.map(DateTime.now, DateTime.formatIso);
const randomUUIDv4 = crypto.randomUUIDv4.pipe(
Effect.mapError(
(cause) =>
new CodexErrors.CodexAppServerTransportError({
detail: "Failed to generate Codex runtime identifier.",
cause,
}),
),
);
const randomUUIDv4 = (purpose: CodexErrors.CodexAppServerIdentifierPurpose) =>
crypto.randomUUIDv4.pipe(
Effect.mapError(
(cause) =>
new CodexErrors.CodexAppServerIdentifierGenerationError({
purpose,
cause,
}),
),
);

const sessionCreatedAt = yield* nowIso;
const initialSession = {
Expand All @@ -788,7 +783,7 @@ export const makeCodexSessionRuntime = (

const emitEvent = (event: Omit<ProviderEvent, "id" | "provider" | "createdAt">) =>
Effect.gen(function* () {
const id = yield* randomUUIDv4;
const id = yield* randomUUIDv4("provider-event");
return yield* offerEvent({
id: EventId.make(id),
provider: PROVIDER,
Expand Down Expand Up @@ -956,7 +951,7 @@ export const makeCodexSessionRuntime = (

yield* client.handleServerRequest("item/commandExecution/requestApproval", (payload) =>
Effect.gen(function* () {
const requestId = ApprovalRequestId.make(yield* randomUUIDv4);
const requestId = ApprovalRequestId.make(yield* randomUUIDv4("command-approval-request"));
const turnId = TurnId.make(payload.turnId);
const itemId = ProviderItemId.make(payload.itemId);
const decision = yield* Deferred.make<ProviderApprovalDecision>();
Expand Down Expand Up @@ -1012,7 +1007,9 @@ export const makeCodexSessionRuntime = (

yield* client.handleServerRequest("item/fileChange/requestApproval", (payload) =>
Effect.gen(function* () {
const requestId = ApprovalRequestId.make(yield* randomUUIDv4);
const requestId = ApprovalRequestId.make(
yield* randomUUIDv4("file-change-approval-request"),
);
const turnId = TurnId.make(payload.turnId);
const itemId = ProviderItemId.make(payload.itemId);
const decision = yield* Deferred.make<ProviderApprovalDecision>();
Expand Down Expand Up @@ -1068,7 +1065,7 @@ export const makeCodexSessionRuntime = (

yield* client.handleServerRequest("item/tool/requestUserInput", (payload) =>
Effect.gen(function* () {
const requestId = ApprovalRequestId.make(yield* randomUUIDv4);
const requestId = ApprovalRequestId.make(yield* randomUUIDv4("user-input-request"));
const turnId = TurnId.make(payload.turnId);
const itemId = ProviderItemId.make(payload.itemId);
const answers = yield* Deferred.make<ProviderUserInputAnswers>();
Expand Down Expand Up @@ -1293,7 +1290,11 @@ export const makeCodexSessionRuntime = (
const rawResponse = yield* client.raw.request("turn/start", params);
const response = yield* decodeV2TurnStartResponse(rawResponse).pipe(
Effect.mapError((error) =>
toProtocolParseError("Invalid turn/start response payload", error),
CodexErrors.CodexAppServerProtocolParseError.fromSchemaError(
"decode-response-payload",
error,
{ method: "turn/start" },
),
),
);
const turnId = TurnId.make(response.turn.id);
Expand Down
Loading
Loading