import { query as sdkQuery } from "@anthropic-ai/claude-agent-sdk"; import { type WorktreeSetupSpec } from "./git-worktree.js"; import { SessionLogger } from "./session-logger.js"; import { WandStorage } from "./storage.js"; import { ExecutionMode, ProcessEvent, SessionProvider, SessionRunner, SessionSnapshot, SessionSource, WandConfig } from "./types.js"; import { type StructuredExecHost } from "./structured-exec-host.js"; import type { StructuredRunnerAdapter } from "./structured-runner.js"; export interface StructuredSessionManagerRunners { claudeCli?: StructuredRunnerAdapter; codex?: StructuredRunnerAdapter; opencode?: StructuredRunnerAdapter; grok?: StructuredRunnerAdapter; qoder?: StructuredRunnerAdapter; pi?: StructuredRunnerAdapter; } interface CreateStructuredSessionOptions { cwd: string; mode: ExecutionMode; provider?: SessionProvider; runner?: SessionRunner; worktreeEnabled?: boolean; worktreeSpec?: WorktreeSetupSpec; /** 用户指定的模型(别名或完整 ID)。留空则 spawn 时不加 --model。 */ model?: string; /** 用户预设的思考深度。留空 / null 视为 off。 */ thinkingEffort?: SessionSnapshot["thinkingEffort"]; sessionSource?: SessionSource; automationId?: string; /** 所属工作空间 ID(多标签 / 分屏项目)。 */ workspaceId?: string; /** 所属工作空间任务 ID(任务 = 独立 worktree + 一组标签)。 */ workspaceTaskId?: string; /** * 恢复用的初始会话 id: * - Codex:历史 thread id,首条消息即 `codex exec ... resume ` 续接。 * - Claude:历史 session id,首条消息即 `--resume` / SDK resume 续接。 * 留空表示新建会话。 */ claudeSessionId?: string; } /** * 返回最近一次真正提交给结构化会话的用户输入。 * * 排队非空时,队尾才是“上一条提交”;否则回看当前正在处理的最后一个 user turn。 * 这里只接受可无损还原成字符串的 text / tool_result,避免把图片等结构化内容误判 * 成更早的纯文本输入。 */ export declare function getLastSubmittedStructuredInput(snapshot: Pick): string | null; /** 仅用于 in-flight 排队分支:连续两次内容相同则把后一次视为输入重放。 */ export declare function isDuplicateStructuredQueueInput(snapshot: Pick, input: string): boolean; export declare class StructuredSessionManager { private readonly storage; private readonly config; private readonly logger; private readonly sdkQueryFactory; private readonly execHost?; private readonly sessions; private readonly pendingRunnerExecutions; private readonly pendingSdkAbort; /** * Active SDK Query handle per session, kept around so we can call * `query.interrupt()` for a graceful stop instead of aborting via signal. * Only populated while an SDK call is in flight. */ private readonly pendingSdkQueries; private readonly interruptedWith; private readonly interruptedSkills; /** * Sessions where the current interrupt is a "queue promote" (用户从排队条点了「立即」 * 把队首插队到 now)。退出处理三个分支默认会把 queuedMessages 清空——因为常规的 * interrupt 语义是"算了,做这个",把队列也作废。但 queue-promote 的语义是 * "先做这条,剩下的队列还要继续",所以这里打个标记,让退出 handler 保留 queue。 * 收到后必须 delete 掉,避免下一次普通 interrupt 误带 flag。 */ private readonly preserveQueueOnInterrupt; /** Last wall-clock time (ms) a streaming checkpoint reached SQLite. */ private readonly lastStreamSaveAt; private readonly streamCheckpointTimers; private readonly streamCheckpointDirty; /** * Idempotency keys we've already accepted, mapped to their wall-clock timestamp. * Android WebView 在进程恢复时偶尔会重发上一个未收到响应的 POST(HTTP/2 stream * reset 等场景),客户端 JS 没有重试逻辑也拦不住。这里用 (sessionId, key) 永 * 久去重,重复就抛错让前端弹 toast 提示,**不**做任何处理。timestamp 仅用于 * map 大小溢出时按时间裁剪。 */ private readonly seenIdempotencyKeys; private emitEvent; private archiveTimer; private readonly topicCoordinator; private readonly streamEmitTimers; private readonly claudeCliRunner; private readonly codexRunner; private readonly openCodeRunner; private readonly grokRunner; private readonly qoderRunner; private readonly piRunner; /** Structured CLI runs that were mid-flight when the previous web process died. */ private pendingRecoveryIds; private disposed; constructor(storage: WandStorage, config: WandConfig, logger?: SessionLogger | null, sdkQueryFactory?: typeof sdkQuery, runners?: StructuredSessionManagerRunners, execHost?: StructuredExecHost | undefined); private archiveExpiredSessions; setEventEmitter(emitEvent: (event: ProcessEvent) => void): void; /** Stop every runner and flush terminal state before storage is closed. */ dispose(): void; /** Called once after startup wiring; safe to skip when no host or candidates. */ recoverDetachedRuns(): Promise; private resumeDetachedRun; private finalizeRecoveredRun; private trackStreamEmitTimer; private clearStreamEmitTimer; /** Mark streaming payload dirty and enforce both leading and trailing checkpoints. */ private saveStreamingSnapshot; private flushStreamingCheckpoint; private clearStreamingCheckpoint; private cancelStreamingCheckpointTimer; private saveAuthoritativeSession; private checkpointSessionMessages; list(): SessionSnapshot[]; /** Return lightweight snapshots for the session list (no output/messages). */ listSlim(): SessionSnapshot[]; get(id: string): SessionSnapshot | null; setSessionTopic(id: string, title: string, description: string): SessionSnapshot; private setSessionTopicGenerating; /** * Update worktree merge progress on the canonical in-memory snapshot before * persisting it. A null result means this manager does not own the session. */ setWorktreeMergeState(id: string, status: SessionSnapshot["worktreeMergeStatus"], info: SessionSnapshot["worktreeMergeInfo"]): SessionSnapshot | null; private maybeGenerateSessionTopic; createSession(options: CreateStructuredSessionOptions): SessionSnapshot; sendMessage(id: string, input: string, opts?: { interrupt?: boolean; idempotencyKey?: string; preserveQueue?: boolean; queueAlreadyRemoved?: boolean; skills?: string[]; }): Promise; /** * Reorder the pending queued messages. `order` is a permutation of the current * indices, e.g. `[2, 0, 1]` means "move the third queued message to the front, * push the original first to position #2". Throws if the permutation is * malformed (length mismatch / duplicate / out-of-range). 不允许在 inFlight * 期间改"已经被 flushNextQueuedMessage 拿走的队首",但本方法只动 queue 数组 * 本身,flushNext 在另一段时序里读 sessions.get(...) 当前快照,已经天然安全。 */ reorderQueuedMessages(sessionId: string, order: number[]): SessionSnapshot; /** Remove a single queued message only while index and text still identify the same item. */ deleteQueuedMessage(sessionId: string, index: number, expectedText?: string): SessionSnapshot; /** * Remove one queued message by index before sending it. Keeping this operation * on the server prevents clients from re-sending the text while the original * queue entry remains available for the automatic flush path. */ promoteQueuedMessage(sessionId: string, index: number, expectedText?: string, idempotencyKey?: string): Promise; /** Clear all queued messages. No-op when queue is already empty. */ clearQueuedMessages(sessionId: string): SessionSnapshot; /** Update the selected model for a structured session. Takes effect on the next spawn. */ setSessionModel(sessionId: string, model: string | null): SessionSnapshot; /** * Update the thinking-effort level for a structured session. Takes effect on * the next spawn / next message (SDK runner injects `thinking`, Claude CLI * runner passes `--effort`, codex runner overrides `model_reasoning_effort`). */ setSessionThinkingEffort(sessionId: string, effort: SessionSnapshot["thinkingEffort"]): SessionSnapshot; /** * Switch the execution mode of a structured session mid-flight. Takes effect on * the next message/query — permission policy, append-system-prompt and CLI flags * are all re-derived from session.mode per turn. Mirrors setSessionModel; also * re-syncs autoApprovePermissions so the permission posture matches the new mode. */ setSessionMode(sessionId: string, mode: ExecutionMode): SessionSnapshot; /** Toggle auto-approve for the session. */ toggleAutoApprove(sessionId: string): SessionSnapshot; /** Resolve a specific escalation by requestId. */ resolveEscalation(sessionId: string, requestId: string, resolution: unknown): SessionSnapshot; stop(id: string): SessionSnapshot; delete(id: string): void; private requireSession; /** True only while this exact turn still owns the session's mutable state. */ private isCurrentRequest; private currentSessionForRequest; /** Delete a handle only if it still belongs to the execution doing cleanup. */ private releasePendingRunnerExecution; private releasePendingSdkAbort; private releasePendingSdkQuery; private emitStructuredSnapshot; private flushNextQueuedMessage; private emit; private incrementApprovalStats; private runCodexStreaming; private runGrokStreaming; private runOpenCodeStreaming; /** * Spawn `claude -p --output-format stream-json` and parse NDJSON lines as * they arrive, emitting incremental WebSocket events so the UI can render * text / thinking / tool_use blocks in real-time. * * Permission handling: * - Non-root + full-access/managed: --permission-mode bypassPermissions * - Non-root + auto-edit: --permission-mode acceptEdits * - Root: --permission-mode acceptEdits + --allowedTools (extends approval * outside CWD). stdin is always "ignore" — no ACP bidirectional control. */ private runClaudeStreaming; /** * Use @anthropic-ai/claude-agent-sdk instead of spawning claude -p directly. * The SDK still spawns the claude binary but provides typed AsyncGenerator * messages, so we skip NDJSON parsing. Options are 1:1 with the CLI flags. * * Streaming is enabled via includePartialMessages: true — the SDK emits * SDKPartialAssistantMessage (type: "stream_event") with BetaRawMessageStreamEvent * payloads for incremental text/thinking/tool_use updates, followed by a final * SDKAssistantMessage with the authoritative complete content. */ private runClaudeSdkStreaming; private compactContentBlocks; private buildCompletedAssistantMessages; private resolveQueuedMessagesAfterInterrupt; private resolveQueuedMessageSkillsAfterInterrupt; private normalizeToolResultContent; /** * 组装结构化 runner 退出失败时的可读错误字符串。 * * 痛点:之前 claude -p / codex exec 异常退出只把"stderr.trim() || `... exited * with code N`"塞给 UI。如果 stderr 是空的,用户在前端只能看到 "EXIT 1" 这种 * 没有任何上下文的串,根本不知道是网络错误、参数错误还是 binary 找不着。 * * 这里固定把"provider + 退出码 / 信号"放在最前面,再把 stderr / NDJSON 错误 * 事件 / 最后一段 stdout 之类的上下文跟在后面,方便定位。 */ private formatStructuredExitError; private finishStructuredFailure; /** Extract usage from an SDKResultSuccess message (sdk runner). */ private extractSdkUsage; } export {};