跳到主要内容

数据截至 (上游 commit 7fb95fe9048f)

Agent 仿真:把 agent 放进剧本里跑一遍

30 秒导读: 前五章讲的是"agent 已经跑过了,我们怎么把它看清楚"。这一章反过来:LangWatch 自己去把 agent 跑一遍——用一个模拟用户按剧本跟它对话,再用一个裁判 agent 判定通过与否。关键不在于"多了个测试功能",而在于两件事:裁判判的时候能读到被测 agent 真实产生的 span;整个仿真过程写回同一条事件管线(第 2 章),所以它和线上流量共用一套存储、投影和 UI。


1. 这是什么(零基础也能懂)

一句话定义: 仿真(simulation)= 给 agent 写一段"剧本",让 LangWatch 派一个模拟用户照剧本跟它聊,聊完派一个裁判按验收条件判分。

它解决什么问题。 普通的 LLM 观测平台是"录像机":agent 在线上跑,平台把 span 收下来给你看。但录像机回答不了"我改了 prompt,退款流程还灵吗"——因为你没有流量,或者不敢拿真实用户试。仿真是"主动出题":平台自己造流量。

剧本长什么样。 一条剧本就是数据库里的一行,字段少得惊人(platform/app/prisma/schema.prisma:2221 model Scenario):

字段干什么
name剧本名,给人看
situation情境描述,喂给模拟用户当"你是谁、你要干嘛"
criteria验收条件数组,喂给裁判当打分标准
labels标签,会被写进子进程的 OTEL resource attributes
simulatorModel / judgeModel可选,单条剧本覆盖项目默认的模拟用户/裁判模型

注意这里没有"期望输出"字段。仿真不比对字符串,它比对的是"裁判读完整段对话和证据后认不认"。

一次运行里有三个角色。 这是理解全章的心智模型:

角色是什么谁提供
被测体(agent)一个 AgentAdapter,包着你的 HTTP 服务 / 平台里的 prompt / 代码 agent / 工作流你配置的 target
模拟用户(user simulator)一个 LLM,拿 situation 演用户,一轮一轮地问ScenarioRunner.userSimulatorAgent
裁判(judge)一个 LLM,拿 criteria 判 success / failure / inconclusiveScenarioRunner.judgeAgent(http target 时由 SDK 叠加远端 span 回查,见 §7)

三者被塞进同一个 agents 数组交给 SDK,由 @langwatch/scenario 这个外部 SDK 驱动轮转 (platform/app/src/server/scenarios/execution/scenario-child-process.ts:153 ScenarioRunner.run):

// 真实源码节选,scenario-child-process.ts:159-163
agents: [
adapter, // 被测体
ScenarioRunner.userSimulatorAgent({ model: simulatorModel }), // 模拟用户
judgeAgent, // 裁判
],

一句话直觉: 把它当成给 agent 做的"情景面试"——HR(模拟用户)按题库提问,面试官(裁判)不光听你怎么答,还要调你的工作记录(span)核对你是不是真做了。

判定结果只有三种(platform/app/src/server/scenarios/scenario-event.enums.ts:4 Verdict):success / failure / inconclusive

本章不重复讲评估器目录与执行(见 04-evaluation.md),也不重复 trace 汇总(见 03-trace-processing.md)。仿真里的"打分"是一个跑在 SDK 内部的裁判 agent,和第 4 章的评估器体系是两套东西。


2. 顶层全景(它大概怎么转)

怎么读这张图:从左到右是一次运行的时间顺序;注意右侧的箭头拐回了左侧的"事件管线"——这就是本章最重要的结构特征,仿真不是旁路系统,它是事件管线的又一个生产者。图里的 subscriber(订阅者)就是第 2 章讲的、挂在投影落库成功之后才触发的副作用钩子(旧称 reactor,ADR-098 改名);而"看到 queued 就把 run 编排进执行池"这类丢一次不可接受的活,交给了带持久状态与精确一次发件箱的 process manager(进程管理器,见 02-event-sourcing.md §4.4)。

[1] 触发 [2] 排队(事件溯源) [3] 执行(隔离)
┌──────────────┐ queueRun ┌─────────────────┐ submit ┌────────────────┐
│ UI / API │────命令────▶│ simulation │─────────▶│ 执行池 │
│ 跑一个套件 │ │ processing 管线 │ process │ 并发 3,超额缓冲 │
└──────────────┘ └────────┬────────┘ manager └───────┬────────┘
│ 折叠投影 │ spawn
▼ ▼
┌─────────────────┐ ┌────────────────┐
│ ClickHouse │ │ 子进程 │
│ simulation_runs │ │ 独立 OTEL │
└────────▲────────┘ │ 跑 SDK 剧本 │
│ └───────┬────────┘
[5] 落库 │ [4] 事件回流 │ HTTP POST
└────────────────────────────┘
/api/scenario-events

部件一句话职责:

部件干什么在哪个文件
SuiteRunService.startRun把"套件 × 目标 × 重复次数"展开成 N 个 run,预生成 ID 并逐个 queueRunplatform/app/src/server/app-layer/suites/suite-run.service.ts:123
simulation-processing 管线命令 → 事件 → 折叠投影 → subscriber/process manager,run 的唯一真相源platform/app/src/server/event-sourcing/pipelines/simulation-processing/pipeline.ts:90
simulationRunExecution process manager看到 queued 事件就经事务性发件箱派发 execute intent,把 job 丢进执行池;另管取消广播兜底与 stall 看门狗platform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.ts:206
ScenarioExecutionPool并发闸门 + 子进程登记表 + 取消标记platform/app/src/server/scenarios/execution/execution-pool.ts:52
scenario-child-process.ts子进程入口:stdin 收料 → 跑 SDK → flush OTEL → stdout 吐结果platform/app/src/server/scenarios/execution/scenario-child-process.ts:98
/api/scenario-events接住 SDK 回报的剧本事件,翻译成管线命令platform/app/src/app/api/scenario-events/[[...route]]/app.ts:33
SimulationRunStateFoldProjection把事件折叠成一行 run 状态(状态、消息、判定、成本)platform/app/src/server/event-sourcing/pipelines/simulation-processing/projections/simulationRunState.foldProjection.ts:261

