Skip to content

Latest commit

 

History

History
392 lines (317 loc) · 14.9 KB

File metadata and controls

392 lines (317 loc) · 14.9 KB

第 3 章 Query Loop:对话主循环

一句"帮我修复这个 bug"从你按下回车开始,到 Claude 写完代码并说"修好了"为止,背后是上百次 stream event 和工具调用。本章拆开这条 Turn Loop——以源码行号为准。

3.1 三层职责分工

CCB 把"对话循环"切成三层,每层各管一件事:

文件 角色 行数
UI 层 src/screens/REPL.tsx 渲染消息、捕获输入、弹工具权限确认 6541
编排层 src/QueryEngine.ts 会话状态机、压缩、归因、turn 边界 1365
协议层 src/query.ts 单 turn loop:发请求 → 流式接收 → 执行工具 → 决定要不要再来一轮 2042

下游还有 src/services/api/claude.ts(3558 行)负责真正构造 HTTP 请求和解析 SSE,那是第 5 章。本章只关心消息怎么从 REPL 走到 query() 再走回 UI

graph TB
    User[用户输入]
    PromptInput[PromptInput<br/>组件]
    REPL[REPL.tsx]
    QE["QueryEngine.submitMessage()<br/>line 217"]
    ASK["ask() 顶级函数<br/>line 1256"]
    Q[query.ts query()]
    QL[queryLoop 主循环]
    API[claude.ts queryModel<br/>line 1039]
    Trans[transitions.ts<br/>Terminal=11 / Continue=7]

    User --> PromptInput --> REPL
    REPL -->|新 prompt| QE
    REPL -->|或一次性场景| ASK
    ASK -->|内部用| QE
    QE --> Q --> QL
    QL --> Trans
    QL -->|每次迭代| API
    API -->|stream events| QL
    QL -->|tool_use| Tool[Tool.call → ToolResult]
    Tool -->|tool_result| QL
    QL -->|loop end| QE
    QE -->|message stream| REPL
    REPL -->|render| User
Loading

⚠️ 重要更正:QueryEngine 的主入口submitMessage() 不是 ask()ask()QueryEngine.ts:1256 的一个顶级 helper function,它内部 yield* engine.submitMessage(...),供一次性场景使用(如 print 模式)。

3.2 Terminal 和 Continue:单独的类型文件

第一版书我以为 Terminalquery.ts 里。实际上它和 Continue 都在独立的小文件 src/query/transitions.ts

// src/query/transitions.ts:1-11
export type Terminal =
  | { reason: 'completed' }
  | { reason: 'blocking_limit' }
  | { reason: 'image_error' }
  | { reason: 'model_error'; error?: unknown }
  | { reason: 'aborted_streaming' }
  | { reason: 'aborted_tools' }
  | { reason: 'prompt_too_long' }
  | { reason: 'stop_hook_prevented' }
  | { reason: 'hook_stopped' }
  | { reason: 'max_turns'; turnCount: number }

11 种 Terminal(不是我之前写的"end_turn 等几种")。completed 是正常结束(注意不是 end_turn),其余是各种异常退出。

// src/query/transitions.ts:13-20
export type Continue =
  | { reason: 'collapse_drain_retry'; committed: number }
  | { reason: 'reactive_compact_retry' }
  | { reason: 'max_output_tokens_escalate' }
  | { reason: 'max_output_tokens_recovery'; attempt: number }
  | { reason: 'stop_hook_blocking' }
  | { reason: 'token_budget_continuation' }
  | { reason: 'next_turn' }

7 种 Continue——每一种都对应 queryLoop 内部的一个具体重试场景:

Continue 原因 触发场景
next_turn 正常的下一轮(执行完 tool_result 后继续)
max_output_tokens_recovery 模型回复被截断,加倍 maxOutputTokensOverride 重试(≤ 3 次)
max_output_tokens_escalate recovery 仍不够,进一步升级
reactive_compact_retry 收到 prompt_too_long 错误,紧急压缩后重试
collapse_drain_retry Context collapse 排水重试(feature CONTEXT_COLLAPSE)
stop_hook_blocking Stop hook 阻止结束,强制再来一轮
token_budget_continuation API task_budget 没用完,继续

