Skip to content

Commit 96146d2

Browse files
committed
[agent] feat: Anthropic SDK adapter — executeStream via @anthropic-ai/sdk
- Add @anthropic-ai/sdk@0.105.0 dependency - base.ts: add optional executeStream() to ProviderAdapter interface - anthropic.ts: implement executeStream() using SDK messages.create(stream: true) - Type-safe params, auto-retry on 429/5xx, SDK-managed SSE parsing - Existing buildRequest/parseStream kept as fallback for web-search sidecar - server.ts: prefer executeStream when available, fall back to fetch+parseStream
1 parent 367c642 commit 96146d2

5 files changed

Lines changed: 151 additions & 26 deletions

File tree

bun.lock

Lines changed: 15 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

package.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
"release:watch": "bun scripts/release.ts watch"
3030
},
3131
"dependencies": {
32+
"@anthropic-ai/sdk": "^0.105.0",
3233
"zod": "^4.0.0"
3334
},
3435
"devDependencies": {

src/adapters/anthropic.ts

Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import Anthropic from "@anthropic-ai/sdk";
2+
import type { RawMessageStreamEvent } from "@anthropic-ai/sdk/resources/messages.mjs";
13
import type { ProviderAdapter } from "./base";
24
import { debugDroppedFrame } from "../debug";
35
import type {
@@ -312,5 +314,106 @@ export function createAnthropicAdapter(provider: OcxProviderConfig): ProviderAda
312314
});
313315
return events;
314316
},
317+
318+
async *executeStream(parsed: OcxParsedRequest, signal?: AbortSignal): AsyncGenerator<AdapterEvent> {
319+
const { system, messages } = messagesToAnthropicFormat(parsed, isOAuth);
320+
const tools = toolsToAnthropicFormat(parsed, isOAuth);
321+
322+
const sdkParams: Record<string, unknown> = {
323+
model: parsed.modelId,
324+
messages,
325+
max_tokens: parsed.options.maxOutputTokens ?? DEFAULT_MAX_TOKENS,
326+
stream: true,
327+
};
328+
if (isOAuth) {
329+
sdkParams.system = [
330+
{ type: "text", text: CLAUDE_CODE_SYSTEM_INSTRUCTION },
331+
...(system ? [{ type: "text", text: system }] : []),
332+
];
333+
} else if (system) {
334+
sdkParams.system = system;
335+
}
336+
if (tools) sdkParams.tools = tools;
337+
if (parsed.options.temperature !== undefined) sdkParams.temperature = parsed.options.temperature;
338+
if (parsed.options.topP !== undefined) sdkParams.top_p = parsed.options.topP;
339+
if (parsed.options.stopSequences) sdkParams.stop_sequences = parsed.options.stopSequences;
340+
341+
if (parsed.options.reasoning) {
342+
const maxOut = parsed.options.maxOutputTokens ?? DEFAULT_MAX_TOKENS;
343+
const wantBudget = reasoningBudget(parsed.options.reasoning);
344+
const maxTokens = Math.min(REASONING_MAX_TOKENS_CEILING, Math.max(maxOut, wantBudget + OUTPUT_HEADROOM));
345+
const budget = Math.max(MIN_THINKING_BUDGET, Math.min(wantBudget, maxTokens - OUTPUT_FLOOR));
346+
sdkParams.max_tokens = maxTokens;
347+
sdkParams.thinking = { type: "enabled", budget_tokens: budget };
348+
delete sdkParams.temperature;
349+
delete sdkParams.top_p;
350+
}
351+
352+
if (parsed.options.toolChoice) {
353+
const tc = parsed.options.toolChoice;
354+
if (tc === "auto") sdkParams.tool_choice = { type: "auto" };
355+
else if (tc === "none") sdkParams.tool_choice = { type: "none" };
356+
else if (tc === "required") sdkParams.tool_choice = { type: "any" };
357+
else if (typeof tc === "object" && "name" in tc) sdkParams.tool_choice = { type: "tool", name: isOAuth ? applyClaudeToolPrefix(tc.name) : tc.name };
358+
}
359+
360+
const headers: Record<string, string> = {};
361+
if (isOAuth) headers["anthropic-beta"] = ANTHROPIC_OAUTH_BETA;
362+
if (provider.headers) Object.assign(headers, provider.headers);
363+
364+
const client = new Anthropic({
365+
apiKey: provider.apiKey ?? "",
366+
baseURL: `${provider.baseUrl}/v1`,
367+
maxRetries: 2,
368+
defaultHeaders: headers,
369+
...(isOAuth ? { authToken: provider.apiKey } : {}),
370+
});
371+
372+
const stream = await client.messages.create(
373+
sdkParams as unknown as Parameters<typeof client.messages.create>[0],
374+
{ signal },
375+
);
376+
377+
let currentToolName = "";
378+
for await (const event of stream as AsyncIterable<RawMessageStreamEvent>) {
379+
switch (event.type) {
380+
case "content_block_start": {
381+
if (event.content_block.type === "tool_use") {
382+
const name = isOAuth ? stripClaudeToolPrefix(event.content_block.name) : event.content_block.name;
383+
currentToolName = name;
384+
yield { type: "tool_call_start", id: event.content_block.id, name };
385+
}
386+
break;
387+
}
388+
case "content_block_delta": {
389+
if (event.delta.type === "text_delta") {
390+
yield { type: "text_delta", text: event.delta.text };
391+
} else if (event.delta.type === "thinking_delta") {
392+
yield { type: "thinking_delta", thinking: event.delta.thinking };
393+
} else if (event.delta.type === "input_json_delta") {
394+
yield { type: "tool_call_delta", arguments: event.delta.partial_json };
395+
}
396+
break;
397+
}
398+
case "content_block_stop": {
399+
if (currentToolName) {
400+
yield { type: "tool_call_end" };
401+
currentToolName = "";
402+
}
403+
break;
404+
}
405+
case "message_delta": {
406+
const u = event.usage;
407+
yield {
408+
type: "done",
409+
usage: u ? usageFromAnthropic({ output_tokens: u.output_tokens }) : undefined,
410+
};
411+
break;
412+
}
413+
default:
414+
break;
415+
}
416+
}
417+
},
315418
};
316419
}