主线走一遍(高层):

  1. 用户点"跑套件" → SuiteService.run 解析引用、剔除已归档项(platform/app/src/server/suites/suite.service.ts:286)。
  2. SuiteRunService.startRunstartSuiteRun(套件级),再为每个组合 queueSimulationRunCommand(run 级)。
  3. 管线把 queued 事件折叠成 ClickHouse 里状态为 QUEUED 的行——列表立刻就有东西可看,不用等执行。
  4. process manager 的 execute intent 把 run 提交到本 pod 的执行池(outbox 持有重试,丢消息不再可能);池有空位就 spawn 子进程。
  5. 子进程里 SDK 驱动三个 agent 轮转,每一步都 POST 回 /api/scenario-events
  6. 路由把事件翻译成 startRun / messageSnapshot / textMessageEnd / finishRun 命令,再次进管线,折叠更新同一行。

3. 四层 ID:谁包着谁

仿真里有四个 ID,搞混就看不懂后面所有代码。它们的包含关系是严格四层:

scenarioSetId 一个"集合",UI 左侧栏的一项
└─ batchRunId 一次批量运行(点一次"跑"产生一个)
└─ scenarioRunId 一条剧本在这一批里的一次具体执行 ← 事件聚合根
└─ (scenarioId 指向剧本定义本身,跨批次复用)
ID粒度生成于
scenarioId剧本定义(Postgres 行)ScenarioRepository.create,KSUID
scenarioSetId集合/套件命名空间见下方内部命名空间
batchRunId一次批量运行generateBatchRunId(),platform/app/src/server/scenarios/scenario.ids.ts:11
scenarioRunId一次执行,事件溯源的 aggregateIdgenerateScenarioRunId(),scenario.ids.ts:16

集合 ID 的命名空间技巧。 用户可以自己起集合名(SDK/CI 直接传字符串),平台自己也要造集合。为了不撞车,平台造的集合一律加前缀 __internal__(platform/app/src/server/scenarios/internal-set-id.ts:13),再靠后缀区分用途:

形态含义判定函数
任意用户字符串SDK / CI 跑出来的外部集合isInternalSetId 返回 false
__internal__<projectId>__on-platform-scenarios平台上单条"保存并运行"isOnPlatformSet(internal-set-id.ts:31)
__internal__<suiteId>__suite一个测试套件isSuiteSetId(platform/app/src/server/suites/suite-set-id.ts:21)

这个后缀不只是命名——suiteRunSync 订阅者就是isSuiteSetId 判断要不要把这条 run 同步进套件聚合(platform/app/src/server/event-sourcing/pipelines/simulation-processing/subscribers/suiteRunSync.subscriber.ts:61)。命名空间在这里承担了路由职责。

最容易忽略的一招:scenarioRunId 是预生成的。 套件调度时就把每条 run 的 ID 造好,写进 queueRun 事件;子进程再把同一个 ID 通过 RunOptions.runId 交给 SDK(scenario-child-process.ts:187)。于是 SDK 自己回报的事件和平台先写的 queued 事件落在同一个聚合根上,折叠成一行,而不是两行对不上的记录(platform/app/src/server/app-layer/suites/suite-run.service.ts:191-195 的注释直接写明了这个意图)。历史上还留着一个按 ID 去重的合并函数 mergeRunData(platform/app/src/server/scenarios/scenario-run.utils.ts:26)作为兜底。


4. 事件模型:剧本跑出来的东西长什么样

它要解决的小问题: 一次仿真是个持续几十秒到几分钟的过程,UI 要能边跑边看。所以不能等结束才写一条结果,得把过程拆成事件流。

九种事件,但只有五种进库。 类型定义在 platform/app/src/server/scenarios/scenario-event.enums.ts:11 ScenarioEventType:

事件含义持久化?
SCENARIO_RUN_STARTED运行开始,带 name/description/metadata
SCENARIO_MESSAGE_SNAPSHOT当前完整对话快照
SCENARIO_TEXT_MESSAGE_START某条消息开始(占位)
SCENARIO_TEXT_MESSAGE_END某条消息完成,带 traceId
SCENARIO_RUN_FINISHED结束,带 results(verdict / metCriteria / unmetCriteria)
SCENARIO_TEXT_MESSAGE_CONTENT流式增量否,只广播
SCENARIO_TOOL_CALL_START / _ARGS / _END工具调用流式过程否,只广播

分界线就一个函数:isStreamingEvent(platform/app/src/app/api/scenario-events/[[...route]]/app.ts:358)。流式事件只推给前端做打字机效果,不进事件溯源——不为了"实时手感"污染真相源,是很值得抄的取舍。

事件 schema 建在 AG-UI 上。 所有剧本事件继承同一个 base(platform/app/src/server/scenarios/schemas/event-schemas.ts:48 baseScenarioEventSchema),它扩展的是 @ag-ui/core 的事件基类,再加上四个 ID 字段。用户 metadata 用 .passthrough() 放行,平台自己的字段收在严格校验的 langwatch 命名空间里(event-schemas.ts:64 langwatchMetadataSchema)——用户扩展和平台内部字段互不干扰。

从 HTTP 事件到管线命令。 路由做三件事:Zod 校验 → 把内联媒体外置到 stored-objects → 翻译成命令(app.ts:291 dispatchSimulationEvent):

收到的事件派发的命令
RUN_STARTEDsimulations.startRun
MESSAGE_SNAPSHOTsimulations.messageSnapshot(顺带抽出所有 trace_id)
TEXT_MESSAGE_START/ENDsimulations.textMessageStart/textMessageEnd
RUN_FINISHEDsimulations.finishRun

MESSAGE_SNAPSHOT 那一步把每条消息的 trace_id 收集成 traceIds 数组(app.ts:319-321)——这是仿真和 trace 世界的缝合线,第 9 节的成本回填全靠它。

折叠成一行状态。 投影 SimulationRunStateData 就是 ClickHouse simulation_runs 表的形状(platform/app/src/server/event-sourcing/pipelines/simulation-processing/projections/simulationRunState.foldProjection.ts:144),状态迁移一目了然:

queued ──▶ QUEUED ──▶ started ──▶ IN_PROGRESS ──┬─▶ finished ─▶ SUCCESS / FAILURE

└─▶ (无 finished 事件) ─▶ stall 看门狗写 ERROR("stalled")