3.3 query() 与 queryLoop 的关系

// src/query.ts:275
export async function* query(
  params: QueryParams,
): AsyncGenerator<
  StreamEvent | RequestStartEvent | Message | TombstoneMessage | ToolUseSummaryMessage,
  Terminal
>

query() 是个包装器:负责创建 / 复用 Langfuse trace、try/finally 清理(详见 3.6)、最后调用 queryLoop()

// src/query.ts:392
async function* queryLoop(
  params: QueryParams,
  consumedCommandUuids: string[],
  consumedAutonomyCommands: QueuedCommand[],
)

queryLoop 是真正的循环。

queryLoop 内部状态(line 420-431)

let state: State = {
  messages: params.messages,
  toolUseContext: params.toolUseContext,
  maxOutputTokensOverride: params.maxOutputTokensOverride,
  autoCompactTracking: undefined,
  stopHookActive: undefined,
  maxOutputTokensRecoveryCount: 0,
  hasAttemptedReactiveCompact: false,
  turnCount: 1,
  pendingToolUseSummary: undefined,
  transition: undefined,
}

transition 字段记录上次为什么进入这一轮——让测试可以断言"recovery 路径确实触发了"。

迭代主体:while(true) + 7 个 continue(line 459-2042)

每次迭代的伪流程:

flowchart TD
    Start[迭代开始<br/>destructure state] --> Build[构造 messages + system + tools]
    Build --> Send[调 claude.ts queryModel]
    Send --> Stream{逐 SSE chunk}
    Stream -->|content_block_delta text| YieldText[yield 文本片段 → UI]
    Stream -->|content_block_delta thinking| YieldThink[yield thinking → UI]
    Stream -->|tool_use 流式完成| StreamExec{StreamingToolExecutor<br/>开启?}
    StreamExec -->|是| ParallelStart[并行启动 tool 执行]
    StreamExec -->|否| Collect[收集 tool_use blocks]
    Stream -->|stream end stop_reason| StopReason{stop_reason?}
    StopReason -->|end_turn| Done[Terminal completed]
    StopReason -->|tool_use| Exec[runTools 串行 / StreamingExec 收尾]
    Exec -->|结果| Append[append tool_result 到 messages]
    Append -->|maxTurns 到顶| Stop[Terminal max_turns]
    Append -->|stop_hook 阻止| StopHook[Continue stop_hook_blocking]
    Append -->|预测 token 超| AutoCompact[autocompact → continue]
    Append --> Next[Continue next_turn]
    StopReason -->|max_tokens| Recovery{recovery 次数 < 3?}
    Recovery -->|是| RecoveryRetry[Continue max_output_tokens_recovery]
    Recovery -->|否| TerminalRecover[Terminal blocking_limit]
    Send -->|API 抛 prompt_too_long| Reactive{已尝试过 RC?}
    Reactive -->|否| RC[Continue reactive_compact_retry]
    Reactive -->|是| TerminalLong[Terminal prompt_too_long]
Loading

3.4 并行 tool 执行:StreamingToolExecutor

第一版我写"isReadOnly && isConcurrencySafe 就并行"是猜的。真实机制

// src/query.ts:744-751
const streamingToolExecutor = config.gates.streamingToolExecution
  ? createStreamingToolExecutor(...)
  : null

streamingToolExecutorfeature gate 控制(不是简单的工具属性判定):

  • 启用时:模型流式发出 tool_use block 的瞬间就启动执行(query.ts:1084-1087),不等整个 stream 结束。完成的 tool_result 立即 yield 给 UI(line 1094-1105)。
  • 未启用时:fallback 到 runTools()(line 1656),等流彻底停止后串行执行。

这是个真实的性能优化:流式执行能让用户在模型还在生成的同时看到工具开始跑(如 Bash 命令的 stdout)。

3.5 yieldMissingToolResultBlocks:流中断的补救(line 148-178)

function* yieldMissingToolResultBlocks(
  ...
)

