import { query } from "@anthropic-ai/claude-agent-sdk"; import { asError } from "../../utils/errors"; import { randomUUID } from "crypto"; import { existsSync } from "fs"; import { join } from "path"; import { homedir } from "os"; import type { AgentBackend, AgentSession, AgentSessionContext, AgentEvent } from "../types"; import type { Attachment } from "../../types/attachment"; import { SdkNormalizer } from "./claude-normalize"; import { MessageStream } from "../message-stream"; import { getSdkSkillsSetting } from "../../core/skills"; import { getSdkHooks } from "../../core/sdk-hooks"; import { getConfig } from "../../utils/config"; import { resolveClaudeCredential, credentialEnv } from "../credentials"; import { sleep } from "../../utils/retry"; /** The shape of the SDK `query()` handle the session consumes. Injected so the * session is unit-testable without spawning Claude. */ export type QueryHandle = AsyncIterable & { close(): void }; export type QueryFn = (args: { prompt: unknown; options: unknown }) => QueryHandle; const MAX_SEND_RETRIES = 2; const DEFAULT_RETRY_DELAYS = [3_000, 8_000]; /** The SDK persists sessions at ~/.claude/projects//.jsonl. */ function sessionFileExists(sessionId: string, cwd: string): boolean { const encoded = cwd.replace(/\//g, "-"); return existsSync(join(homedir(), ".claude", "projects", encoded, `${sessionId}.jsonl`)); } /** Resolve a context/config model to the SDK's `model` option ("default" → unset). */ export function resolveSdkModel(model?: string | null): string | undefined { const m = model || getConfig().model; return m && m !== "default" ? m : undefined; } export class ClaudeBackend implements AgentBackend { readonly name = "claude" as const; private queryFn: QueryFn; constructor(deps?: { queryFn?: QueryFn }) { this.queryFn = deps?.queryFn ?? (query as unknown as QueryFn); } async openSession(ctx: AgentSessionContext): Promise { return new ClaudeSession(ctx, this.queryFn); } async canResume(backendSessionId: string, cwd: string): Promise { return sessionFileExists(backendSessionId, cwd); } } /** * A warm Claude session: one `query()` subprocess + `MessageStream` reused * across turns (the latency optimization). Each `send()` pushes a turn and * yields its normalized events until a terminal `result`/`error`. * * Invariants (from the plan review): * - exactly ONE `session` event per `send()`, even across an internal retry * (a retry resumes the same session id, so the post-retry init is swallowed); * - retry teardown+restart is internal and atomic w.r.t. `abort()`. */ class ClaudeSession implements AgentSession { private _sessionId: string | null; private handle: QueryHandle | null = null; private iterator: AsyncIterator | null = null; private stream: MessageStream | null = null; private aborted: string | null = null; private retryCount = 0; private readonly retryDelays: number[]; constructor( private ctx: AgentSessionContext & { retryDelaysMs?: number[] }, private queryFn: QueryFn, ) { this._sessionId = typeof ctx.resume === "string" ? ctx.resume : null; this.retryDelays = ctx.retryDelaysMs ?? DEFAULT_RETRY_DELAYS; } get backendSessionId(): string | null { return this._sessionId; } private startQuery(): void { this.stream = new MessageStream(); const options: Record = { systemPrompt: this.ctx.systemPrompt, cwd: this.ctx.cwd, permissionMode: "bypassPermissions", skills: getSdkSkillsSetting(), hooks: getSdkHooks(), }; // Interactive (chat) sessions stream partials and load project/user settings; // headless one-shot jobs keep the leaner option set they had pre-refactor. if (this.ctx.interactive) { options.includePartialMessages = true; options.settingSources = ["project", "user"]; } const model = resolveSdkModel(this.ctx.model); if (model) options.model = model; if (this._sessionId) { options.resume = this._sessionId; } else { options.sessionId = randomUUID(); // Interactive sessions also forbid auto-continue of a prior session in the // same cwd; jobs always run with a unique id and never auto-continued. if (this.ctx.interactive) options.continue = false; } // Hand the CLI a credential Nia owns when one is configured. Without this // it inherits ~/.claude/.credentials.json, which only refreshes when a // human runs `claude` on this machine — the coupling that had Nia // answering as codex for sixteen days. const credential = resolveClaudeCredential(getConfig()); if (credential.envVar) { options.env = credentialEnv(credential, process.env as Record); } if (this.ctx.outputSchema) options.outputFormat = { type: "json_schema", schema: this.ctx.outputSchema }; if (this.ctx.mcpServers) options.mcpServers = this.ctx.mcpServers; if (this.ctx.subagents && Object.keys(this.ctx.subagents).length > 0) options.agents = this.ctx.subagents; this.handle = this.queryFn({ prompt: this.stream, options }); this.iterator = this.handle[Symbol.asyncIterator](); } async *send(text: string, attachments?: Attachment[]): AsyncIterable { let sawSession = false; while (true) { if (!this.iterator || !this.stream) this.startQuery(); this.stream!.push(text, attachments); const normalizer = new SdkNormalizer(resolveClaudeCredential(getConfig()).kind); let retry = false; while (true) { let res: IteratorResult; try { res = await this.iterator!.next(); } catch (err) { if (this.aborted) throw new Error(this.aborted); throw asError(err); } if (this.aborted) throw new Error(this.aborted); if (res.done) { if (this.aborted) throw new Error(this.aborted); throw new Error("stream ended without result"); } for (const ev of normalizer.consume(res.value)) { if (ev.type === "session") { this._sessionId = ev.backendSessionId; if (!sawSession) { sawSession = true; yield ev; } continue; } if (ev.type === "error" && ev.retryable && this.retryCount < MAX_SEND_RETRIES) { this.retryCount++; yield { type: "thinking", delta: "retrying after API error..." }; await this.teardown(); await sleep(this.retryDelays[this.retryCount - 1] ?? 8_000); retry = true; break; } // A retryable error that survived every retry is a capacity problem // with THIS model, not proof the provider is gone — scope it to the // model so a second model on the same provider still gets its turn. const out = ev.type === "error" && ev.retryable && this.retryCount >= MAX_SEND_RETRIES ? { ...ev, failover: ev.failover ?? ("model" as const) } : ev; yield out; if (out.type === "result" || out.type === "error") { this.retryCount = 0; return; } } if (retry) break; // restart the outer loop: startQuery() resumes the same session id } } } abort(reason: string): void { this.aborted = reason; this.handle?.close(); } private async teardown(): Promise { this.stream?.end(); this.handle?.close(); this.stream = null; this.handle = null; this.iterator = null; } async close(): Promise { await this.teardown(); } }