对应的 handler:handleSimulationRunQueued:293handleSimulationRunStarted:313handleSimulationRunFinished:520

两个细节值得单看:

  • 快照乱序保护。LastSnapshotOccurredAt 旧的快照直接丢弃(simulationRunState.foldProjection.ts:352)——一行 if 挡掉了消息数组被旧状态覆盖的经典事故。
  • 消息体积封顶 64 KiB。 超限就换成一句带前缀的说明文字并打 warn 日志(simulationRunState.foldProjection.ts:49 MAX_MESSAGE_CONTENT_BYTES:58 capOversizedString)。注释写明了肇事场景:语音剧本把 base64 音频塞进 Messages.Content,导致列表接口整个炸掉。它选择让问题可见且有界,而不是静默地把 90 MB 行写进 ClickHouse。

管线的命令与 subscriber 注册见 02-event-sourcing.md;这里只需要知道 run 的聚合类型是 simulation_run,aggregateId 就是 scenarioRunId(platform/app/src/server/event-sourcing/pipelines/simulation-processing/pipeline.ts:95,platform/app/src/server/event-sourcing/pipelines/simulation-processing/commands.ts:30)。


5. 编排:从"点一下运行"到子进程起飞

5.1 为什么不直接在 web 进程里跑

一次仿真会真的调 LLM、真的发 HTTP 请求、真的产生 OTEL span。如果和主进程共用一个 TracerProvider,被测 agent 的 span 会和平台自己的 span 搅在一起,而且每条 run 需要用不同的 API key 和 endpoint 上报。

LangWatch 的答案很直白:一条 run = 一个子进程。隔离靠环境变量完成——父进程只给子进程一份白名单 env(platform/app/src/server/scenarios/execution/child-environment.ts:69 buildChildProcessEnv),塞进 LANGWATCH_API_KEYLANGWATCH_ENDPOINTOTEL_RESOURCE_ATTRIBUTES;子进程 import @langwatch/scenario 时 SDK 在模块加载期读这些变量,自己建一个独立的 TracerProvider(scenario-child-process.ts:15-19 的注释把这条链讲得很清楚)。

5.2 一次执行的完整时序

process manager: 收到 queued 事件
│ (跳过:已请求取消 / 事件里没有 target——后者直接以 ERROR 收尾)
▼ execute intent(事务性发件箱,丢失会重试)
pool.submit(job)
│ 满员 → 压进 _pending 数组

executeScenarioRun
├─ prefetchScenarioData 拉剧本/项目/adapter 数据/三个模型参数
├─ 检查点① pool.wasCancelled? → 写 CANCELLED 终态,收工
├─ spawn 子进程,stdin 写 JSON job data
│ ├─ 15 分钟超时定时器
│ ├─ stdout/stderr 逐行转成父进程结构化日志
│ └─ close 事件
│ ├─ 检查点② wasCancelled? → cancelled
│ ├─ exit code ≠ 0 → failed(带 stderr)
│ └─ exit code = 0 → success
└─ 失败/取消 → ScenarioFailureHandler 补写终态事件

关键锚点:executeScenarioRun(scenario.processor.ts:295)、spawnScenarioChildProcess(:382)、超时(:453,CHILD_PROCESS.TIMEOUT_MS = 15 分钟,platform/app/src/server/scenarios/scenario.constants.ts:38)。

为什么要在 spawn 前先 prefetch。 子进程里没有 Prisma、没有 ClickHouse 连接,拿不到剧本和模型凭据。所以父进程先把一切查好打成一个 JSON 从 stdin 灌进去(prefetchScenarioData,platform/app/src/server/scenarios/execution/data-prefetcher.ts:315)。这个函数还负责解析三个模型角色——被测 prompt 用自己的模型或 scenarios.agent_under_test 角色默认(data-prefetcher.ts:512-518),模拟用户和裁判各自走"套件覆盖 → 剧本覆盖 → 项目默认"的三级 cascade(data-prefetcher.ts:520-529)。

5.3 退出码语义:一个容易写错的地方

测试判定 failure → 子进程 exit 0 (运行本身是成功的)
真的崩了/网络挂了 → 子进程 exit 1

源码里专门写了注释强调(scenario-child-process.ts:7-8)。把"测试没通过"和"执行出错"混成同一个退出码,会让父进程无法区分"agent 表现不好"和"基础设施坏了"——前者要展示判定理由,后者要报警。

子进程在退出前必须 flushOtelTraces()(scenario-child-process.ts:205):SDK 不暴露 observability handle,所以代码直接从全局拿 TracerProvider,还要处理 ProxyTracerProvider 包着真 provider 的情况(:229-231)。flush 失败只 warn 不失败——遥测不能反过来杀掉业务结果

5.4 并发池:237 行的调度器

ScenarioExecutionPool 没有用任何队列库,核心就是四个成员变量(platform/app/src/server/scenarios/execution/execution-pool.ts:52-63):

成员作用
_running: Map<scenarioRunId, ChildProcess>跑着的子进程,取消广播靠它按 ID 找人
_runningJobs: Map<scenarioRunId, ExecutionJobData>已放行但子进程还没注册完的 job——drain 时靠它给"死在 spawn 窗口里"的 run 补终态
_pending: ExecutionJobData[]满员时的缓冲
_cancelled: Set<scenarioRunId>取消标记,三个检查点共用

submit(:137)有空位就起,没空位就压队;deregisterChild(:127)在子进程退出时触发 dequeueNext(:211),而 dequeueNext 一次只放行一个(:233return),下一个等下一次完成再放。并发数来自 SCENARIO_WORKER.CONCURRENCY = 3(platform/app/src/server/scenarios/scenario.constants.ts:28),每个 worker pod 一个池实例(platform/app/src/server/workers/startWorkers.ts:82)——总并发靠 pod 数横向扩。

5.5 两个工程细节

冷启动。 生产环境优先跑预编译 bundle dist/scenario-child-process.js;找不到就回退 tsx,并打一条 error 级日志明说"会导致约 4 分钟冷启动,请跑 build"(platform/app/src/server/scenarios/execution/child-process-spawn.ts:151)。降级但不崩溃,同时把代价喊出来