作用:流被错误中断时,对每个孤立的 tool_use 块(没有对应 tool_result)生成一条 is_error: true 的 tool_result。

调用点

  • query.ts:1143 —— fallback 失败时
  • query.ts:1230 —— API 异常
  • query.ts:1297 —— 其他错误路径

不补这一刀,消息历史里会出现"模型说要调 X 但没结果",下次再发会被 Anthropic API 拒绝。

3.6 query() 的 finally:JSC 性能修复

// src/query.ts (finally 块内)
// Clear JSC's native Performance buffers. OTel (otperformance) references
// globalThis.performance which stores marks/measures/resource timings in a
// C++ Vector that never shrinks. Long-running sessions accumulate hundreds
// of MB of dead capacity even after spans are flushed and nullified.
const gPerf = globalThis.performance
if (gPerf && typeof gPerf.clearMarks === 'function') {
  try {
    gPerf.clearMarks()
    gPerf.clearMeasures?.()
    gPerf.clearResourceTimings?.()
  } catch {}
}

这跟第 2 章讲的 performanceShim.ts配合关系——shim 把 mark/measure 重定向到 JS Map,但 OTel/Langfuse 还是会写 marks,每次 turn 结束必须 clearMarks 主动释放

另外还会切闭包链:

// src/query.ts (finally 块内)
if (paramsWithTrace !== params) {
  paramsWithTrace.toolUseContext.langfuseTrace = null
  paramsWithTrace.toolUseContext.langfuseRootTrace = null
  paramsWithTrace.toolUseContext.langfuseBatchSpan = null
}

阻止 toolUseContext 持续引用 SpanImpl 树(每个 span 可能挂着几百 KB 序列化的对话历史)。

3.7 QueryEngine:把 query() 包成"会话"

// src/QueryEngine.ts:192
export class QueryEngine {
  // ...
  // src/QueryEngine.ts:208 constructor
  constructor(config: QueryEngineConfig)
  
  // src/QueryEngine.ts:217 主入口
  async *submitMessage(
    prompt: string | ContentBlockParam[],
    options?: { ... }
  )
}

// src/QueryEngine.ts:1256 顶级 helper(一次性场景)
export async function* ask({ ... })

QueryEngineConfig(line 138-181)真实字段:

export type QueryEngineConfig = {
  cwd: string
  tools: Tools
  commands: Command[]
  mcpClients: MCPServerConnection[]
  agents: AgentDefinition[]
  canUseTool: CanUseToolFn
  getAppState: () => AppState
  setAppState: (f: (prev: AppState) => AppState) => void
  initialMessages?: Message[]
  readFileCache: FileStateCache
  customSystemPrompt?: string
  appendSystemPrompt?: string
  userSpecifiedModel?: string
  fallbackModel?: string
  thinkingConfig?: ThinkingConfig
  maxTurns?: number
  maxBudgetUsd?: number
  taskBudget?: { total: number }
  jsonSchema?: Record<string, unknown>
  verbose?: boolean
  replayUserMessages?: boolean
  handleElicitation?: ToolUseContext['handleElicitation']
  includePartialMessages?: boolean
  setSDKStatus?: (status: SDKStatus) => void
  abortController?: AbortController
  orphanedPermission?: OrphanedPermission
  snipReplay?: (...) => ... | undefined
}

它做的事情(按职责):

  1. 维持 messages 数组在多 turn 间累积
  2. 文件历史快照(line 53-56 imports,line 653-662 实际调用)—— fileHistoryMakeSnapshot()Edit 工具的 oldString diff 校验;
  3. commit attribution(line 47 imports,line 395-398 更新回调)—— git commit 时拼接 Co-Authored-By
  4. 跨 turn autoCompact tracking(line 1374 持久化在 state.messages 中);
  5. subagent 调度(用 runAgent.ts 配合)。

query() 调用点(line 688-699)

for await (const message of query({
  messages, systemPrompt, userContext, systemContext,
  canUseTool: wrappedCanUseTool,
  toolUseContext: processUserInputContext,
  fallbackModel, querySource: 'sdk', maxTurns, taskBudget
}))

