Skip to content
Draft
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
49 changes: 48 additions & 1 deletion apps/server/src/provider/Layers/CodexProvider.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,53 @@
import { assert, it } from "@effect/vitest";
import * as Effect from "effect/Effect";

import { applyPreferredCodexDefaultModel, mapCodexModelCapabilities } from "./CodexProvider.ts";
import {
applyPreferredCodexDefaultModel,
mapCodexModelCapabilities,
refreshCodexSkillsAfterPluginSync,
} from "./CodexProvider.ts";

it.effect("waits for plugin inventory before force-reloading skills", () =>
Effect.gen(function* () {
const calls: Array<{ method: string; payload: unknown }> = [];
const skillsResponse = {
data: [{ cwd: "/tmp/project", errors: [], skills: [] }],
};
const client = {
request: (method: string, payload: unknown) => {
calls.push({ method, payload });
switch (method) {
case "plugin/installed":
return Effect.succeed({ marketplaces: [] });
case "skills/list":
return Effect.succeed(skillsResponse);
default:
return Effect.die(new Error(`Unexpected request: ${method}`));
}
},
} as unknown as Parameters<typeof refreshCodexSkillsAfterPluginSync>[0]["client"];

const response = yield* refreshCodexSkillsAfterPluginSync({
client,
cwd: "/tmp/project",
});

assert.strictEqual(response, skillsResponse);
assert.deepStrictEqual(calls, [
{
method: "plugin/installed",
payload: { cwds: ["/tmp/project"] },
},
{
method: "skills/list",
payload: {
cwds: ["/tmp/project"],
forceReload: true,
},
},
]);
}),
);