src/adapters/base.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,4 +17,8 @@ export interface ProviderAdapter {
1717

1818
parseStream(response: Response): AsyncGenerator<AdapterEvent>;
1919
parseResponse?(response: Response): Promise<AdapterEvent[]>;
20+
21+
/** Optional: adapter handles its own fetch + streaming via a provider SDK. When present,
22+
* server.ts uses this instead of buildRequest → fetch → parseStream. */
23+
executeStream?(parsed: OcxParsedRequest, signal?: AbortSignal): AsyncGenerator<AdapterEvent>;
2024
}

src/server.ts

Lines changed: 28 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -204,19 +204,35 @@ async function handleResponses(
204204
});
205205
}
206206

207-
const request = adapter.buildRequest(parsed, { headers: req.headers });
208-
209-
// Abort the upstream fetch if the client (Codex) disconnects mid-stream, so a cancelled turn does
210-
// not leak the upstream connection or keep draining tokens. The bridge's cancel() fires upstream.abort() (RC2).
211207
const upstream = new AbortController();
212208
linkAbortSignal(upstream, options.abortSignal);
209+
210+
if (parsed.stream && adapter.executeStream) {
211+
const eventStream = adapter.executeStream(parsed, upstream.signal);
212+
const toolNsMap = new Map<string, { namespace: string; name: string }>();
213+
const freeformToolNames = new Set<string>();
214+
const toolSearchToolNames = new Set<string>();
215+
for (const t of parsed.context.tools ?? []) {
216+
if (t.namespace) toolNsMap.set(namespacedToolName(t.namespace, t.name), { namespace: t.namespace, name: t.name });
217+
if (t.freeform) freeformToolNames.add(t.name);
218+
if (t.toolSearch) toolSearchToolNames.add(t.name);
219+
}
220+
const sseStream = bridgeToResponsesSSE(
221+
eventStream, parsed.modelId, toolNsMap, freeformToolNames, toolSearchToolNames,
222+
() => upstream.abort(), 2_000,
223+
options.forceEmptyResponseId ? { responseId: "" } : undefined,
224+
);
225+
return new Response(sseStream, {
226+
headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no" },
227+
});
228+
}
229+
230+
// Fallback: buildRequest → fetch → parseStream (used by all non-SDK adapters and web-search sidecar)
231+
const request = adapter.buildRequest(parsed, { headers: req.headers });
213232
let upstreamResponse: Response;
214233
try {
215234
upstreamResponse = await fetch(request.url, {
216-
method: request.method,
217-
headers: request.headers,
218-
body: request.body,
219-
signal: upstream.signal,
235+
method: request.method, headers: request.headers, body: request.body, signal: upstream.signal,
220236
});
221237
} catch (err) {
222238
return formatErrorResponse(502, "upstream_error", `Provider unreachable: ${err instanceof Error ? err.message : String(err)}`);
@@ -229,8 +245,6 @@ async function handleResponses(
229245

230246
if (parsed.stream) {
231247
const eventStream = adapter.parseStream(upstreamResponse);
232-
// Map flattened MCP tool names back to {namespace, name} so the bridge can restore the
233-
// namespace field Codex needs to route the call to the right MCP server.
234248
const toolNsMap = new Map<string, { namespace: string; name: string }>();
235249
const freeformToolNames = new Set<string>();
236250
const toolSearchToolNames = new Set<string>();
@@ -240,31 +254,19 @@ async function handleResponses(
240254
if (t.toolSearch) toolSearchToolNames.add(t.name);
241255
}
242256
const sseStream = bridgeToResponsesSSE(
243-
eventStream,
244-
parsed.modelId,
245-
toolNsMap,
246-
freeformToolNames,
247-
toolSearchToolNames,
248-
() => upstream.abort(),
249-
2_000,
257+
eventStream, parsed.modelId, toolNsMap, freeformToolNames, toolSearchToolNames,
258+
() => upstream.abort(), 2_000,
250259
options.forceEmptyResponseId ? { responseId: "" } : undefined,
251260
);
252261
return new Response(sseStream, {
253-
headers: {
254-
"Content-Type": "text/event-stream",
255-
"Cache-Control": "no-cache",
256-
"Connection": "keep-alive",
257-
"X-Accel-Buffering": "no",
258-
},
262+
headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no" },
259263
});
260264
}
261265

262266
if (adapter.parseResponse) {
263267
const events = await adapter.parseResponse(upstreamResponse);
264268
const json = buildResponseJSON(events, parsed.modelId);
265-
return new Response(JSON.stringify(json), {
266-
headers: { "Content-Type": "application/json" },
267-
});
269+
return new Response(JSON.stringify(json), { headers: { "Content-Type": "application/json" } });
268270
}
269271

270272
return formatErrorResponse(500, "internal_error", "Non-streaming not supported by this adapter");

0 commit comments

Comments
 (0)