日志接缝。 父进程的 logger 绑了 scenarioRunId/batchRunId 等上下文,子进程默认拿不到,查日志就串不起来。解法是把上下文 JSON 编码进环境变量 LANGWATCH_LOG_CONTEXT,子进程解码后 logger.child(context)(platform/app/src/server/scenarios/execution/child-logger.ts:33 / :77)。解码失败只往 stderr 写一行警告,绝不抛(:63)。


6. 被测目标怎么接进来(adapter 层)

它要解决的小问题: SDK 只认一种东西——一个有 call(input) => stringAgentAdapter。所以每种"被测体"都要被包成这个形状。

四种 target:

type被测体是什么适配器
prompt平台里存的一条 prompt 配置SerializedPromptConfigAdapter
http你自己部署的远端 agent 服务SerializedHttpAgentAdapter
code平台里的代码 agentSerializedCodeAgentAdapter
workflowoptimization studio 的工作流SerializedWorkflowAgentAdapter

6.1 prompt target:把平台里的 prompt 直接当被测体

最短的一条路(platform/app/src/server/scenarios/execution/serialized-adapters/prompt-config.adapter.ts:59 call):取 prompt 配置 → system prompt 和自带 messages 先过一遍 Liquid 渲染(:77-87)→ 拼消息时会话历史只在模板没有自己引用 messages 时才追加,模板自己排了历史就不再重复塞(:89-95)→ generateText 调模型 → 返回文本。

意义在于:改 prompt 和测 prompt 在同一个平台里闭环,不用把 prompt 复制到测试代码里。

6.2 http target:把远端 agent 包成 adapter

SerializedHttpAgentAdapter.call(platform/app/src/server/scenarios/execution/serialized-adapters/http-agent.adapter.ts:134)一次调用六步:

读 agent 配置 → 渲染 URL 模板 → 拼 headers(自定义 + 认证 + 追踪)
→ 渲染 body 模板 → ssrfSafeFetch 发出去 → JSONPath 抠回复

(旧的同进程版 HttpAgentAdapter 还在 platform/app/src/server/scenarios/adapters/http-agent.adapter.ts:48,但本 commit 下没有生产调用方;子进程路径走的都是 serialized 版。)

三个 Liquid 引擎,不是一个。 这是本节最值得学的设计(platform/app/src/server/scenarios/execution/http-template-engine.ts:89 / :102 / :115):

引擎默认转义为什么
urlLiquidencodeURIComponent插进 URL 的值要 URL 编码
bodyLiquidJSON 字符串转义对话里一个换行/引号就能把请求体 JSON 撑破
headerLiquid纯文本直出header 值既不是 URL 也不是 JSON,套谁的转义都不对

bodyLiquid 还需要一个例外:messages 数组本来就是序列化好的 JSON,再转义一次就变成字符串了。解法是包一层标记类 RawJson(:30),outputEscape 见到它就原样输出。三个引擎都注册了 | raw 过滤器让用户按表达式手动豁免(:92:108:116)。

模板里能用的变量由 buildTemplateContext 决定(:142):messages(RawJson)、threadIdinput(最后一条用户消息;结构化内容也包成 RawJson),以及本轮的 traceId / traceparent(见第 7 节,追踪上下文也能进模板)。

字段映射与别名自动匹配。 用户的 agent 输入参数不一定叫 input,可能叫 queryquestionuser_messageSCENARIO_FIELD_ALIASES(platform/app/src/server/scenarios/execution/resolve-field-mappings.ts:117)列出了三个规范字段的常见别名,computeBestMatchMappings(:160)据此自动配对;实在配不上但只有一个输入,就默认映射到 input(:181-187)。同一文件还保留了改名前的旧字段名映射表 LEGACY_FIELD_NAMES(:12),老配置不会因为改名失效。

多输入工作流的前置校验。 工作流有多个输入却没配映射时,老行为是"第一个输入拿到用户消息,其余全是空字符串",跑完才发现不对。validateWorkflowAgentMappings(platform/app/src/server/scenarios/execution/validate-workflow-mappings.ts:22)在跑之前就抛 BAD_REQUEST,并在错误文案里告诉用户去哪儿配。把一个"跑完才知道错"变成"点之前就报错"

6.3 adapter 怎么跨进程

子进程里没有数据库,所以父进程传过去的不是 adapter 对象(函数无法序列化),而是一包数据 TargetAdapterData。子进程用注册表按 type 重建 adapter(platform/app/src/server/scenarios/execution/serialized-adapter.registry.ts:42 SERIALIZED_ADAPTER_FACTORIES:93 createAdapter)。加一种新目标只需往这张表里加一行,不改 createAdapter 本身。

旧的 serialized.adapters.ts 向后兼容再导出层已移除——四种 adapter 实现直接住在 serialized-adapters/ 目录,由 barrel 文件逐个再导出(platform/app/src/server/scenarios/execution/serialized-adapters/index.ts:8-11)。

模型也是同理:传 LiteLLMParams,子进程用 createModelFromParams 现造一个 OpenAI-compatible provider,base URL 指向 NLP 服务的代理(platform/app/src/server/scenarios/execution/model.factory.ts:156),每个参数以 x-litellm-* header 透传。

认证配置用的是判别联合 + 策略表(platform/app/src/server/scenarios/adapters/auth.strategies.ts:52 AUTH_STRATEGIES),支持 none / bearer / api_key / basic 四种。


7. 裁判怎么读到真实 span(本章的核心)

7.1 问题:只看对话文本的裁判是瞎的

假设剧本的验收条件是"必须真的发起了退款"。如果裁判只读对话,agent 回一句"好的,我已经为您办理退款"就能骗过它。要判"做没做",裁判必须能看到 agent 的 span。

同进程的场景好办:SDK 有个 JudgeSpanCollector,OTEL 仪表化在同一个进程里把 span 直接喂给它。但 http target 的被测 agent 在别人的服务器上,span 走的是用户自己的 OTEL SDK,先到 LangWatch 存起来。裁判怎么拿?

7.2 解法:每轮一个 trace,判决前由 SDK 回查(span 收集已移入 SDK)

怎么读这张图:上半是"去",下半是"回";把两边接上的不再是平台里的桥接代码,而是"每轮 trace 的 id 被打在消息上"这个事实

① adapter 每一轮发请求前
injectTraceContextHeaders → headers 里注入 W3C traceparent
(本轮的 traceId 同时被 SDK 运行时打在每条消息上)