注意 canUseTool包装过的——submitMessage 在外面套一层钩子做权限校验,再传给 query()

3.8 REPL.tsx:6541 行的"消费方"

// src/screens/REPL.tsx:828
export function REPL({ ... 25+  props ... }: Props): React.ReactNode

它通过 useMergedTools()(line 272)、useCanUseTool()(line 207)等 hook 接入 QueryEngine。在 line 1165 有个 QueryGuard 实例(防止并发 turn),在 line 1616 处理 earlyInput 回放(第 2 章讲过)。

真正消费 query 事件流的 for await 循环在 line 3403+:

for await (const event of query({ ... })) {
  // 根据 event.type 增量 setState
}

事件类型清单:

事件 时机 UI 反应
RequestStartEvent HTTP 已发出 显示 "Thinking..." spinner
message_start 模型开始回复 替换 spinner 为打字效果
content_block_delta (text_delta) 一段文本增量 逐字渲染
content_block_delta (thinking_delta) thinking 增量 折叠区域内逐字渲染
content_block_delta (input_json_delta) tool_use 参数增量 进度条 / 参数预览
content_block_stop 块完整 tool_use 块可立即询问权限
message_stop 整条 assistant 消息结束 显示 stop_reason
Message (yielded from queryLoop) 完整 user/assistant/tool_result 加入 Messages 列表
TombstoneMessage 被 compact 掉的消息 替换为"已压缩"占位

3.9 Task.ts ≠ task management 工具

第一版我误把 src/Task.ts 写成"任务管理"。实际它定义的是 后台任务(background task) 的类型,跟 TaskCreateTool / TodoWrite 等用户级 task 是两套东西:

// src/Task.ts (示意)
type TaskType =
  | 'local_bash'         // 后台 Bash 命令
  | 'local_agent'        // 后台 Agent 子会话
  | 'remote_agent'       // 远端 Agent
  | 'in_process_teammate' // 同进程队友
  | 'local_workflow'     // Workflow 脚本
  | 'monitor_mcp'        // MCP 监控
  | 'dream'              // /dream 记忆整理

type TaskStatus = 'pending' | 'running' | 'completed' | 'failed' | 'killed'

isTerminalTaskStatus(s) 判定 completed | failed | killed。这套类型给 daemon ps/logs/attach/kill 系列命令以及 --bg 模式用。

3.10 实战:观察一次完整 turn

DEBUG=cc:query bun run dev
> 在当前目录写一个 hello.ts,然后用 bun 运行它

日志输出顺序:

  1. [query] ownsTrace=true
  2. RequestStartEvent
  3. 一连串 content_block_delta(thinking + 文本)
  4. tool_use{Bash} collected(或者并行 StreamingToolExecutor 已经启动了)
  5. canUseTool() → REPL 弹权限确认 → 用户回车
  6. tool_result appended
  7. 第 2 个迭代开始,又一次 RequestStartEvent
  8. Terminal { reason: 'completed' }

3.11 本章关键文件

文件 行数 关键行号
src/query.ts 2042 query():275 queryLoop():392 State 类型:260 MAX_OUTPUT_TOKENS_RECOVERY_LIMIT:193 yieldMissingToolResultBlocks:148 StreamingToolExecutor gate:744-751
src/query/transitions.ts 21 Terminal 11 reasons / Continue 7 reasons(独立小文件
src/QueryEngine.ts 1365 QueryEngine class:192 constructor:208 submitMessage:217 query() 调用:688-699 ask():1256 fileHistoryMakeSnapshot:653
src/screens/REPL.tsx 6541 REPL():828 useCanUseTool:207 useMergedTools:272 QueryGuard:1165 query 消费循环:3403+
src/Task.ts 80+ 后台 TaskType 7 种 / TaskStatus 5 种

3.12 配套架构图

📊 diagrams/ch03-query-loop.html — queryLoop 完整状态机(11+7 transition 矩阵)


上一章 ← 第 2 章 启动与入口分发 | 下一章 → 第 4 章 Tool System