Skip to content

Commit a8b8e5a

Browse files
committed
fix(workflow): make native schema fallback session-aware
1 parent 7d797e3 commit a8b8e5a

8 files changed

Lines changed: 313 additions & 72 deletions

src/local-agent-adapters.ts

Lines changed: 49 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,8 @@ import type { LocalAgentProvider } from "./local-agent-profiles.js";
66
import { removeDevspaceNodeModulesBinFromPath } from "./local-agent-path.js";
77
import {
88
createCodexSdkLocalAgentRuntime,
9+
isNativeSchemaUnsupportedFailure,
10+
ProviderSchemaUnsupportedError,
911
type LocalAgentRunInput,
1012
type LocalAgentRunResult,
1113
} from "./local-agent-runtime.js";
@@ -59,47 +61,56 @@ class ClaudeLocalAgentAdapter implements LocalAgentAdapter {
5961
async run(input: LocalAgentRunInput): Promise<LocalAgentRunResult> {
6062
const { query } = await import("@anthropic-ai/claude-agent-sdk");
6163
const claudeExecutable = process.env.CLAUDE_COMMAND ?? resolveExecutable("claude");
62-
const messages = query({
63-
prompt: input.prompt,
64-
options: {
65-
cwd: input.workspace,
66-
model: input.model,
67-
...(input.effort ? { thinking: { type: "adaptive" } as const, effort: input.effort as EffortLevel } : {}),
68-
resume: input.providerSessionId,
69-
permissionMode: "bypassPermissions",
70-
allowDangerouslySkipPermissions: true,
71-
env: claudeCommandEnvironment(process.env),
72-
...(claudeExecutable ? { pathToClaudeCodeExecutable: claudeExecutable } : {}),
73-
...claudeOutputFormatOptions(input.schema),
74-
},
75-
});
64+
try {
65+
const messages = query({
66+
prompt: input.prompt,
67+
options: {
68+
cwd: input.workspace,
69+
model: input.model,
70+
...(input.effort
71+
? { thinking: { type: "adaptive" } as const, effort: input.effort as EffortLevel }
72+
: {}),
73+
resume: input.providerSessionId,
74+
permissionMode: "bypassPermissions",
75+
allowDangerouslySkipPermissions: true,
76+
env: claudeCommandEnvironment(process.env),
77+
...(claudeExecutable ? { pathToClaudeCodeExecutable: claudeExecutable } : {}),
78+
...claudeOutputFormatOptions(input.schema),
79+
},
80+
});
7681

77-
let providerSessionId = input.providerSessionId ?? null;
78-
let finalResponse = "";
79-
let structured: unknown | undefined;
80-
const items: unknown[] = [];
81-
for await (const message of messages) {
82-
items.push(message);
83-
const record = message as Record<string, unknown>;
84-
if (typeof record.session_id === "string") providerSessionId = record.session_id;
85-
if (record.type !== "result") continue;
86-
const resultError = claudeResultError(record);
87-
if (resultError) throw new Error(resultError);
88-
const extracted = extractClaudeResultPayload(record);
89-
if (extracted) {
90-
finalResponse = extracted.finalResponse;
91-
structured = extracted.structured;
82+
let providerSessionId = input.providerSessionId ?? null;
83+
let finalResponse = "";
84+
let structured: unknown | undefined;
85+
const items: unknown[] = [];
86+
for await (const message of messages) {
87+
items.push(message);
88+
const record = message as Record<string, unknown>;
89+
if (typeof record.session_id === "string") providerSessionId = record.session_id;
90+
if (record.type !== "result") continue;
91+
const resultError = claudeResultError(record);
92+
if (resultError) throw new Error(resultError);
93+
const extracted = extractClaudeResultPayload(record);
94+
if (extracted) {
95+
finalResponse = extracted.finalResponse;
96+
structured = extracted.structured;
97+
}
9298
}
93-
}
9499

95-
finalResponse = requireFinalResponse("Claude", finalResponse);
96-
return {
97-
provider: this.provider,
98-
providerSessionId,
99-
finalResponse,
100-
items,
101-
...(structured !== undefined ? { structured } : {}),
102-
};
100+
finalResponse = requireFinalResponse("Claude", finalResponse);
101+
return {
102+
provider: this.provider,
103+
providerSessionId,
104+
finalResponse,
105+
items,
106+
...(structured !== undefined ? { structured } : {}),
107+
};
108+
} catch (error) {
109+
if (input.schema && isNativeSchemaUnsupportedFailure(error)) {
110+
throw new ProviderSchemaUnsupportedError(this.provider, error);
111+
}
112+
throw error;
113+
}
103114
}
104115
}
105116

src/local-agent-runtime.test.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import type { RunResult, ThreadOptions } from "@openai/codex-sdk";
33
import {
44
CodexSdkLocalAgentRuntime,
55
createCodexSdkLocalAgentRuntime,
6+
isNativeSchemaUnsupportedFailure,
67
} from "./local-agent-runtime.js";
78

89
const emptyTurn = (finalResponse: string): RunResult => ({
@@ -127,3 +128,11 @@ assert.deepEqual(codex.resumed, [
127128

128129
const created = await createCodexSdkLocalAgentRuntime(undefined, () => new FakeCodex());
129130
assert.equal(created.provider, "codex");
131+
132+
assert.equal(
133+
isNativeSchemaUnsupportedFailure(
134+
new Error("Invalid output schema: keyword is not supported"),
135+
),
136+
true,
137+
);
138+
assert.equal(isNativeSchemaUnsupportedFailure(new Error("authentication failed")), false);

src/local-agent-runtime.ts

Lines changed: 46 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import type {
55
RunResult,
66
SandboxMode,
77
ThreadOptions,
8+
TurnOptions,
89
} from "@openai/codex-sdk";
910

1011
export type LocalAgentWriteMode = "read_only" | "allowed" | "full_access";
@@ -35,14 +36,42 @@ export interface LocalAgentRuntime {
3536
run(input: LocalAgentRunInput): Promise<LocalAgentRunResult>;
3637
}
3738

38-
interface CodexTurnOptions {
39-
outputSchema?: unknown;
40-
signal?: AbortSignal;
39+
export class ProviderSchemaUnsupportedError extends Error {
40+
constructor(
41+
readonly provider: string,
42+
readonly cause: unknown,
43+
) {
44+
super(`${provider} does not support the requested native output schema: ${errorMessage(cause)}`);
45+
this.name = "ProviderSchemaUnsupportedError";
46+
}
47+
}
48+
49+
export function isProviderSchemaUnsupportedError(
50+
error: unknown,
51+
): error is ProviderSchemaUnsupportedError {
52+
return error instanceof ProviderSchemaUnsupportedError;
53+
}
54+
55+
export function isNativeSchemaUnsupportedFailure(error: unknown): boolean {
56+
const message = errorMessage(error).toLowerCase();
57+
const mentionsSchema =
58+
/output[ _-]?schema/.test(message) ||
59+
/json[ _-]?schema/.test(message) ||
60+
/structured[ _-]?output/.test(message) ||
61+
/output[ _-]?format/.test(message);
62+
const unsupported =
63+
/not supported/.test(message) ||
64+
/unsupported/.test(message) ||
65+
/invalid (?:output|json )?schema/.test(message) ||
66+
/schema (?:is )?invalid/.test(message) ||
67+
/unknown (?:field|parameter|option)/.test(message) ||
68+
/not available/.test(message);
69+
return mentionsSchema && unsupported;
4170
}
4271

4372
interface CodexThreadLike {
4473
readonly id: string | null;
45-
run(prompt: string, turnOptions?: CodexTurnOptions): Promise<RunResult>;
74+
run(prompt: string, turnOptions?: TurnOptions): Promise<RunResult>;
4675
}
4776

4877
interface CodexClientLike {
@@ -88,7 +117,15 @@ export class CodexSdkLocalAgentRuntime implements LocalAgentRuntime {
88117
? this.codex.resumeThread(input.providerSessionId, options)
89118
: this.codex.startThread(options);
90119
const turnOptions = input.schema ? { outputSchema: input.schema } : undefined;
91-
const turn = await thread.run(input.prompt, turnOptions);
120+
let turn: RunResult;
121+
try {
122+
turn = await thread.run(input.prompt, turnOptions);
123+
} catch (error) {
124+
if (input.schema && isNativeSchemaUnsupportedFailure(error)) {
125+
throw new ProviderSchemaUnsupportedError(this.provider, error);
126+
}
127+
throw error;
128+
}
92129

93130
return {
94131
provider: this.provider,
@@ -120,3 +157,7 @@ async function defaultCodexFactory(): Promise<CodexFactory> {
120157
const module = await import("@openai/codex-sdk");
121158
return (options) => new module.Codex(options) as Codex;
122159
}
160+
161+
function errorMessage(error: unknown): string {
162+
return error instanceof Error ? error.message : String(error);
163+
}

src/workflow-api.ts

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import {
1919
export interface WorkflowProviderRunInput {
2020
provider: string;
2121
prompt: string;
22+
providerSessionId?: string;
2223
model?: string;
2324
effort?: string;
2425
workspace: string;
@@ -351,12 +352,12 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi {
351352
schema: agentOpts.schema,
352353
prompt,
353354
provider,
354-
run: (p) =>
355+
run: (p, options) =>
355356
deps.runProvider({
356357
...providerBase,
357358
prompt: p,
358-
// Keep schema on adapter for codex/claude native+repair attempts.
359-
schema: agentOpts.schema,
359+
providerSessionId: options.providerSessionId,
360+
...(options.mode === "native" ? { schema: agentOpts.schema } : {}),
360361
}),
361362
onRetry: ({ attempt, errors, mode }) => {
362363
deps.journal.appendEvent({

src/workflow-cli.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -326,6 +326,7 @@ export async function runWorkflowWorker(
326326
const providerResult = await runLocalAgentProvider(input.provider, {
327327
prompt: input.prompt,
328328
workspace: input.workspace,
329+
providerSessionId: input.providerSessionId,
329330
model: input.model,
330331
effort: input.effort,
331332
writeMode: "allowed",

src/workflow-engine.test.ts

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -279,6 +279,59 @@ import { createStubBudget } from "./workflow-types.js";
279279
await rm(dir, { recursive: true, force: true });
280280
}
281281

282+
// ---------------------------------------------------------------------------
283+
// schema retry: native schema only on first attempt + provider session reuse
284+
// ---------------------------------------------------------------------------
285+
{
286+
const dir = await mkdtemp(join(tmpdir(), "wf-schema-retry-"));
287+
const store = new WorkflowStore(dir);
288+
const run = store.createRun({
289+
name: "schema-retry",
290+
source: "inline",
291+
scriptPath: "inline",
292+
scriptHash: "h",
293+
workspaceRoot: dir,
294+
});
295+
const calls: WorkflowProviderRunInput[] = [];
296+
const api = createWorkflowApi({
297+
runId: run.id,
298+
journal: store,
299+
meta: { name: "schema-retry", description: "d" },
300+
args: undefined,
301+
concurrency: 1,
302+
signal: new AbortController().signal,
303+
workspaceRoot: dir,
304+
enabledProviders: ["codex"],
305+
runProvider: async (input) => {
306+
calls.push(input);
307+
if (calls.length === 1) {
308+
return {
309+
finalResponse: '{"n":"bad"}',
310+
structured: { n: "bad" },
311+
providerSessionId: "sess-1",
312+
};
313+
}
314+
return { finalResponse: '{"n":2}', providerSessionId: "sess-1" };
315+
},
316+
});
317+
318+
const out = await api.agent("give n", {
319+
schema: {
320+
type: "object",
321+
properties: { n: { type: "number" } },
322+
required: ["n"],
323+
},
324+
});
325+
assert.deepEqual(out, { n: 2 });
326+
assert.ok(calls[0]?.schema);
327+
assert.equal(calls[0]?.providerSessionId, undefined);
328+
assert.equal(calls[1]?.schema, undefined);
329+
assert.equal(calls[1]?.providerSessionId, "sess-1");
330+
331+
store.close();
332+
await rm(dir, { recursive: true, force: true });
333+
}
334+
282335
// ---------------------------------------------------------------------------
283336
// executeWorkflow end-to-end with sandbox + nest depth
284337
// ---------------------------------------------------------------------------

0 commit comments

Comments
 (0)