② 用户的 agent 采纳 traceparent,产生的 span 挂在同一轮 trace 上,
经自己的 OTEL 管线上报给 LangWatch(见 01-ingestion.md)


③ 平台只负责配置:对 http target 打开 fetchRemoteTraces,
等待预算按该项目实测的 ingest lag 算出,prefetch 时发给子进程


④ SDK 裁判开口判判决前
收齐所有消息上的 trace id → 从 trace API 拉回 span(settle-wait)
→ 滤掉基础设施 span、与本地已收集的去重
→ 拿不齐就加一条合成 error span 再判

① 注入是按轮的。 SerializedHttpAgentAdapter.call 每轮都调 injectTraceContextHeaders(platform/app/src/server/scenarios/execution/serialized-adapters/http-agent.adapter.ts:138-139):用 OTEL propagation 往 headers 里写 traceparent,顺手把本轮 trace ID 抓下来——这个函数现在住在共享包里(packages/observability/src/trace/traceContext.ts:23),抓的时候会剔掉 OTEL 的全零无效 ID(:47INVALID_TRACE_ID 检查)。本轮的 traceId / traceparent 同时还作为模板变量开放(http-agent.adapter.ts:141-145)。相比之下,旧实现只在 adapter 上存一个 capturedTraceId,每轮覆盖——多轮剧本里裁判只能看到最后一轮的远端 span(ADR-097 自己点名了这个缺陷,dev/docs/adr/097-scenario-remote-trace-judging.md:1)。

③ 平台退成"配置方"。 远端 span 裁判这套机制整体搬进了 @langwatch/scenario SDK(两种语言,config 开关),平台侧只剩一个函数:buildRemoteTraceRunConfig(platform/app/src/server/scenarios/execution/remote-trace-run-config.ts:42)——target 是 http 才返回 fetchRemoteTraces: true,否则返回空对象(:53)。等待预算不再拍脑袋:prefetch 时按这个项目在 stored_spans 上的实测摄入滞后算出 clamp(1.25 × p95 + 5s, 10s, 30s),默认 30 秒、进程内缓存 1 小时(platform/app/src/server/scenarios/execution/ingest-lag.service.ts:66 resolveTraceWaitTimeoutMs:41 默认值),随 job data 发给子进程(data-prefetcher.ts:613)。

④ 裁判变成两阶段,这是延迟契约的关键。 对话中途的裁判调用只做"继续 / 进入判决"的决策,一个 span 都不拉;真正判决那一次先 settle-wait:在共享预算内每秒轮询,直到 trace 里至少有一条远端 span、且拉到的 span 的父级都能对上(父级没齐 = trace 还在路上)。预算耗尽但有部分 span → 保留全部并加一条合成 langwatch.span_collection.error span 标明"trace 不完整";一条都没有或拉取硬失败 → 合成 span 如实报告"什么都没收集到"。裁判还有一次性的 wait_for_traces 工具可以多等一个 traceWaitExtensionMs(平台传的是 30 秒上限,remote-trace-run-config.ts:30 TRACE_WAIT_CAP_MS:61)。以上行为都在 SDK 内;平台的注释与 ADR 是权威描述(remote-trace-run-config.ts:5-9dev/docs/adr/097-scenario-remote-trace-judging.md:1)。

旧实现的平台侧四个文件——remote-span-judge-agent.tsbridge-trace-id.tsremote-span-collector.tstrace-api-span-query.ts(外加 synthetic-error-span.ts)——已随这次搬迁全部移除。它们的两个设计动机仍然成立,只是由 SDK 继承:裁判不该拿"模拟用户和裁判自己产生的 span"当被测方的行为证据(滤噪),以及"没拿到证据"必须明说,不能给裁判一个空集合让它误判(合成 error span)。

这一整套只对 http target 启用。 其它 target 用标准 ScenarioRunner.judgeAgent(scenario-child-process.ts:145-149),证据来源是同进程 collector——它们本来就跑在同一进程里。


8. 可靠性:跑不完的时候怎么办

一次仿真"没有 RUN_FINISHED 事件"有四种成因,处理方式各不相同:

成因谁发现怎么收尾终态
子进程超时(>15 分钟)父进程定时器kill + resolve 失败 → failure handlerERROR
子进程崩溃(exit ≠ 0)close 事件带 stderr 写失败事件ERROR
用户取消取消广播SIGTERM + 写取消事件CANCELLED
worker 整个没了stall 看门狗(process manager 的定时唤醒)到点写终态事件ERROR(reason "stalled")

8.1 失速不再"读时派生",而是由看门狗写成终态

旧实现里 resolveRunStatus(stall-detection.ts,已移除)是个纯函数:有 finishedStatus 就用它;没有就看"距最后一个事件过了多久",超过阈值在读侧派生 STALLED。问题是读时派生的状态不可查询、不可聚合,而且每条读路径都得重算一遍。

现在的答案是 process manager 的唤醒(wake):每条 run 的 simulationRunExecution 进程实例在 queued 和每次活动时都把唤醒截止设为 lastActivity + STALL_THRESHOLD_MS(platform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.ts:301 handleRunActivity);到点仍无活动,wake handler 就发一个 finish intent,写下终态 ERRORerror: "stalled"(:318 simulationRunExecutionWake,具体分支在 :351-357)。阈值仍是子进程超时的 2 倍(platform/app/src/server/scenarios/scenario.constants.ts:51 STALL_THRESHOLD_MS = CHILD_PROCESS.TIMEOUT_MS * 2,即 30 分钟)。唤醒由 wake worker 驱动,不依赖当初干活的那个 pod 还活着——"worker 都没了,谁去写我死了"这个问题,交给了一个带持久状态的编排器

两个配套:ScenarioRunStatus.STALLED 枚举仅为外部 API/历史行保留,新代码已不再产生它(scenario-event.enums.ts:37-42 的注释写明 "nothing derives STALLED at read time anymore");历史遗留的、永远没等到终态的旧 run 由一次性回填任务收口,读侧查询在 platform/app/src/server/event-sourcing/pipelines/simulation-processing/repositories/stalledSimulationRuns.clickhouse.repository.ts:1(阈值 24 小时 + 10 万行保险丝,防止误伤在途 run)。