it("maps current Codex model capability fields", () => {
const capabilities = mapCodexModelCapabilities({
Expand Down
33 changes: 31 additions & 2 deletions apps/server/src/provider/Layers/CodexProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,34 @@ export interface CodexAppServerProviderSnapshot {
readonly skills: ReadonlyArray<ServerProviderSkill>;
}

export interface CodexSkillDiscoveryClient {
readonly request: {
(
method: "plugin/installed",
payload: CodexSchema.V2PluginInstalledParams,
): Effect.Effect<CodexSchema.V2PluginInstalledResponse, CodexErrors.CodexAppServerError>;
(
method: "skills/list",
payload: CodexSchema.V2SkillsListParams,
): Effect.Effect<CodexSchema.V2SkillsListResponse, CodexErrors.CodexAppServerError>;
};
}

export const refreshCodexSkillsAfterPluginSync = Effect.fn("refreshCodexSkillsAfterPluginSync")(
function* (input: { readonly client: CodexSkillDiscoveryClient; readonly cwd: string }) {
// Codex materializes remote plugins asynchronously after initialization.
// Reading the installed-plugin inventory is the readiness barrier for the
// subsequent forced disk scan; that transition emits no `skills/changed`.
yield* input.client.request("plugin/installed", {
cwds: [input.cwd],
});
return yield* input.client.request("skills/list", {
cwds: [input.cwd],
forceReload: true,
});
},
);

const REASONING_EFFORT_LABELS: Readonly<Record<string, string>> = {
none: "None",
minimal: "Minimal",
Expand Down Expand Up @@ -391,8 +419,9 @@ const probeCodexAppServerProvider = Effect.fn("probeCodexAppServerProvider")(fun

const [skillsResponse, models] = yield* Effect.all(
[
client.request("skills/list", {
cwds: [input.cwd],
refreshCodexSkillsAfterPluginSync({
client,
cwd: input.cwd,
}),
requestAllCodexModels(client),
],
Expand Down
117 changes: 87 additions & 30 deletions apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -393,27 +393,79 @@ describe("isRecoverableThreadResumeError", () => {
});

describe("openCodexThread", () => {
it.effect("refreshes skills after plugin sync before starting a thread", () =>
Effect.gen(function* () {
const calls: Array<{ method: string; payload: unknown }> = [];
const started = makeThreadOpenResponse("fresh-thread");
const client = {
request: (method: string, payload: unknown) => {
calls.push({ method, payload });
switch (method) {
case "plugin/installed":
return Effect.succeed({ marketplaces: [] });
case "skills/list":
return Effect.succeed({
data: [{ cwd: "/tmp/project", errors: [], skills: [] }],
});
case "thread/start":
return Effect.succeed(started);
default:
return Effect.die(new Error(`Unexpected request: ${method}`));
}
},
} as unknown as Parameters<typeof openCodexThread>[0]["client"];

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

NodeAssert.equal(opened.thread.id, "fresh-thread");
NodeAssert.deepStrictEqual(
calls.map((call) => call.method),
["plugin/installed", "skills/list", "thread/start"],
);
NodeAssert.deepStrictEqual(calls[0]?.payload, { cwds: ["/tmp/project"] });
NodeAssert.deepStrictEqual(calls[1]?.payload, {
cwds: ["/tmp/project"],
forceReload: true,
});
}),
);

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 = {
request: <M extends "thread/start" | "thread/resume">(
method: M,
payload: CodexRpc.ClientRequestParamsByMethod[M],
) => {
request: (method: string, payload: unknown) => {
calls.push({ method, payload });
if (method === "thread/resume") {
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "thread not found",
}),
);
switch (method) {
case "plugin/installed":
return Effect.succeed({ marketplaces: [] });
case "skills/list":
return Effect.succeed({
data: [{ cwd: "/tmp/project", errors: [], skills: [] }],
});
case "thread/resume":
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "thread not found",
}),
);
case "thread/start":
return Effect.succeed(started);
default:
return Effect.die(new Error(`Unexpected request: ${method}`));
}
return Effect.succeed(started as CodexRpc.ClientRequestResponsesByMethod[M]);
},
};
} as unknown as Parameters<typeof openCodexThread>[0]["client"];

const opened = yield* openCodexThread({
client,
Expand All @@ -428,31 +480,36 @@ describe("openCodexThread", () => {
NodeAssert.equal(opened.thread.id, "fresh-thread");
NodeAssert.deepStrictEqual(
calls.map((call) => call.method),
["thread/resume", "thread/start"],
["plugin/installed", "skills/list", "thread/resume", "thread/start"],
);
}),
);

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",
}),
);
request: (method: string) => {
switch (method) {
case "plugin/installed":
return Effect.succeed({ marketplaces: [] });
case "skills/list":
return Effect.succeed({
data: [{ cwd: "/tmp/project", errors: [], skills: [] }],
});
case "thread/resume":
return Effect.fail(
new CodexErrors.CodexAppServerRequestError({
code: -32603,
errorMessage: "timed out waiting for server",
}),
);
case "thread/start":
return Effect.succeed(makeThreadOpenResponse("fresh-thread"));
default:
return Effect.die(new Error(`Unexpected request: ${method}`));
}
return Effect.succeed(
makeThreadOpenResponse("fresh-thread") as CodexRpc.ClientRequestResponsesByMethod[M],
);
},
};
} as unknown as Parameters<typeof openCodexThread>[0]["client"];

const error = yield* openCodexThread({
client,
Expand Down
23 changes: 16 additions & 7 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,11 @@ import * as CodexErrors from "effect-codex-app-server/errors";
import * as CodexRpc from "effect-codex-app-server/rpc";
import * as EffectCodexSchema from "effect-codex-app-server/schema";

import { buildCodexInitializeParams } from "./CodexProvider.ts";
import {
buildCodexInitializeParams,
refreshCodexSkillsAfterPluginSync,
type CodexSkillDiscoveryClient,
} from "./CodexProvider.ts";
import { codexSessionAppServerArgs } from "./codexLaunchArgs.ts";
import { expandHomePath } from "../../pathExpansion.ts";
import { buildCodexDeveloperInstructions } from "../CodexDeveloperInstructions.ts";
Expand Down Expand Up @@ -447,22 +451,27 @@ type CodexThreadOpenResponse =

type CodexThreadOpenMethod = "thread/start" | "thread/resume";

interface CodexThreadOpenClient {
interface CodexThreadOpenClient extends CodexSkillDiscoveryClient {
readonly request: <M extends CodexThreadOpenMethod>(
method: M,
payload: CodexRpc.ClientRequestParamsByMethod[M],
) => Effect.Effect<CodexRpc.ClientRequestResponsesByMethod[M], CodexErrors.CodexAppServerError>;
}

export const openCodexThread = (input: {
export const openCodexThread = Effect.fn("openCodexThread")(function* (input: {
readonly client: CodexThreadOpenClient;
readonly threadId: ThreadId;
readonly runtimeMode: RuntimeMode;
readonly cwd: string;
readonly requestedModel: string | undefined;
readonly serviceTier: CodexServiceTier | undefined;
readonly resumeThreadId: string | undefined;
}): Effect.Effect<CodexThreadOpenResponse, CodexErrors.CodexAppServerError> => {
}): Effect.fn.Return<CodexThreadOpenResponse, CodexErrors.CodexAppServerError> {
yield* refreshCodexSkillsAfterPluginSync({
client: input.client,
cwd: input.cwd,
});

const resumeThreadId = input.resumeThreadId;
const startParams = buildThreadStartParams({
cwd: input.cwd,
Expand All @@ -472,10 +481,10 @@ export const openCodexThread = (input: {
});

if (resumeThreadId === undefined) {
return input.client.request("thread/start", startParams);
return yield* input.client.request("thread/start", startParams);
}

return input.client
return yield* input.client
.request("thread/resume", {
threadId: resumeThreadId,
...startParams,
Expand All @@ -491,7 +500,7 @@ export const openCodexThread = (input: {
}).pipe(Effect.andThen(input.client.request("thread/start", startParams))),
),
);
};
});

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