openpencil/apps/web/server/api/ai/chat.ts
Kayshen Xu e9b0d0d822 V0.7.2-bugfix (#109)
* Stabilize synced main for AI handoff, drag nesting, and Electron dev (#104)

* docs(readme): update cover screenshot

* fix: stabilize electron dev sync and codex env passthrough

* Preserve nested frame behavior during drag reparenting

Reparenting across containers used raw local coordinates and root-only clipping assumptions, which made nodes jump visually and caused dragged frames to lose clip/corner semantics after nesting. This adapts the drag-reparent fix to the current upstream store architecture, keeps frame/shape nodes from auto-detaching on canvas drags, and promotes formerly root-only frame clipping to explicit clipContent when nested.

Constraint: Latest upstream workspace checkout is incomplete locally (missing workspaces/deps), so full upstream verification could not be rerun in this environment
Rejected: Keep using raw local x/y during parent changes | fails for auto-layout/padding-rendered positions
Rejected: Make all nested frames clip unconditionally | would change non-clipping containers
Confidence: medium
Scope-risk: moderate
Reversibility: clean
Directive: Preserve visual-position conversion through rendered coordinates when parent changes; local coordinates alone are insufficient once layout participates
Not-tested: Fresh full workspace typecheck/test/build on latest upstream checkout (blocked by missing workspace/dependency setup in this local clone)

* Keep AI codegen requests bounded while exporting asset bundles

The AI codegen pipeline needed two stability fixes: exported design images had to flow through chunk/assembly prompts as reusable asset hints, and oversized chat payloads needed a local guard before hitting provider limits. This commit wires asset extraction into the planning pipeline, threads exported asset paths into prompt assembly, and rejects obviously overlarge chat requests with an actionable client-side error.

Constraint: This branch is split out from a larger local fix stack, so only codegen/prompt/context files are included here
Constraint: Provider request limits are approximate locally, so the payload guard must be conservative rather than exact
Rejected: Inline base64 assets directly into prompts | explodes request size and repeats the same payload per chunk
Rejected: Let provider errors handle oversized payloads | too slow and opaque for users
Confidence: high
Scope-risk: moderate
Reversibility: clean
Directive: Keep asset references flowing as stable ./assets paths and enforce payload limits before fetch to avoid silent request bloat
Tested: bun x tsc -p apps/web/tsconfig.json --noEmit; cd apps/web && bun --bun vitest run src/services/ai/__tests__/context-optimizer.test.ts src/services/ai/__tests__/codegen-assets.test.ts src/services/ai/__tests__/structure-bundle.test.ts; bun run build
Not-tested: Manual end-to-end AI generation with live providers

* Explain sanitized design views instead of leaving AI to guess

The sanitized structure bundle already stabilized asset paths, but it still exposed low-level image/layout/component fields that models had to interpret on their own. This change adds explicit consumer-view enrichment for fills, layout, text, variables, themes, and component semantics, carries original image size through the Figma import path, and augments sanitized bundles with summary/highlight guidance for downstream AI consumers.

Constraint: This branch is intentionally stacked on the asset-bundle PR because it extends the sanitized/codegen asset pipeline rather than replacing it
Constraint: Figma import data is not always complete, so original image size must be preserved when present and inferred only as a fallback downstream
Rejected: Keep sanitized.json as a pure field-level dump | still leaves AI to misread transforms, layout, and component relationships
Rejected: Put all explain text directly in asset extraction helpers | mixes resource stabilization with semantic enrichment responsibilities
Confidence: high
Scope-risk: moderate
Reversibility: clean
Directive: Treat consumer-view enrichment as a distinct layer on top of stable asset extraction; future AI-facing semantics should land there instead of leaking into unrelated pipeline code
Tested: bun x tsc -p apps/web/tsconfig.json --noEmit; cd apps/web && bun --bun vitest run src/services/ai/__tests__/consumer-view-enrichment.test.ts src/services/ai/__tests__/codegen-assets.test.ts src/services/ai/__tests__/structure-bundle.test.ts ../../packages/pen-figma/src/figma-fill-mapper.test.ts; bun run build
Not-tested: Manual prompt-to-code generation quality with live provider responses

* Restore code-panel bundle exports for AI handoff flows

The code generation backend still produced asset manifests and AI structure bundles, but the code panel UI no longer exposed those export paths after later sync work. This commit reconnects the panel to bundle export actions, restores ZIP download behavior when generated code includes exported assets, and locks the affordances with focused panel tests.

Constraint: Other local fixes are still in progress in the working tree, so this commit is intentionally limited to the code-panel export surface
Rejected: Rebuild export support in a separate panel | users expect the export actions to remain where generation results are shown
Confidence: high
Scope-risk: narrow
Reversibility: clean
Directive: Keep code-panel UI aligned with codegen asset/bundle backends whenever generation result shape changes
Tested: cd apps/web && bun --bun vitest run src/components/panels/code-panel.test.tsx src/services/ai/__tests__/codegen-assets.test.ts src/services/ai/__tests__/structure-bundle.test.ts; bun run build
Not-tested: Manual click-through of AI Bundle and Download ZIP in the desktop/web UI

* Unblock electron dev startup in the incomplete local workspace

The local workspace was failing before the app could even start: the skills plugin hard-required js-yaml from a node_modules layout that was not present, Vite dev under Bun hit Nitro NodeResponse incompatibilities, and the web tsconfig was missing path mappings for local packages. This commit removes the unnecessary js-yaml dependency from the skills loader, runs Vite under Node for dev startup, hardens readiness probing with socket checks, and points TypeScript/Vite at the in-repo package sources.

Constraint: The current local clone has incomplete hoisted/workspace installation state, so dev startup must not depend on root package links being perfectly present
Constraint: Bun + Nitro dev currently mis-handle NodeResponse in this environment, so the safest startup path is Node-hosted Vite
Rejected: Keep js-yaml and require everyone to fix local hoisting first | still leaves electron:dev broken in the current environment
Rejected: Continue running Vite dev through Bun | reproduces the NodeResponse/Parse Error failure on /api and /editor requests
Confidence: high
Scope-risk: narrow
Reversibility: clean
Directive: Keep the dev launcher biased toward resilient local startup, even when the workspace install shape is imperfect
Tested: bun -e import('./packages/pen-ai-skills/vite-plugin-skills.ts').then(() => console.log('SKILL_PLUGIN_IMPORT_OK')); bun electron:dev verified Vite ready, MCP/Electron compiled, Electron launched, MCP sync log emitted
Not-tested: Long-running interactive desktop session after startup

* fix(figma): preserve cropped image fill transforms

The synced branch started exporting original image dimensions but dropped the
existing crop transform semantics from the shared image-fill type and both
Figma mappers. That broke the new regression test and stripped metadata that
AI consumer-view/bundle code already relies on.

Constraint: keep app and package Figma mappers in lockstep
Rejected: loosen the new regression test | would hide a real metadata regression
Confidence: high
Scope-risk: narrow
Reversibility: clean
Directive: when extending image fill metadata, update shared pen-types and both Figma mapper copies together
Tested: bun --bun run test (148/149 files passed; only server/__tests__/sse-keepalive.test.ts blocked by missing agent_napi.node), cd apps/web && bun --bun vitest run src/canvas/skia/drag-reparent-policy.test.ts src/components/panels/layer-dnd-utils.test.ts src/stores/document-position-utils.test.ts src/components/panels/code-panel.test.tsx ../../packages/pen-renderer/src/__tests__/document-flattener.test.ts ../../packages/pen-figma/src/figma-fill-mapper.test.ts, cd apps/web && bun --bun vitest run src/services/ai/__tests__/codegen-assets.test.ts src/services/ai/__tests__/structure-bundle.test.ts src/services/ai/__tests__/consumer-view-enrichment.test.ts, cd apps/web && bun --bun vitest run src/utils/__tests__/security.test.ts, bun test scripts/loopback-no-proxy.test.ts, npx tsc --noEmit, bun --bun run build
Not-tested: server/__tests__/sse-keepalive.test.ts without a locally built @zseven-w/agent-native addon

* docs(editor): normalize new PR comments to English

The PR had a handful of newly introduced Chinese code comments in dev, sync, and AI helper paths. This follow-up keeps the implementation unchanged while translating those comments to English so the PR stays consistent with the repository comment-language expectation.

Constraint: The request was limited to comment language cleanup after the conflict-resolution merge, so behavior had to remain unchanged
Rejected: Leave the mixed-language comments in place | conflicts with the PR requirement for English comments
Rejected: Broader repository-wide translation sweep | unnecessary scope expansion beyond the PR-introduced comments
Confidence: high
Scope-risk: narrow
Reversibility: clean
Directive: Keep code comments in English on this branch, even when local notes or working memory are in another language
Tested: bun test scripts/loopback-no-proxy.test.ts apps/desktop/__tests__/dev-utils.test.ts; cd apps/web && bun --bun vitest run server/__tests__/mcp-sync-state-active.test.ts src/canvas/skia/__tests__/skia-interaction.test.ts; npx tsc --noEmit; branch-diff comment scan for Han characters in comment lines
Not-tested: Manual runtime behavior, since this change only rewrote comments

* style(editor): apply repository formatting expected by CI

The PR was failing the CI Format check after the conflict-resolution and comment-normalization follow-ups. This commit applies the repository formatter output to the files touched by the branch so CI sees the exact formatting it expects, without changing behavior.

Constraint: The failing GitHub Actions job stopped at Format check, so the fix had to match oxfmt output rather than introduce functional changes
Rejected: Leave the branch as-is and rely on local formatting differences being acceptable | CI explicitly rejects the current formatting
Rejected: Broader code cleanup beyond formatter output | unnecessary scope while repairing the failing check
Confidence: high
Scope-risk: narrow
Reversibility: clean
Directive: After conflict resolution or comment-only edits on this repo, run bun run format:check before pushing because formatter expectations are stricter than the existing file style in some touched files
Tested: bun run format:check; bun run lint; npx tsc --noEmit
Not-tested: Full test suite after this formatting-only commit (previous run showed formatting was the first CI blocker)

* refactor(editor): remove proxy-specific dev workarounds from PR

The PR no longer needs the loopback proxy bypass layer, so this cleanup removes the proxy-specific dev entrypoint, environment bootstrap, helper module, and its tests while keeping the unrelated Electron and AI handoff changes intact.

Constraint: Removal had to be limited to proxy-related code on PR #104 without undoing the other merged fixes on the branch
Rejected: Keep the helper and stop using it | leaves proxy-specific maintenance surface and tests in the PR
Rejected: Revert the entire Electron dev file to upstream earlier than necessary | would risk dropping unrelated local conflict-resolution choices beyond the proxy scope
Confidence: high
Scope-risk: moderate
Reversibility: clean
Directive: If proxy handling is reintroduced later, keep it out of this PR unless there is a dedicated, separately justified change for it
Tested: bun run format:check; bun run lint; npx tsc --noEmit
Not-tested: Manual electron:dev behavior after removing the proxy-specific launcher path

* docs(ai): translate JSON-facing semantic descriptions to English

The PR still emitted Chinese semantic description strings inside the AI consumer-view and structure-bundle JSON outputs. This change translates those JSON-facing runtime descriptions and updates the affected tests so exported AI-facing structure data is consistently English.

Constraint: The request was limited to JSON description strings, so the change had to preserve the same semantics and structure while only translating output text
Rejected: Leave Chinese test fixtures and runtime descriptions in place | conflicts with the requirement for English JSON descriptions
Rejected: Broader i18n cleanup outside these AI JSON description paths | unnecessary scope expansion beyond the requested exported-description surface
Confidence: high
Scope-risk: moderate
Reversibility: clean
Directive: Keep AI/exported JSON explanation strings in English unless a future change explicitly adds localized output modes
Tested: cd apps/web && bun --bun vitest run src/services/ai/__tests__/consumer-view-enrichment.test.ts src/services/ai/__tests__/structure-bundle.test.ts src/services/ai/__tests__/codegen-assets.test.ts; bun run format:check; npx tsc --noEmit
Not-tested: Full app runtime flows that consume these JSON descriptions outside the covered unit tests

* refactor(ai): remove remaining network-proxy handling

The current project still carried Anthropic proxy-specific heuristics and environment handling outside the PR-specific cleanup. Since the earlier crashes and connectivity issues were unrelated to proxying, this removes the remaining network-proxy branches, model remapping, and TLS override advice while leaving unrelated request flows intact.

Constraint: The cleanup needed to remove proxy-specific logic without disturbing unrelated transport concepts such as app-internal API proxy routes or React proxy objects used in tests
Rejected: Keep the proxy heuristics as dormant fallback logic | preserves misleading operational guidance and dead maintenance surface
Rejected: Rename every remaining literal use of the word proxy in the repo | would overreach into unrelated concepts like internal API proxying and JS Proxy-based test setup
Confidence: medium
Scope-risk: moderate
Reversibility: clean
Directive: If endpoint-specific compatibility logic is needed later, add it as explicit endpoint handling rather than generic proxy heuristics
Tested: bun run format:check; bun run lint; npx tsc --noEmit; repo-wide search for network-proxy env references after cleanup
Not-tested: End-to-end Claude connection flows against custom base URLs after removing proxy-specific remapping

* fix(electron): keep Node-backed dev launch for Nitro compatibility

Comparing against upstream commit 7271a03 confirms the current Electron dev fix is not the same idea as the original Bun-based launcher. The upstream version starts Vite with Bun, while the observed failure shows Nitro now crashes in that path with "Vite environment nitro is unavailable". This keeps the non-proxy Node-backed launcher because it fixes the actual regression without restoring the removed proxy code.

Constraint: The request preferred reverting to the upstream original only if the intent matched, but the current Nitro/Electron failure proves the upstream Bun launcher is no longer equivalent in behavior
Rejected: Restore the exact 7271a03 Bun launcher | reproduces the Nitro dev-worker crash and ERR_EMPTY_RESPONSE in Electron
Rejected: Reintroduce the old proxy workaround bundle | unrelated to the reproduced failure and already removed by request
Confidence: high
Scope-risk: narrow
Reversibility: clean
Directive: Keep Electron dev on the Node-backed Vite launcher unless Nitro/Bun dev compatibility is revalidated with a real startup test
Tested: bun run electron:dev (reached Electron launch after Vite/MCP/Electron compile steps); bun test apps/desktop/__tests__/dev-utils.test.ts; bun run format:check; npx tsc --noEmit
Not-tested: Full interactive manual editor workflow after Electron launch

---------

Co-authored-by: Fini <fini.yang@gmail.com>

* fix(ai,cli): openai-compat turn-2, StepFun reasoning+451, Mac CLI discovery

Round up the v0.7.2 stability fixes for AI connectivity and local CLI
detection that surfaced during real user runs against GLM, StepFun, and
Mac users on nvm/fnm/pnpm/bun/mise/asdf/fish shells.

Provider (via @zseven-w/agent-native v0.3.0 submodule bump):
- OpenAI-compat providers can now complete multi-turn tool-calling loops:
  the request builder translates Anthropic-shaped message history
  (tool_use / tool_result blocks, thinking) into OpenAI's tool_calls +
  role="tool" form so turn 2 no longer 400s. system_prompt is finally
  injected instead of being silently dropped.
- The SSE parser accepts `delta.reasoning` (StepFun step_plan) alongside
  `reasoning_content` (GLM / DeepSeek / Qwen), and also streams tool_call
  fragments, which unblocks GLM / dashscope and stops the
  firstTextTimeout → fetch abort → std.http panic → Bun segfault cascade.
- HTTP 451 (StepFun content-safety) surfaces as InvalidRequest with a
  specific "content blocked by provider safety filter" message instead
  of an opaque error_server.

Server route + client watchdog:
- /api/ai/chat forwards the provider's last_error string
  (result.errors[0]) so users see "HTTP 451 content blocked" rather than
  "Provider error: error_server".
- streamChat clears firstTextTimeout on thinking chunks (when
  thinkingResetsTimeout=true), so models that stream long reasoning
  before any text aren't falsely killed as "stuck".

Orchestrator sub-agent resilience:
- Failed sub-agents (empty response / unparseable output) now retry once
  with a minimal ~3KB kernel prompt (schema + jsonl-format only). Only
  the failing subtask re-runs — successful earlier sections are kept.
- Deterministic refusals (HTTP 400/401/429/451, "content blocked",
  "censorship", "authentication failed") short-circuit the retry ladder
  so a 4-minute StepFun safety scan isn't spent twice in a row.

Local CLI discovery (Mac users on managed shells):
- New server/utils/cli-resolver-helpers.ts exports probeViaLoginShell()
  and posixUserBinDirs(). Login-shell probe asks $SHELL (or zsh/bash
  fallback — fish added at /opt/homebrew/bin/fish and friends) with
  `-ilc 'command -v <cli>'` so nvm/pnpm/bun/mise/asdf/volta/fnm shims
  are visible even when Electron scrubs the inherited PATH.
- resolveClaudeCli / resolveGeminiCli / resolveCopilotCli and the
  inline codex/opencode resolvers in connect-agent.ts all run the same
  PATH → login-shell → npm-prefix → user-bin candidates ladder. Each
  step logs via serverLog to ~/.openpencil/logs/server-YYYY-MM-DD.log
  for remote diagnosis.

Builtin provider preset:
- Add StepFun Coding Plan (api.stepfun.com/step_plan/v1, label "StepFun
  Coding Plan") alongside the existing StepFun preset.

Version bump 0.7.1 → 0.7.2 across all workspaces.

---------

Co-authored-by: RaisCui <857943+raiscui@users.noreply.github.com>
Co-authored-by: Fini <fini.yang@gmail.com>
2026-04-14 21:42:56 +08:00

1145 lines
42 KiB
TypeScript

import { defineEventHandler, readBody, setResponseHeaders } from 'h3';
import { writeFile, mkdtemp, rm } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { resolveClaudeCli } from '../../utils/resolve-claude-cli';
import { runCodexExec } from '../../utils/codex-client';
import { startSSEKeepAlive } from '../../utils/sse-keepalive';
import {
buildClaudeAgentEnv,
buildSpawnClaudeCodeProcess,
getClaudeAgentDebugFilePath,
} from '../../utils/resolve-claude-agent-env';
import { normalizeOptionalBaseURL, requireOpenAICompatBaseURL } from './provider-url';
// SENSITIVE_LOG_PATTERN + readDebugTail are now canonical in @zseven-w/pen-mcp.
// Re-export here to keep existing consumers (tests, other modules) working.
import { SENSITIVE_LOG_PATTERN, readDebugTail } from '@zseven-w/pen-mcp';
export { SENSITIVE_LOG_PATTERN };
/** Allowed media types for image attachments */
export const ALLOWED_MEDIA_TYPES = new Set(['image/png', 'image/jpeg', 'image/gif', 'image/webp']);
/** Resolve file extension from media type, falling back to 'png' for disallowed types */
export function resolveMediaExtension(mediaType: string): string {
return ALLOWED_MEDIA_TYPES.has(mediaType) ? mediaType.split('/')[1] : 'png';
}
interface ChatAttachmentWire {
name: string;
mediaType: string;
data: string; // base64
}
interface ChatBody {
system: string;
messages: Array<{
role: 'user' | 'assistant';
content: string;
attachments?: ChatAttachmentWire[];
}>;
model?: string;
provider?: 'anthropic' | 'openai' | 'opencode' | 'copilot' | 'gemini' | 'builtin';
thinkingMode?: 'adaptive' | 'disabled' | 'enabled';
thinkingBudgetTokens?: number;
effort?: 'low' | 'medium' | 'high' | 'max';
/** For builtin provider: direct API key (not CLI-based) */
builtinApiKey?: string;
/** For builtin provider: API root base URL (e.g. https://api.openai.com/v1) */
builtinBaseURL?: string;
/** For builtin provider: 'anthropic' or 'openai-compat' */
builtinType?: 'anthropic' | 'openai-compat';
}
function buildClaudeExitHint(rawError: string, debugTail?: string[]): string | undefined {
if (!/process exited with code 1/i.test(rawError)) return undefined;
const hints: string[] = [];
if (debugTail && debugTail.length > 0) {
const text = debugTail.join('\n');
if (
/Failed to save config with lock: Error: EPERM|operation not permitted, .*\.claude\.json/i.test(
text,
)
) {
hints.push(
'Claude Code cannot write ~/.claude.json (permission denied). ' +
'On Windows, try running as Administrator or manually create the file: echo {} > %USERPROFILE%\\.claude.json',
);
}
if (
/Connection error|Could not resolve host|Failed to connect|ECONNREFUSED|ETIMEDOUT/i.test(text)
) {
hints.push(
'Upstream API connection failed. Check DNS and network reachability to your configured endpoint.',
);
}
if (/ANTHROPIC_CUSTOM_HEADERS present: false, has Authorization header: false/i.test(text)) {
hints.push(
'No API auth header detected. Run "claude login" to authenticate, ' +
'or set ANTHROPIC_API_KEY in ~/.claude/settings.json ' +
'(env: { "ANTHROPIC_API_KEY": "sk-..." }).',
);
}
if (/invalid.*api.?key|unauthorized|401|authentication/i.test(text)) {
hints.push(
'API key authentication failed. Verify your ANTHROPIC_API_KEY is correct and has not expired.',
);
}
if (/ENOTFOUND|getaddrinfo/i.test(text)) {
hints.push(
'DNS resolution failed for the API endpoint. Check that your configured endpoint is correct.',
);
}
if (/certificate|CERT_|ssl|tls/i.test(text)) {
hints.push(
'TLS/SSL certificate error. Check the endpoint certificate chain and your local trust settings.',
);
}
}
// If no debug info available, provide generic Windows guidance
if (hints.length === 0) {
const isWin = process.platform === 'win32';
if (isWin) {
hints.push(
'Claude Code process crashed on Windows. Common fixes: ' +
'(1) Ensure ~/.claude.json exists: echo {} > %USERPROFILE%\\.claude.json ' +
'(2) Check your Claude authentication and endpoint configuration in ~/.claude/settings.json.',
);
} else {
return undefined;
}
}
return `${rawError}\n${hints.join('\n')}`;
}
/**
* Streaming chat endpoint.
* Routes to the appropriate provider SDK based on the `provider` field.
* Requires explicit provider and model; no fallback routing.
*/
export default defineEventHandler(async (event) => {
const body = await readBody<ChatBody>(event);
if (!body?.messages || body?.system == null) {
setResponseHeaders(event, { 'Content-Type': 'application/json' });
return { error: 'Missing required fields: system, messages' };
}
if (!body.provider) {
setResponseHeaders(event, { 'Content-Type': 'application/json' });
return { error: 'Missing provider. Provider fallback is disabled.' };
}
if (!body.model?.trim()) {
setResponseHeaders(event, { 'Content-Type': 'application/json' });
return { error: 'Missing model. Model fallback is disabled.' };
}
if (
body.provider !== 'anthropic' &&
body.provider !== 'openai' &&
body.provider !== 'opencode' &&
body.provider !== 'copilot' &&
body.provider !== 'gemini' &&
body.provider !== 'builtin'
) {
setResponseHeaders(event, { 'Content-Type': 'application/json' });
return { error: 'Missing or unsupported provider. Provider fallback is disabled.' };
}
setResponseHeaders(event, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
});
if (body.provider === 'builtin') return streamViaBuiltin(body);
if (body.provider === 'anthropic') return streamViaAgentSDK(body, body.model);
if (body.provider === 'opencode') return streamViaOpenCode(body, body.model);
if (body.provider === 'copilot') return streamViaCopilot(body, body.model);
if (body.provider === 'gemini') return streamViaGemini(body, body.model);
return streamViaCodex(body, body.model);
});
// Keep-alive ping interval (ms) — must stay below Bun's 10s idle timeout,
// but shouldn't be so aggressive that long-lived nested SSE streams create
// unnecessary write pressure on Bun dev.
const KEEPALIVE_INTERVAL_MS = 5_000;
function getAgentThinkingConfig(
body: ChatBody,
): { type: 'adaptive' | 'disabled' } | { type: 'enabled'; budgetTokens?: number } | undefined {
if (!body.thinkingMode) return undefined;
if (body.thinkingMode === 'enabled') {
return { type: 'enabled', budgetTokens: body.thinkingBudgetTokens };
}
return { type: body.thinkingMode };
}
/**
* Save base64 attachments to temp files. Returns { tempDir, files[] } — caller must clean up tempDir.
*
* When `insideProject` is true, files are saved under `.openpencil-tmp/` in the
* current working directory so that Claude Code Agent SDK (which restricts reads
* to the project directory in plan mode) can access them.
*/
async function saveAttachmentsToTempFiles(
attachments: ChatAttachmentWire[],
insideProject = false,
): Promise<{ tempDir: string; files: string[] }> {
let tempDir: string;
if (insideProject) {
const { mkdirSync, chmodSync } = await import('node:fs');
const baseDir = join(process.cwd(), '.openpencil-tmp');
mkdirSync(baseDir, { recursive: true, mode: 0o700 });
chmodSync(baseDir, 0o700);
tempDir = await mkdtemp(join(baseDir, 'attach-'));
} else {
tempDir = await mkdtemp(join(tmpdir(), 'openpencil-attach-'));
}
const files: string[] = [];
for (const att of attachments) {
const ext = resolveMediaExtension(att.mediaType);
const filePath = join(tempDir, `${files.length}.${ext}`);
await writeFile(filePath, Buffer.from(att.data, 'base64'));
files.push(filePath);
}
return { tempDir, files };
}
/** Collect all attachments from the last user message */
function getLastUserAttachments(body: ChatBody): ChatAttachmentWire[] {
const lastUser = [...body.messages].reverse().find((m) => m.role === 'user');
return lastUser?.attachments ?? [];
}
/**
* Strip "NEVER use tools" and similar instructions from system prompt
* when we need Claude Code Agent SDK to use its Read tool for image analysis.
*/
function stripNoToolsRestriction(systemPrompt: string): string {
return systemPrompt.replace(/^.*NEVER use tools.*$/gim, '').replace(/\n{3,}/g, '\n\n');
}
/** Stream via Claude Agent SDK (uses local Claude Code OAuth login, no API key needed) */
function streamViaAgentSDK(body: ChatBody, requestedModel?: string) {
let activeQuery: { close(): void } | undefined;
let cancelled = false;
const stream = new ReadableStream({
async start(controller) {
const encoder = new TextEncoder();
const safeEnqueue = (payload: Record<string, unknown>) => {
if (cancelled) return false;
try {
controller.enqueue(encoder.encode(`data: ${JSON.stringify(payload)}\n\n`));
return true;
} catch {
cancelled = true;
return false;
}
};
const safeClose = () => {
if (cancelled) return;
cancelled = true;
try {
controller.close();
} catch {
/* already closed */
}
};
// Keep emitting pings for the full stream lifetime. Some providers pause
// for >10s between text deltas, and Bun will otherwise kill the SSE socket.
const pingTimer = startSSEKeepAlive(() => {
safeEnqueue({ type: 'ping', content: '' });
}, KEEPALIVE_INTERVAL_MS);
let debugFile: string | undefined;
let attachTempDir: string | undefined;
try {
const { query } = await import('@anthropic-ai/claude-agent-sdk');
// Build prompt from the last user message
const lastUserMsg = [...body.messages].reverse().find((m) => m.role === 'user');
let prompt = lastUserMsg?.content ?? '';
// If the last user message has image attachments, save to temp files
// inside the project directory so Claude Code has read permission.
const attachments = getLastUserAttachments(body);
const hasImageAttachments = attachments.length > 0;
if (hasImageAttachments) {
const saved = await saveAttachmentsToTempFiles(attachments, true);
attachTempDir = saved.tempDir;
const imageRefs = saved.files
.map(
(f) =>
`First, use the Read tool to read the image file at "${f}". Then analyze it and respond to the user.`,
)
.join('\n');
prompt = imageRefs + '\n\n' + (prompt || 'Describe what you see in the image.');
}
// Remove CLAUDECODE env to allow running from within a CC terminal
const env = buildClaudeAgentEnv();
debugFile = getClaudeAgentDebugFilePath();
const model = requestedModel;
const claudePath = resolveClaudeCli();
const spawnProcess = buildSpawnClaudeCodeProcess();
const thinking = getAgentThinkingConfig(body);
// When images are attached, strip the "NEVER use tools" restriction from
// the system prompt so Claude Code will use its Read tool to view images.
const effectiveSystemPrompt = hasImageAttachments
? stripNoToolsRestriction(body.system)
: body.system;
// When images are attached, use result-based flow (like validate.ts):
// let Claude Code read the image via its Read tool internally, then
// only emit the final result text. This avoids streaming intermediate
// tool-use preamble like "I need to read the file first".
if (hasImageAttachments) {
const runImageQuery = async (): Promise<string> => {
const q = query({
prompt,
options: {
systemPrompt: effectiveSystemPrompt,
...(model ? { model } : {}),
maxTurns: 3,
plugins: [],
permissionMode: 'plan',
persistSession: false,
...(body.effort ? { effort: body.effort } : {}),
...(thinking ? { thinking } : {}),
env,
...(debugFile ? { debugFile } : {}),
...(claudePath ? { pathToClaudeCodeExecutable: claudePath } : {}),
...(spawnProcess ? { spawnClaudeCodeProcess: spawnProcess } : {}),
},
});
activeQuery = q;
try {
for await (const message of q) {
if (cancelled) return '';
if (message.type === 'result') {
const isErrorResult =
'is_error' in message && Boolean((message as { is_error?: boolean }).is_error);
if (message.subtype === 'success' && !isErrorResult) {
return message.result ?? '';
}
const errors = 'errors' in message ? (message.errors as string[]) : [];
const resultText = 'result' in message ? String(message.result ?? '') : '';
const errContent =
errors.join('; ') || resultText || `Query ended with: ${message.subtype}`;
throw new Error(errContent);
}
}
return '';
} finally {
activeQuery = undefined;
q.close();
}
};
const resultText = await runImageQuery();
if (resultText) {
safeEnqueue({ type: 'text', content: resultText });
}
} else {
// Normal text-only chat: stream partial messages as before
const runQuery = async () => {
const q = query({
prompt,
options: {
systemPrompt: effectiveSystemPrompt,
...(model ? { model } : {}),
maxTurns: 1,
includePartialMessages: true,
tools: [],
plugins: [],
permissionMode: 'plan',
persistSession: false,
...(body.effort ? { effort: body.effort } : {}),
...(thinking ? { thinking } : {}),
env,
...(debugFile ? { debugFile } : {}),
...(claudePath ? { pathToClaudeCodeExecutable: claudePath } : {}),
...(spawnProcess ? { spawnClaudeCodeProcess: spawnProcess } : {}),
},
});
activeQuery = q;
try {
for await (const message of q) {
if (cancelled) return;
if (message.type === 'stream_event') {
const ev = message.event;
if (ev.type === 'content_block_delta') {
if (ev.delta.type === 'text_delta') {
safeEnqueue({ type: 'text', content: ev.delta.text });
} else if (ev.delta.type === 'thinking_delta') {
safeEnqueue({
type: 'thinking',
content: (ev.delta as any).thinking,
});
}
}
} else if (message.type === 'result') {
const isErrorResult =
'is_error' in message && Boolean((message as { is_error?: boolean }).is_error);
if (message.subtype !== 'success' || isErrorResult) {
const errors = 'errors' in message ? (message.errors as string[]) : [];
const resultText = 'result' in message ? String(message.result ?? '') : '';
const content =
errors.join('; ') || resultText || `Query ended with: ${message.subtype}`;
safeEnqueue({ type: 'error', content });
}
}
}
} finally {
activeQuery = undefined;
q.close();
}
};
await runQuery();
}
safeEnqueue({ type: 'done', content: '' });
} catch (error) {
const rawContent = error instanceof Error ? error.message : 'Unknown error';
const tail = await readDebugTail(debugFile);
const hintedContent = buildClaudeExitHint(rawContent, tail);
// Append debug log tail so the user can see what Claude Code actually reported
let content = hintedContent ?? rawContent;
if (tail && tail.length > 0 && /process exited with code/i.test(rawContent)) {
const debugSnippet = tail.slice(-10).join('\n');
content += `\n\n[Debug log]:\n${debugSnippet}`;
}
safeEnqueue({ type: 'error', content });
} finally {
clearInterval(pingTimer);
try {
activeQuery?.close();
} catch {
/* ignore */
}
activeQuery = undefined;
if (attachTempDir) {
rm(attachTempDir, { recursive: true, force: true }).catch(() => {});
}
safeClose();
}
},
cancel() {
cancelled = true;
try {
activeQuery?.close();
} catch {
/* ignore */
}
activeQuery = undefined;
},
});
return new Response(stream);
}
/** Error name → user-friendly label mapping */
const OPENCODE_ERROR_LABELS: Record<string, string> = {
APIError: 'API error',
ProviderAuthError: 'Authentication failed',
UnknownError: 'Unknown error',
MessageOutputLengthError: 'Response too long',
MessageAbortedError: 'Request aborted',
StructuredOutputError: 'Output format error',
ContextOverflowError: 'Context too long',
};
/**
* Extract a human-readable message from an OpenCode error object.
* Handles structured errors like { name: "APIError", data: { message: "..." } }
* and nested JSON in message strings.
*/
export function formatOpenCodeError(error: unknown): string {
if (!error) return 'Unknown error';
if (typeof error === 'string') return error;
const err = error as Record<string, any>;
// Structured OpenCode error: { name, data: { message, ... } }
if (err.name && err.data?.message) {
const label = OPENCODE_ERROR_LABELS[err.name] ?? err.name;
let msg: string = err.data.message;
// Try to extract nested error message from JSON in the message string
// e.g. 'Unauthorized: {"error":{"code":"invalid_api_key","message":"invalid access token"}}'
const jsonStart = msg.indexOf('{');
if (jsonStart > 0) {
try {
const nested = JSON.parse(msg.slice(jsonStart));
const nestedMsg = nested?.error?.message ?? nested?.message;
if (nestedMsg) {
const prefix = msg.slice(0, jsonStart).replace(/:\s*$/, '').trim();
msg = prefix ? `${prefix}: ${nestedMsg}` : nestedMsg;
}
} catch {
/* not JSON, use as-is */
}
}
return `${label} — ${msg}`;
}
// Plain { message } object
if (err.message) return err.message;
// Fallback: truncated JSON
const json = JSON.stringify(error);
return json.length > 200 ? json.slice(0, 200) + '…' : json;
}
/** Parse an OpenCode model string ("providerID/modelID") into its parts */
function parseOpenCodeModel(model?: string): { providerID: string; modelID: string } | undefined {
if (!model || !model.includes('/')) return undefined;
const idx = model.indexOf('/');
return { providerID: model.slice(0, idx), modelID: model.slice(idx + 1) };
}
// Note: OpenCode SDK does not support `reasoning` in promptAsync/prompt params.
// The `reasoning` field was silently dropped by buildClientParams. Removed.
/** Wrap an async generator with a timeout — yields values until timeout fires */
async function* streamWithTimeout<T>(
stream: AsyncGenerator<T>,
timeoutPromise: Promise<{ done: true; value: undefined }>,
): AsyncGenerator<T> {
while (true) {
const result = (await Promise.race([stream.next(), timeoutPromise])) as IteratorResult<T>;
if (result.done) break;
yield result.value;
}
}
function streamViaCodex(body: ChatBody, model?: string) {
const stream = new ReadableStream({
async start(controller) {
const encoder = new TextEncoder();
const pingTimer = startSSEKeepAlive(() => {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'ping', content: '' })}\n\n`),
);
}, KEEPALIVE_INTERVAL_MS);
let attachTempDir: string | undefined;
try {
const lastUserMsg = [...body.messages].reverse().find((m) => m.role === 'user');
const prompt = lastUserMsg?.content ?? '';
// Save image attachments to temp files for Codex CLI
const attachments = getLastUserAttachments(body);
let imageFiles: string[] | undefined;
if (attachments.length > 0) {
const saved = await saveAttachmentsToTempFiles(attachments);
attachTempDir = saved.tempDir;
imageFiles = saved.files;
}
const result = await runCodexExec(prompt, {
model,
systemPrompt: body.system,
thinkingMode: body.thinkingMode,
thinkingBudgetTokens: body.thinkingBudgetTokens,
effort: body.effort,
imageFiles,
});
if (result.error) {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content: result.error })}\n\n`),
);
return;
}
if (result.text) {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'text', content: result.text })}\n\n`),
);
}
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'done', content: '' })}\n\n`),
);
} catch (error) {
const content = error instanceof Error ? error.message : 'Unknown error';
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content })}\n\n`),
);
} finally {
clearInterval(pingTimer);
if (attachTempDir) {
rm(attachTempDir, { recursive: true, force: true }).catch(() => {});
}
controller.close();
}
},
});
return new Response(stream);
}
/** Stream via OpenCode SDK using event subscription for real-time streaming */
function streamViaOpenCode(body: ChatBody, model?: string) {
const stream = new ReadableStream({
async start(controller) {
const encoder = new TextEncoder();
const pingTimer = startSSEKeepAlive(() => {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'ping', content: '' })}\n\n`),
);
}, KEEPALIVE_INTERVAL_MS);
let ocServer: { close(): void } | undefined;
try {
const { getOpencodeClient } = await import('../../utils/opencode-client');
const oc = await getOpencodeClient();
const ocClient = oc.client;
ocServer = oc.server;
// Create a session for this conversation
const { data: session, error: sessionError } = await ocClient.session.create({
title: 'OpenPencil Chat',
});
if (sessionError || !session) {
throw new Error(
`Failed to create OpenCode session: ${formatOpenCodeError(sessionError)}`,
);
}
// Inject system prompt as context (no AI reply)
const { error: sysPromptError } = (await ocClient.session.prompt({
sessionID: session.id,
noReply: true,
parts: [{ type: 'text', text: body.system }],
})) as any;
if (sysPromptError) {
console.error(
'[AI] OpenCode system prompt injection failed:',
formatOpenCodeError(sysPromptError),
);
}
// Build prompt from the last user message
const lastUserMsg = [...body.messages].reverse().find((m) => m.role === 'user');
const prompt = lastUserMsg?.content ?? '';
const parsed = parseOpenCodeModel(model);
if (model && !parsed) {
console.warn(
`[AI] OpenCode: could not parse model string "${model}", sending without model override`,
);
}
// Build parts array, adding image attachments if present
const attachments = getLastUserAttachments(body);
const parts: Array<Record<string, unknown>> = [
...attachments.map((a) => ({
type: 'image',
url: `data:${a.mediaType};base64,${a.data}`,
})),
{ type: 'text', text: prompt || 'Analyze these images.' },
];
// Build prompt payload with optional model and reasoning
const promptPayload: Record<string, unknown> = {
sessionID: session.id,
...(parsed ? { model: parsed } : {}),
parts,
};
// Subscribe to event stream for real-time deltas.
// IMPORTANT: The SSE connection is lazy — it only connects when
// iteration starts. We must start consuming BEFORE sending the
// prompt to avoid a race where events are emitted before the
// SSE connection is established.
const eventResult = await ocClient.event.subscribe();
const eventStream = eventResult.stream;
const sessionId = session.id;
const STREAM_TIMEOUT_MS = 180_000;
// Start eagerly consuming the event stream into a buffer.
// This triggers the SSE HTTP connection immediately.
const eventBuffer: unknown[] = [];
let streamDone = false;
let notifyFn: (() => void) | null = null;
const notify = () => {
if (notifyFn) {
const fn = notifyFn;
notifyFn = null;
fn();
}
};
// eslint-disable-next-line @typescript-eslint/no-floating-promises
void (async () => {
const timeoutPromise = new Promise<{ done: true; value: undefined }>((resolve) =>
setTimeout(() => resolve({ done: true, value: undefined }), STREAM_TIMEOUT_MS),
);
try {
for await (const event of streamWithTimeout(eventStream, timeoutPromise)) {
eventBuffer.push(event);
notify();
}
} finally {
streamDone = true;
notify();
}
})();
// Give the SSE connection a moment to establish before sending prompt
await new Promise<void>((resolve) => setTimeout(resolve, 100));
// Now send the prompt — SSE connection should already be active
const { error: asyncError } = await ocClient.session.promptAsync(promptPayload as any);
if (asyncError) {
const detail = formatOpenCodeError(asyncError);
console.error('[AI] OpenCode promptAsync error:', detail);
throw new Error(detail);
}
// Consume buffered events + wait for new ones
let emittedText = false;
let eventCount = 0;
let shouldBreak = false;
while (!shouldBreak) {
// Wait for events if buffer is empty
if (eventBuffer.length === 0) {
if (streamDone) break;
await new Promise<void>((resolve) => {
notifyFn = resolve;
});
continue;
}
const event = eventBuffer.shift();
if (!event || !('type' in (event as any))) continue;
const eventType = (event as any).type as string;
eventCount++;
// Stream text deltas for our session
if (eventType === 'message.part.delta') {
const props = (event as any).properties;
if (props?.sessionID === sessionId && props.field === 'text') {
const data = JSON.stringify({ type: 'text', content: props.delta });
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
emittedText = true;
}
// Forward reasoning deltas as thinking chunks
if (props?.sessionID === sessionId && props.field === 'reasoning') {
const data = JSON.stringify({ type: 'thinking', content: props.delta });
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
}
continue;
}
// Session went idle — response complete
if (eventType === 'session.idle') {
const props = (event as any).properties;
if (props?.sessionID === sessionId) {
shouldBreak = true;
}
continue;
}
// Session error
if (eventType === 'session.error') {
const props = (event as any).properties;
if (props?.sessionID === sessionId || !props?.sessionID) {
const errMsg = formatOpenCodeError(props?.error);
console.error('[AI] OpenCode session error:', errMsg);
const data = JSON.stringify({ type: 'error', content: errMsg });
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
shouldBreak = true;
}
continue;
}
}
// Fallback: if no text was streamed, try reading session messages directly
if (!emittedText) {
try {
const { data: messages } = (await ocClient.session.messages({
sessionID: sessionId,
})) as any;
if (messages && Array.isArray(messages)) {
// Find the last assistant message (each item has { info, parts })
const assistantMsg = [...messages]
.reverse()
.find((m: any) => m.info?.role === 'assistant');
if (assistantMsg?.parts) {
for (const part of assistantMsg.parts) {
if (part.type === 'text' && part.text) {
const data = JSON.stringify({ type: 'text', content: part.text });
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
emittedText = true;
}
}
}
}
} catch {
// fallback failed — will emit error below
}
}
if (!emittedText) {
const data = JSON.stringify({
type: 'error',
content:
'OpenCode returned an empty response. The model may not have generated any output.',
});
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
}
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'done', content: '' })}\n\n`),
);
} catch (error) {
const content = error instanceof Error ? error.message : 'Unknown error';
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content })}\n\n`),
);
} finally {
const { releaseOpencodeServer } = await import('../../utils/opencode-client');
releaseOpencodeServer(ocServer);
clearInterval(pingTimer);
controller.close();
}
},
});
return new Response(stream);
}
/** Map ChatBody effort to Copilot SDK ReasoningEffort */
function mapCopilotReasoningEffort(
effort?: 'low' | 'medium' | 'high' | 'max',
): 'low' | 'medium' | 'high' | 'xhigh' | undefined {
if (!effort) return undefined;
if (effort === 'max') return 'xhigh';
return effort;
}
/** Stream via Gemini CLI (`gemini -p -o stream-json`) — CLI handles its own auth */
function streamViaGemini(body: ChatBody, model?: string) {
const stream = new ReadableStream({
async start(controller) {
const encoder = new TextEncoder();
const pingTimer = startSSEKeepAlive(() => {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'ping', content: '' })}\n\n`),
);
}, KEEPALIVE_INTERVAL_MS);
try {
const { streamGeminiExec } = await import('../../utils/gemini-client');
// Build prompt from messages
const lastUserMsg = [...body.messages].reverse().find((m) => m.role === 'user');
const prompt = lastUserMsg?.content ?? '';
const { stream: geminiStream } = streamGeminiExec(prompt, {
model,
systemPrompt: body.system,
});
for await (const event of geminiStream) {
if (event.type === 'text') {
const data = JSON.stringify({ type: 'text', content: event.content });
try {
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
} catch {
/* stream closed */
}
} else if (event.type === 'error') {
const data = JSON.stringify({ type: 'error', content: event.content });
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
}
// 'done' is handled after loop
}
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'done', content: '' })}\n\n`),
);
} catch (error) {
const content = error instanceof Error ? error.message : 'Unknown error';
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content })}\n\n`),
);
} finally {
clearInterval(pingTimer);
controller.close();
}
},
});
return new Response(stream);
}
/** Stream via GitHub Copilot SDK (@github/copilot-sdk) */
function streamViaCopilot(body: ChatBody, model?: string) {
const stream = new ReadableStream({
async start(controller) {
const encoder = new TextEncoder();
const pingTimer = startSSEKeepAlive(() => {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'ping', content: '' })}\n\n`),
);
}, KEEPALIVE_INTERVAL_MS);
let copilotClient: { stop(): Promise<unknown> } | undefined;
try {
const { CopilotClient, approveAll } = await import('@github/copilot-sdk');
// Use standalone copilot binary to avoid Bun's node:sqlite issue
const { resolveCopilotCli, resolveCliPathForSdk } =
await import('../../utils/copilot-client');
const rawCliPath = resolveCopilotCli();
// On Windows, .cmd wrappers cause "spawn EINVAL" — resolve to .js entry point
const cliPath = rawCliPath ? resolveCliPathForSdk(rawCliPath) : undefined;
const client = new CopilotClient({
autoStart: true,
...(cliPath ? { cliPath } : {}),
});
copilotClient = client;
await client.start();
const session = await client.createSession({
...(model ? { model } : {}),
streaming: true,
onPermissionRequest: approveAll,
systemMessage: { mode: 'replace', content: body.system },
...(body.effort ? { reasoningEffort: mapCopilotReasoningEffort(body.effort) } : {}),
});
const lastUserMsg = [...body.messages].reverse().find((m) => m.role === 'user');
const prompt = lastUserMsg?.content ?? '';
// Subscribe to streaming deltas
session.on('assistant.message_delta', (event) => {
const deltaContent = (event as any).data?.deltaContent ?? '';
if (deltaContent) {
const data = JSON.stringify({ type: 'text', content: deltaContent });
try {
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
} catch {
/* stream closed */
}
}
});
// Wait for completion
await session.sendAndWait({ prompt }, 120_000);
await session.destroy();
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'done', content: '' })}\n\n`),
);
} catch (error) {
const content = error instanceof Error ? error.message : 'Unknown error';
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content })}\n\n`),
);
} finally {
clearInterval(pingTimer);
if (copilotClient) {
copilotClient.stop().catch(() => {});
}
controller.close();
}
},
});
return new Response(stream);
}
/**
* Stream via builtin provider — direct API key, no CLI tool needed.
* Uses Zig NAPI addon (agent-native) with Anthropic or OpenAI-compatible providers.
*/
function streamViaBuiltin(body: ChatBody) {
const stream = new ReadableStream({
async start(controller) {
const encoder = new TextEncoder();
const BUILTIN_EVENT_IDLE_TIMEOUT_MS = 45_000;
const pingTimer = startSSEKeepAlive(() => {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'ping', content: '' })}\n\n`),
);
}, KEEPALIVE_INTERVAL_MS);
try {
const {
createAnthropicProvider,
createOpenAICompatProvider,
createQueryEngine,
seedMessages,
submitMessage,
nextEvent,
abortEngine,
destroyIterator,
destroyQueryEngine,
destroyProvider,
} = await import('@zseven-w/agent-native');
const apiKey = body.builtinApiKey;
const rawModel = body.model?.trim() ?? '';
// Model string may be "builtin:<providerId>:<actualModel>" — extract the actual model name
const model = rawModel.startsWith('builtin:')
? rawModel.split(':').slice(2).join(':')
: rawModel;
if (!apiKey || !model) throw new Error('Builtin provider requires apiKey and model');
const normalizedBuiltinBaseURL = normalizeOptionalBaseURL(body.builtinBaseURL);
const builtinProvider =
body.builtinType === 'anthropic'
? createAnthropicProvider(apiKey, model, normalizedBuiltinBaseURL)
: createOpenAICompatProvider(
apiKey,
requireOpenAICompatBaseURL(normalizedBuiltinBaseURL),
model,
);
// Pure streaming — no tools, maxTurns=1 prevents agentic looping
const builtinEngine = createQueryEngine({
provider: builtinProvider,
systemPrompt: body.system,
maxTurns: 1,
maxOutputTokens: 16384,
cwd: process.cwd(),
});
// Seed prior conversation history for multi-turn context
const priorMsgs = body.messages
.slice(0, -1)
.filter(
(m: any) =>
(m.role === 'user' || m.role === 'assistant') && typeof m.content === 'string',
);
if (priorMsgs.length > 0) {
seedMessages(builtinEngine, JSON.stringify(priorMsgs));
}
const lastMsg = body.messages[body.messages.length - 1]?.content ?? '';
const builtinIter = await submitMessage(builtinEngine, lastMsg);
// Abort engine if no events arrive within 60s (provider sent 200 but no SSE data)
let gotFirstEvent = false;
const firstEventTimer = setTimeout(() => {
if (!gotFirstEvent) {
console.warn('[builtin] No SSE events received within 60s — aborting engine');
abortEngine(builtinEngine);
}
}, 60_000);
try {
let raw: string | null;
while (
(raw = await waitForBuiltinEvent(
nextEvent,
builtinIter,
() => abortEngine(builtinEngine),
BUILTIN_EVENT_IDLE_TIMEOUT_MS,
)) !== null
) {
if (!gotFirstEvent) {
gotFirstEvent = true;
clearTimeout(firstEventTimer);
}
const evt = JSON.parse(raw);
// Zig events are tagged unions: {"stream_event":{...}} or {"result":{...}}
const se = evt.stream_event;
if (se?.type === 'text_delta' && se.text) {
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'text', content: se.text })}\n\n`),
);
} else if (se?.type === 'thinking_delta' && se.text) {
controller.enqueue(
encoder.encode(
`data: ${JSON.stringify({ type: 'thinking', content: se.text })}\n\n`,
),
);
} else if (evt.result?.is_error) {
// Zig attaches the provider's last_error string in result.errors[0]
// (e.g. "Content blocked by provider safety filter (HTTP 451)...").
// Surface that instead of the opaque subtype so users see the
// actual reason — "content blocked" vs. "rate limit" vs. "auth"
// is information they can act on.
const detail =
(Array.isArray(evt.result.errors) && evt.result.errors[0]) ||
evt.result.subtype ||
'unknown';
const errMsg = `Provider error: ${detail}`;
console.error('[builtin]', errMsg);
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content: errMsg })}\n\n`),
);
}
}
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'done', content: '' })}\n\n`),
);
} finally {
clearTimeout(firstEventTimer);
destroyIterator(builtinIter);
destroyQueryEngine(builtinEngine);
destroyProvider(builtinProvider);
}
} catch (error) {
const content = error instanceof Error ? error.message : 'Unknown error';
controller.enqueue(
encoder.encode(`data: ${JSON.stringify({ type: 'error', content })}\n\n`),
);
} finally {
clearInterval(pingTimer);
controller.close();
}
},
});
return new Response(stream);
}
async function waitForBuiltinEvent<TIterator>(
nextEventFn: (iter: TIterator) => Promise<string | null>,
iter: TIterator,
onTimeout: () => void,
timeoutMs: number,
): Promise<string | null> {
return await new Promise<string | null>((resolve, reject) => {
const timer = setTimeout(() => {
try {
onTimeout();
} catch {
/* ignore */
}
reject(new Error('Builtin provider stalled without output. Please retry.'));
}, timeoutMs);
nextEventFn(iter)
.then((value) => {
clearTimeout(timer);
resolve(value);
})
.catch((err) => {
clearTimeout(timer);
reject(err);
});
});
}