配套的策略仍是不重试:SCENARIO_QUEUE.MAX_ATTEMPTS = 1,注释直接写 "1 = no retries, immediate fail after stall detection"(platform/app/src/server/scenarios/scenario.constants.ts:22)。一次仿真会真的花钱调 LLM,自动重试可能悄悄翻倍账单。

8.2 取消:一条穿过事件溯源的路径

取消要跨进程(点击在 web pod,子进程在某个 worker pod),LangWatch 没有维护"谁在哪个 pod"的表,而是广播 + 本地认领;与旧实现不同的是,广播的派发由 process manager 经事务性发件箱完成,丢了会重试,还有一个 grace 唤醒兜底:

用户点取消

▼ ScenarioCancellationService.cancelJob
读折叠投影:已是终态? ──是──▶ 什么也不做
│否
▼ dispatch cancel_requested 命令(事件溯源)
├─▶ 折叠投影记下 CancellationRequestedAt(幂等,保留首次时间戳)
└─▶ simulationRunExecution process manager:
├─ 相位 = running → cancel intent → Redis 频道 "scenario:cancel"
│ │ (发件箱重试 + 60 秒 grace 唤醒兜底)
│ ▼ 每个 worker pod 都收到
│ pool.runningChildren.get(id) 命中? ──是──▶ child.kill("SIGTERM")
│ └─ 无论命不命中 → pool.markCancelled(id)
│ (grace 到期仍无终态事件 → 看门狗补写 finish(CANCELLED))
└─ 相位 = queued → 立即补一条 finished(CANCELLED),
且只要 execute intent 发出过就照样广播取消

锚点:ScenarioCancellationService.cancelJob(platform/app/src/server/scenarios/cancellation.ts:92)、终态判定 isCancellableStatus(platform/app/src/server/scenarios/scenario-event.enums.ts:65)、CANCELLATION_CHANNEL(platform/app/src/server/scenarios/cancellation-channel.ts:16)、订阅侧(scenario.processor.ts:598)、取消分支(platform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.ts:321 handleCancelRequested)、cancel intent 的 Redis 发布(process-manager/simulationRunExecutionIntentHandlers.ts:107-111cancellation-channel.ts:48 publishCancellation)。

为什么排队中的 run 也要广播: queued 不等于"没派发"——execute intent 在 run 入队那一刻就发出去了,池里可能已经压着这个 job(等空位或在 prefetch),而 pool.wasCancelled 只有取消订阅侧会设置。不广播的话,这条 run 读作 CANCELLED 却照样跑完、照样花钱(process manager 里 handleCancelRequested 的注释把这条链写得很清楚,simulationRunExecution.process.ts:330-339)。旧实现里"调度侧替排队 job 写终态(occurredAt + 1ms 排序)"的特判已随这次重构收进了 process manager。

三个检查点覆盖竞态。 取消信号可能在任何时刻到达,所以 _cancelled 集合被查了三次:进池/出队时(execution-pool.ts:145:216)、prefetch 结束后(scenario.processor.ts:319)、子进程 close 时(:479)。出队时撞见已取消的 job 还会调 _onSkipCancelled 主动补写 finished(CANCELLED),保证"被跳过"也能到终态。

8.3 补终态:不让 run 永远卡在 IN_PROGRESS

ScenarioFailureHandler.ensureFailureEventsEmitted(platform/app/src/server/scenarios/scenario-failure-handler.ts:123)负责在异常路径上派发 finishRun,状态取 CANCELLEDERROR,判定内容由 buildFailureResults(platform/app/src/server/scenarios/scenario-failure-results.ts:29)构造——取消记 INCONCLUSIVE,失败记 FAILURE。依赖的是命令的幂等性(FinishRunCommand 的 idempotencyKey 是 ${tenantId}:${scenarioRunId}:finishRun,platform/app/src/server/event-sourcing/pipelines/simulation-processing/commands/finishRun.command.ts:138),所以重复补也不会写出两条终态。


9. 落库与聚合:一次运行的余波

finishRun 落地时,同一条事件会扇出到多个 subscriber / process manager(注册见 platform/app/src/server/event-sourcing/pipelines/simulation-processing/pipeline.ts:110-129):

部件触发事件干什么
snapshotUpdateBroadcast 订阅者消息类事件推 SSE 给前端刷新
simulationRunExecution process managerqueued / cancel_requested / 活动事件 / 唤醒派发执行池、取消广播(带 grace 兜底)、stall 看门狗
suiteRunSync 订阅者started / finished同步到套件聚合(仅套件集合)
traceMetricsSync 订阅者finished把 trace 的成本/延迟拉回这条 run

(另有一个可选的 customerIoSimulationSync 订阅者做 Customer.io 同步。旧实现里独立的 cancellationBroadcast / scenarioExecution reactor 已不复存在:取消广播与执行派发都并进了 simulationRunExecution process manager——细节见 §8。)

两级聚合。 run 级聚合根是 scenarioRunId;套件级另起一条管线,聚合根是 batchRunId(platform/app/src/server/event-sourcing/pipelines/suite-run-processing/pipeline.ts:36)。跨管线的桥是 suiteRunSync 订阅者——它住在仿真管线上,消费仿真事件,派发套件命令(platform/app/src/server/event-sourcing/pipelines/simulation-processing/subscribers/suiteRunSync.subscriber.ts:51),并在文件注释里写明了这个方向性;套件管线自己一个订阅者都没有(源码注释原话:"No subscriber on this pipeline — cross-pipeline subscribers live on the simulation pipeline",suite-run-processing/pipeline.ts:34)。

成本回填走"拉"不走"推"。 span 可能比仿真事件先到,也可能后到。traceMetricsSync 在 run 结束时,对每个还没有指标的 traceId 派发 computeRunMetrics(platform/app/src/server/event-sourcing/pipelines/simulation-processing/subscribers/traceMetricsSync.subscriber.ts:40-53),由命令自己去读 trace 汇总;trace 还没到就安排延迟重试。这里的异常必须重抛,注释说得很直白:这是最后一次机会,吞掉就永久丢指标(:24-26)。折叠端把每个 trace 的花费按角色拆开累加(simulationRunState.foldProjection.ts:592 handleSimulationRunMetricsComputed)——所以 UI 能告诉你"这次仿真里,模拟用户花了多少、裁判花了多少、被测 agent 花了多少"。

trace 汇总本身怎么算,见 03-trace-processing.md

读侧。 应用层入口是 SimulationRunService(platform/app/src/server/app-layer/simulations/simulation-run.service.ts:17),对着 ClickHouse 仓库提供集合列表、批次历史、单次 run 明细等查询。旧版那条独立的 Elasticsearch 分析路径(scenario-analytics.tscreateScenarioAnalyticsQuery)已移除,ClickHouse 投影成为唯一读路径。


10. 巧妙之处(可以直接抄走的)

  1. 让裁判看证据而不只看话术。 整条 trace 接力(每轮注入 traceparent → trace id 打在消息上 → 判决前 settle-wait 回查)存在的唯一理由,就是让"有没有真的做"变得可判。平台侧锚点:packages/observability/src/trace/traceContext.ts:23remote-trace-run-config.ts:42;机制本身已移入 SDK(ADR-097)。

  2. 拿不到证据就造一条"我没拿到证据"的 span。 而不是给裁判一个空数组让它误判。旧的平台侧 synthetic-error-span.ts 已移除,行为由 SDK 继承,平台注释里留有权威描述(remote-trace-run-config.ts:5-9)。诚实优先的思路,在任何 LLM-as-judge 系统里都成立。

  3. 过滤掉自家基础设施 span。 裁判不能拿自己和模拟用户的 span 当被测方的行为证据——滤噪逻辑同样随裁判移进了 SDK(dev/docs/adr/097-scenario-remote-trace-judging.md:1)。

  4. 实时事件不进真相源。 流式 delta 只广播、不持久化(app.ts:358 isStreamingEvent),UI 手感和数据一致性各得其所。

  5. ID 预生成打通两侧写入。 平台先造 scenarioRunIdqueued,再把同一个 ID 交给 SDK 当 runId,两边落进同一个聚合(suite-run.service.ts:191-195scenario-child-process.ts:187)。

  6. 失速由看门狗写成终态。 没有活着的进程能写"我死了"这条事件——那就让带持久状态的 process manager 定时唤醒,自己补写(simulationRunExecution.process.ts:395)。读时派生的旧实现(stall-detection.ts)已移除。

  7. body / URL / header 用三个转义策略不同的模板引擎。 加上 RawJson 标记类做定点豁免(http-template-engine.ts:30/89/102/115)——这是把"用户输入里的换行撑破请求体 JSON"这一类 bug 从根上消掉。

  8. 降级要吵。 预编译 bundle 缺失时回退 tsx,但打 error 日志并写明"约 4 分钟冷启动 + 修复命令"(child-process-spawn.ts:151)。

  9. 前置校验替代事后困惑。 多输入工作流没配映射直接拒绝启动(validate-workflow-mappings.ts:22)。

  10. 进程边界上手动搬运日志上下文。 一个 env 变量换来跨进程可 join 的日志(child-logger.ts:33)。


11. 边界与局限(诚实清单)

  • 旧版未接线代码已清场。 上一版里"有代码但没生产调用方"的 SimulationRunnerService(simulation-runner.service.ts)与 ScenarioExecutionOrchestrator(execution/orchestrator.ts)已移除;活的路径是 queueRun 命令 → simulationRunExecution process manager → 执行池 → 子进程(见 §2 部件表与 §13 代码地图)。一个反方向的新例子:同进程版 HttpAgentAdapter(platform/app/src/server/scenarios/adapters/http-agent.adapter.ts:48)在本 commit 下没有生产调用方,子进程路径全部走 SerializedHttpAgentAdapter

  • 裁判读 span 只覆盖 http target。 另外三种 target 用标准 judge(scenario-child-process.ts:145-149),证据来源是同进程 collector;buildRemoteTraceRunConfig 对非 http 直接返回空配置(remote-trace-run-config.ts:53)。ADR-097 注明:给 workflow/code target 开远端拉取是后续工作。

  • 远端 span 等待预算按项目实测,上限 30 秒。 预算来自该项目 stored_spans 的摄入滞后 p95(ingest-lag.service.ts:66,默认 30 秒见 :41);判决时裁判还有一次性的 wait_for_traces 延长(remote-trace-run-config.ts:30)。用户的 OTEL 管线再慢,超出部分就由那条合成 error span 标明"trace 不完整"——这是可靠性与等待时间的显式取舍。

  • 不重试。 MAX_ATTEMPTS = 1(platform/app/src/server/scenarios/scenario.constants.ts:22),偶发网络抖动会直接算失败。

  • 消息内容硬上限 64 KiB。 超了就被替换成说明文字(simulationRunState.foldProjection.ts:49),语音/多模态剧本尤其容易踩到。

  • ArchiveSetCommand 的扇出没做完。 源码注释写着:把集合归档事件接进逐 run 的折叠投影还是待办(跟踪号 lw#3636),需要"一个 set 事件 → 多个 run 聚合"的 fanout(platform/app/src/server/event-sourcing/pipelines/simulation-processing/commands.ts:144-146)。

  • 判定语义有损。 inconclusive 在折叠时和 failure 一起被归成 FAILURE(simulationRunState.foldProjection.ts:560-561),原始 verdict 另存字段,但状态列上二者不可区分。


12. 横向对比(同一 shelf 的取舍)

关切本章的做法对照
谁产生流量平台自己造(模拟用户)01-ingestion.md:等外部把 span 送进来
存哪复用第 2 章的命令/事件/投影,不另起表02-event-sourcing.md
谁打分剧本内的裁判 agent,读对话 + span04-evaluation.md:评估器目录,按 trace 打分
隔离手段每 run 一个子进程 + 独立 TracerProvider05-ai-gateway.md:Go 进程里的拦截器链

一句话记忆:第 1/3 章是"看",第 4 章是"评",第 5 章是"管",本章是"考"。


13. 代码地图(导航索引)

主题文件路径符号名
剧本定义platform/app/prisma/schema.prismamodel Scenario
剧本 CRUDplatform/app/src/server/scenarios/scenario.service.tsScenarioServicegetById
ID 生成platform/app/src/server/scenarios/scenario.ids.tsgenerateBatchRunIdgenerateScenarioRunId
内部集合命名空间platform/app/src/server/scenarios/internal-set-id.tsINTERNAL_SET_PREFIXisOnPlatformSetexpandSetIdFilter
套件集合命名空间platform/app/src/server/suites/suite-set-id.tsisSuiteSetIdgetSuiteSetId
事件类型/状态枚举platform/app/src/server/scenarios/scenario-event.enums.tsScenarioEventTypeScenarioRunStatusVerdictisCancellableStatus
事件 schemaplatform/app/src/server/scenarios/schemas/event-schemas.tsbaseScenarioEventSchemalangwatchMetadataSchemascenarioEventSchema
事件入口路由platform/app/src/app/api/scenario-events/[[...route]]/app.tsdispatchSimulationEventisStreamingEvent
套件调度platform/app/src/server/app-layer/suites/suite-run.service.tsSuiteRunService.startRun
管线定义platform/app/src/server/event-sourcing/pipelines/simulation-processing/pipeline.tscreateSimulationProcessingPipeline
命令定义platform/app/src/server/event-sourcing/pipelines/simulation-processing/commands.tsQueueRunCommandFinishRunCommandCancelRunCommandArchiveSetCommand
折叠投影platform/app/src/server/event-sourcing/pipelines/simulation-processing/projections/simulationRunState.foldProjection.tsSimulationRunStateFoldProjectioncapOversizedString
执行编排 process managerplatform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.tshandleRunQueuedhandleCancelRequestedsimulationRunExecutionWake
执行编排 intent 落地platform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecutionIntentHandlers.tscreateExecuteRunHandlercreateCancelExecutionHandlercreateFinishRunHandler
取消/停滞历史回填读侧platform/app/src/server/event-sourcing/pipelines/simulation-processing/repositories/stalledSimulationRuns.clickhouse.repository.tsStalledHistoricalRun
套件同步 subscriberplatform/app/src/server/event-sourcing/pipelines/simulation-processing/subscribers/suiteRunSync.subscriber.tscreateSuiteRunSyncSubscriber
成本回填 subscriberplatform/app/src/server/event-sourcing/pipelines/simulation-processing/subscribers/traceMetricsSync.subscriber.tscreateTraceMetricsSyncSubscriber
快照广播 subscriberplatform/app/src/server/event-sourcing/pipelines/simulation-processing/subscribers/snapshotUpdateBroadcast.subscriber.tscreateSnapshotUpdateBroadcastSubscriber
套件聚合管线platform/app/src/server/event-sourcing/pipelines/suite-run-processing/pipeline.tscreateSuiteRunProcessingPipeline
并发池platform/app/src/server/scenarios/execution/execution-pool.tsScenarioExecutionPoolsubmitdequeueNext
子进程编排platform/app/src/server/scenarios/scenario.processor.tsexecuteScenarioRunspawnScenarioChildProcessstartScenarioProcessor
子进程环境白名单platform/app/src/server/scenarios/execution/child-environment.tsbuildChildProcessEnv
子进程入口platform/app/src/server/scenarios/execution/scenario-child-process.tsexecuteScenarioflushOtelTraces
spawn 命令解析platform/app/src/server/scenarios/execution/child-process-spawn.tsresolveChildProcessSpawn
日志上下文桥platform/app/src/server/scenarios/execution/child-logger.tsencodeScenarioLogContextcreateChildProcessLogger
数据预取platform/app/src/server/scenarios/execution/data-prefetcher.tsprefetchScenarioData
摄入滞后测量(span 等待预算)platform/app/src/server/scenarios/execution/ingest-lag.service.tsresolveTraceWaitTimeoutMs
远端裁判配置(SDK 开关)platform/app/src/server/scenarios/execution/remote-trace-run-config.tsbuildRemoteTraceRunConfigTRACE_WAIT_CAP_MS
HTTP 适配器platform/app/src/server/scenarios/execution/serialized-adapters/http-agent.adapter.tsSerializedHttpAgentAdapter
Prompt 适配器platform/app/src/server/scenarios/execution/serialized-adapters/prompt-config.adapter.tsSerializedPromptConfigAdapter
认证策略platform/app/src/server/scenarios/adapters/auth.strategies.tsAUTH_STRATEGIESapplyAuthentication
模板引擎platform/app/src/server/scenarios/execution/http-template-engine.tsRawJsonbuildTemplateContextrenderBodyTemplate
字段映射platform/app/src/server/scenarios/execution/resolve-field-mappings.tsresolveFieldMappingscomputeBestMatchMappings
工作流映射校验platform/app/src/server/scenarios/execution/validate-workflow-mappings.tsvalidateWorkflowAgentMappings
跨进程 adapter 注册表platform/app/src/server/scenarios/execution/serialized-adapter.registry.tsSERIALIZED_ADAPTER_FACTORIEScreateAdapter
模型工厂platform/app/src/server/scenarios/execution/model.factory.tscreateModelFromParamscreateJudgeModelFromParams
追踪头注入(共享包)packages/observability/src/trace/traceContext.tsinjectTraceContextHeadersgetActiveTraceId
远端裁判 ADRdev/docs/adr/097-scenario-remote-trace-judging.mdADR-097
停滞阈值常量platform/app/src/server/scenarios/scenario.constants.tsSTALL_THRESHOLD_MSSCENARIO_WORKERCHILD_PROCESS
取消服务platform/app/src/server/scenarios/cancellation.tsScenarioCancellationService.cancelJob
取消通道platform/app/src/server/scenarios/cancellation-channel.tsCANCELLATION_CHANNELpublishCancellationsubscribeToCancellations
失败补偿platform/app/src/server/scenarios/scenario-failure-handler.tsScenarioFailureHandler.ensureFailureEventsEmitted
失败判定构造platform/app/src/server/scenarios/scenario-failure-results.tsbuildFailureResults
读侧服务platform/app/src/server/app-layer/simulations/simulation-run.service.tsSimulationRunService
池的装配点platform/app/src/server/workers/startWorkers.tsnew ScenarioExecutionPool

已移除(上一版在本表有行):remote-span-judge-agent.tsbridge-trace-id.tsremote-span-collector.tstrace-api-span-query.tssynthetic-error-span.ts(远端裁判机制移入 SDK,见 §7)、stall-detection.ts(读时派生由 stall 看门狗取代,见 §8.1)、scenario-analytics.ts(ES 分析路径)、simulation-runner.service.ts / execution/orchestrator.ts(未接线代码清场)、background/worker.ts(装配点移到 workers/startWorkers.ts)。