跳到主要内容

数据截至 (上游 commit 7fb95fe9048f)

LangWatch — 架构与原理

30 秒导读: LangWatch 是一个开源的 LLM / AI agent 平台,把「观测、评估、仿真、网关」四件事做在一起。它的底层不是「收到数据写进表」,而是一条事件溯源(event sourcing,把状态变化记成一串不可变事件,状态由事件推导而来)管线:任何输入先变成一条不可变事件写进 ClickHouse 的 event_log 表,之后的 trace 汇总、评估触发、告警、实时推送,全都是从这条日志派生出来的。

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

一句话定义: LangWatch 是一个 LLM / agent 的全生命周期平台——记录 agent 每一步在干什么(观测)、给这些步骤自动打分(评估)、在上线前拿剧本把 agent 跑一遍(仿真),并且可以让所有模型调用都从它的代理走一遍(网关)。

解决什么问题、给谁用: 假设你写了一个客服 agent 放到线上。你想知道三件事:昨天那次它为什么答错了?它答得好不好,能不能自动打分?改了 prompt 之后,有没有把以前对的地方搞坏?LangWatch 就是回答这三个问题的地方——它是给把 LLM 应用做到生产环境的团队用的。

它能做什么,分五块:

能力块白话主要落地位置(相对克隆根)
观测(tracing)收 OpenTelemetry 的 span,拼成可查询的 traceplatform/app/src/server/routes/otel.ts
在线评估 / 护栏trace 一进来就按规则跑评分器,也能当拦路护栏用platform/app/src/server/event-sourcing/pipelines/evaluation-processing/
Agent 仿真用模拟用户 + 裁判,按剧本把 agent 跑一遍platform/app/src/server/scenarios/
工作流 / 实验引擎可视化搭出评估工作流,在 Go 引擎里执行;离线实验批量跑数据集services/nlpgo/app/engine/
AI 网关OpenAI/Anthropic 兼容代理,做虚拟密钥、预算、限流、护栏、跨厂商回退services/aigateway/

用起来什么样: 最小接入就是装 SDK、setup() 一下,再给入口函数加个装饰器(本例取自 sdks/python/README.md:50):

import langwatch

langwatch.setup() # 读 LANGWATCH_API_KEY,装上 OpenTelemetry 导出器

@langwatch.trace() # 这个函数的整段执行 = 一个 trace(一棵 span 树)
async def handle_message(text: str):
... # 里面的 OpenAI 调用会被自动记成子 span

之后 SDK 把 span 以 OTLP 协议 POST 到 /api/otel/v1/traces(platform/app/src/server/routes/otel.ts:74otelIngestAuth),剩下的事都发生在服务端。

一句话直觉: 把 LangWatch 想成会计里的总账 + 报表关系。event_log 是总账——只追加、不涂改、每笔都留底;你在界面上看到的 trace 列表、成本统计、评估分数,全是报表,是从总账重算出来的。报表算错了可以推倒重算(重放),总账不动。

本节不涉及底层代码。记住一句话就够:每个 span 先落成一条不可变事件,其他一切都是这条事件的下游。

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

2.1 多个进程,一个事件管线

怎么读这张图: 左边是数据的两个来源,中间是唯一的服务端(TypeScript),右边是它派生出来的存储;下方是被服务端"喊起来干活"的执行侧进程。

你的 agent ┌──────────────────────────────┐
(Python/TS SDK, OTLP) ───▶│ │
│ LangWatch 服务端 (Node/TS) │───▶ ClickHouse
任何 HTTP 客户端 │ 入口 → 命令 → 事件 → 投影 │ event_log
(走网关代理模型) ─────────▶│ → 订阅者 │ + 派生表
│ └───────────┬──────────────────┘
│ │ 订阅者派活
▼ ▼
┌───────────────┐ ┌────────────────────────────┐
│ AI 网关 (Go) │ │ 执行侧 │
│ aigateway │─── span ▶│ nlpgo(Go 工作流/评估引擎) │
└───────────────┘ 回灌 │ langevals(Py 评分器库) │
└────────────────────────────┘

Go 侧服务打包成同一个二进制,靠第一个命令行参数分派(cmd/service/main.go:46services map,现在有 aigateway / langyagent / nlpgo 三个入口)。

2.2 一条 span 的命运(主线)

怎么读这张图: 从上到下是一条 span 的处理顺序;第 ④ 步是分水岭——它之上是"写总账",它之下全是"算报表"。订阅者(subscriber) 是本库反复出现的术语,指「投影落库成功之后才被触发的副作用挂钩」——它前世的名字是反应器(reactor),2026-08 的 ADR-098 把这个词汇正式退役,机器没换、名字换了(详见第 2 章)。

① HTTP 入口 POST /api/otel/v1/traces 鉴权 → 套餐配额 → 解析 OTLP


② 采集服务 去重锁 (tenant, trace, span) 太老/噪声 span 直接丢弃


③ 命令 recordSpan PII 脱敏 / 成本富化 / token 估算


④ 事件日志 event_log (ClickHouse) 唯一真相,只追加

┌───┴──────────────────────┐
▼ ▼
⑤ fold 投影 ⑥ map 投影
traceSummary spanStorage(另:日志/指标拆到独立管线)
同一 trace 严格串行 逐条独立、可并行
│ │
│ 状态落库成功 │ append 成功
▼ ▼
⑦ 挂 fold 的订阅者 ⑦' 挂 map 的订阅者
评估触发 · 告警 · span 落库推送(spanStorageBroadcast)
WebSocket 推送 ·
仿真/实验指标回填

订阅者可以挂在 fold 上,也可以挂在 map 上。 声明形如 .withSubscriber("名字", { fold: "traceSummary", when, handler })——fold/map 字段指明挂在哪个投影上(platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts:326),路由器为两者各留了一条派发路径(fold 投影提交成功后派发在 projectionRouter.ts:1799,map 在 :719)。trace 管线里 10 个订阅者挂 traceSummary fold,1 个(spanStorageBroadcast)挂 spanStorage map(platform/app/src/server/event-sourcing/pipelines/trace-processing/pipeline.ts:197-313)。

2.3 部件一句话职责

部件干什么在哪个文件符号
OTLP 路由收 OTLP 的 traces/logs/metrics,鉴权 + 限额platform/app/src/server/routes/otel.ts:74otelIngestAuthhandlerManagedAuth
采集服务遍历 OTLP 请求,逐 span 过去重锁再发命令platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:115handleOtlpTraceRequest
命令处理器把一次"意图"变成一条或多条事件platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:170RecordSpanCommand.handle
事件存储把事件写进 ClickHouse event_logplatform/app/src/server/event-sourcing/stores/eventStoreClickHouse.ts:22EventStoreClickHouse
投影路由器把事件分发给 fold / map 投影,再触发订阅者platform/app/src/server/event-sourcing/projections/projectionRouter.ts:886ProjectionRouter.dispatch
fold 执行器加载状态 → apply → 落库;乱序时整段重算platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.ts:355FoldProjectionExecutor.execute
map 执行器一条事件 → 一条记录 → append,无状态platform/app/src/server/event-sourcing/projections/mapProjectionExecutor.ts:32MapProjectionExecutor.execute
组队列同组 FIFO、跨组并行的 Redis 队列platform/app/src/server/event-sourcing/queues/groupQueue/groupQueue.ts:303GroupQueueProcessor
管线注册表一次性把 19 条管线接到运行时platform/app/src/server/event-sourcing/pipelineRegistry.ts:456PipelineRegistry.registerAll
AI 网关Go 拦截器链:限流→策略→选模型→缓存→预算→护栏→追踪services/aigateway/app/app.go:85App.buildInterceptors
工作流引擎解析 DAG、分层调度、每个节点交给块执行器services/nlpgo/app/engine/engine.go:37Engine
仿真执行 process manager看到 queued 事件就把这条 run 编排进执行池platform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.ts:206handleRunQueued
仿真执行池并发闸门 + 子进程登记表 + 取消标记,每条 run 起一个子进程platform/app/src/server/scenarios/execution/execution-pool.ts:52ScenarioExecutionPool

2.4 十九条管线,同一套内核

服务端不是一条管线,而是 19 条结构完全相同的管线,用同一个 definePipeline() 建造器(platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts:637)声明,在 PipelineRegistry.registerAll() 里统一注册——registerAll() 里有 19 次 this.deps.eventSourcing.register(

与本文各章相关的主干是这六条:

管线名聚合类型(aggregateId 是什么)干什么定义位置
trace_processingtrace(traceId)span → trace 汇总 + 存储platform/app/src/server/event-sourcing/pipelines/trace-processing/pipeline.ts:169
evaluation_processingevaluation(evaluationId)一次评估的生命周期与结果platform/app/src/server/event-sourcing/pipelines/evaluation-processing/pipeline.ts
simulation_processingsimulation_run(scenarioRunId)一次仿真跑的消息流与判定platform/app/src/server/event-sourcing/pipelines/simulation-processing/pipeline.ts
suite_run_processingsuite_run一批仿真(套件)的聚合状态platform/app/src/server/event-sourcing/pipelines/suite-run-processing/pipeline.ts
experiment_run_processingexperiment_run离线实验跑的逐条结果与汇总platform/app/src/server/event-sourcing/pipelines/experiment-run-processing/pipeline.ts
log_processing / metric_processinglog / metricOTLP 日志与指标各自的存储(从 trace 管线拆出)platform/app/src/server/event-sourcing/pipelines/log-processing/pipeline.ts:20

其余 13 条(authz-grantsautomationsbilling-reportingblob-maintenancecoding-agent-processinggateway-spend-processinggithub-maintenancegovernance-eventslangy-conversation-processinglangy-maintenanceprocess-manager-maintenancetopic-clustering-processing 等)是同一内核在别的域上的复用,机制完全同构,本文不逐条展开。

事件类型名带一套四段命名法 <来源>.<域>.<聚合类型>.<标识>,例如 lw.obs.trace.span_received(platform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/constants.ts:1,命名法说明见 platform/app/src/server/event-sourcing/domain/aggregateType.ts)。

2.5 主线走一遍(高层,不进代码)

  1. SDK 把 span 用 OTLP 发到 /api/otel/v1/traces;路由先鉴权、再查套餐配额,然后解析请求体。
  2. 采集服务逐 span 处理:起始时间早于 31 天的丢弃、编码 agent 的纯基建噪声 span 过滤掉,剩下的先抢一把 (租户, traceId, spanId) 的去重锁(platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:222ingestNormalizedSpan)。
  3. 过关的 span 变成一条 recordSpan 命令,进 Redis 组队列。
  4. 命令处理器做 PII 脱敏、成本富化、token 估算,然后产出一条 lw.obs.trace.span_received 事件(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:364)。
  5. 事件先写进 event_log(必须成功),再分发给投影(platform/app/src/server/event-sourcing/services/eventSourcingService.ts:308storeEvents)。
  6. fold 投影把这条事件叠进 traceSummary 状态——同一个 traceId 的事件严格串行处理;map 投影把这条事件独立翻译成 stored_spans 的一行。
  7. 投影落库成功后,挂在它上面的订阅者才被触发:挂 fold 的负责派评估任务、判告警、往前端推 WebSocket 更新、回填仿真/实验指标(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1799);挂 map 的负责推 span 落库通知(projectionRouter.ts:727)。

3. 阅读地图

六章按"数据流的先后 + 由浅入深"排列。只想懂原理,读 01→02→03 三章就够;要落到具体产品能力,再按需挑 04/05/06。

顺序章节讲什么什么时候该读
1数据入口:一条 span 怎么被接住OTLP / REST 两个入口、鉴权与配额、去重锁、超大 span 的 S3 甩包、噪声过滤想知道"我发的 span 为什么没出现"
2事件溯源内核:命令、事件、投影、订阅者四个基本概念、definePipeline 建造器、组队列的有序性、乱序重算、重放、subscriber 与 process manager 的分工全库最核心的一章,理解其他所有章的前提
3trace 处理管线:从 span 事件到可查询的 tracetraceSummary fold 怎么攒出成本/时延/输入输出,map 投影,十来个订阅者,读侧怎么拼回一条 trace想知道界面上那些数字是怎么算的
4评估层:一个分数是怎么算出来的评估器的目录与字段映射、前置条件闸门、三个执行后端(TS 进程内 / Python langevals / Go DAG 引擎)、护栏语义;外加 Go 工作流引擎的分层调度与离线实验入口要做在线评估或护栏,或想看可视化工作流/批量实验怎么跑
5AI 网关:一个 Go 拦截器链上的治理层洋葱式拦截器、虚拟密钥、预算、缓存、跨厂商回退、trace 回灌要做多厂商代理与成本治理
6Agent 仿真:把 agent 放进剧本里跑一遍剧本、模拟用户、裁判 agent、子进程执行池、套件聚合要在上线前做回归测试

4. 巧妙之处(可借鉴的技术)

① 「投影落库成功才触发订阅者」这条规矩,替掉了一整套分布式事务。 副作用(发评估、告警、推送)不挂在事件到达时,而挂在投影成功写入之后:fold 路径先 foldExecutor.execute 拿到状态(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1739),成功了才 dispatchToSubscribers(projectionRouter.ts:1799);map 路径在 append 成功的同一个事务闭包里派发(projectionRouter.ts:727)。这样任何订阅者看到的都是已经持久化的状态,不会出现"通知说 trace 完成了,但库里还没有"的空窗。

② fold 状态自己就是进度检查点,不额外存 offset。 注释写得很直白:the fold state in the store serves as the checkpoint(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1688)。状态里带一个 LastEventOccurredAt 字段,既是业务数据也是进度标记,少一套需要和数据同步的元数据。

③ 乱序只在"严格更早"时才重算,相等不重算。 事件的 occurredAt 严格小于状态里记录的最大值时,才从事件日志把这个聚合的历史全拉回来重跑一遍;相等时按到达顺序处理(platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.ts:450-462)。SDK 常常用同一毫秒发出 snapshot 和 finished 两条事件,这条豁免省掉了大量无谓的全量重算。

④ 用单调递增的 UpdatedAt 给 ClickHouse 的去重当裁判。 派生表都是 ReplacingMergeTree(UpdatedAt)(platform/app/src/server/clickhouse/migrations/00002_create_schema.sql:178),同版本列相同就靠这一列定胜负。所以 apply() 强制 Math.max(Date.now(), 上一次 + 1)(platform/app/src/server/event-sourcing/projections/abstractFoldProjection.ts:233)——同一毫秒内连写也保证严格递增,不会出现"新状态被旧状态盖掉"。

⑤ 队列积压时合批把 O(n²) 压成 O(n);聚合级订阅者随批折叠,事件级订阅者仍逐条派发。 同一聚合堆积的事件最多 500 条(2026-08 从 100 提到 500)合成一次"读状态 → 折叠 → 写状态"(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:71foldProjectionExecutor.ts:519executeBatch);派发订阅者时按 jobId 折叠——makeJobId 带事件 id 的(如 customEvaluationSync 要逐条读 span)仍每条事件各派一次,聚合级 key 的(广播、告警)合并成一个任务,注释写明这是"用一次折叠替代 N 次 serialize+gzip+blob 往返,殊途同归"(projectionRouter.ts:1957-1969collapseByJobId)。

⑥ 分层 group key 一把解决"谁该串行、谁该并行"。 队列的分组键是 ${租户}/${拓扑路径}/${聚合},例如 proj_x/fold/traceSummary/reactor/evaluationTrigger/trace:abc(platform/app/src/server/event-sourcing/services/queues/queueManager.ts:182buildGroupKey;路径模板在 queueManager.ts:799——内部历史原因仍叫 reactor/ 段,ADR-098 改名后未动队列拓扑)。同一个 trace 的同一个投影严格 FIFO,不同 trace、不同投影天然并行,租户之间天然隔离——三件事一个字符串搞定,不需要分布式锁。

⑦ 「写日志用全量,喂投影用瘦身」,并且实时与重放共用同一个瘦身函数。 event_log 存完整事件,分发给投影前先过一次 leanForProjection 裁掉超大属性(platform/app/src/server/app-layer/traces/lean-for-projection.ts:186)。同一个函数也在重放路径里调用,文件头注释明确列出两个调用点(live 在 eventSourcingService.ts:365,replay 在 replayExecutor.apply),并说这是为了让投影状态与路径无关——线上算出来的和重放算出来的必须一模一样(platform/app/src/server/app-layer/traces/lean-for-projection.ts:5-13)。

⑧ 宁可抛错,也不写一条"属性被清空"的事件。 超大 span 在入口甩到 S3、命令里只留一个 spoolRef;如果处理这条命令的进程没配 blob 存储,代码直接抛错而不是继续(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:199-204)。理由写在注释里:event_log 是唯一真相,写进去一条空属性事件就是不可逆的静默数据丢失,让命令重试或让配置错误暴露出来都比这个好。

⑨ 遇到失控的超大 trace,丢工作不丢数据。 一条 trace 超过 512 个 span(上限 2026-08 仍是 MAX_PROCESSED_SPANS = 512,platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/traceSummary.foldProjection.ts:96)后,评估触发器停止派活,但 span 照存、trace 照样能查(超限判定在 platform/app/src/server/event-sourcing/pipelines/trace-processing/subscribers/evaluationTrigger.subscriber.ts:116-126);而且"跨过阈值"的告警日志只在恰好等于阈值那一次打,避免自身成为放大器。

⑩ 评估自己产生的 span 不能再触发评估。 评估工作流发出的 span 带 langwatch.reserved.causality_depth 属性,触发器见到深度 ≥ 1 就跳过(因果环守卫见 platform/app/src/server/event-sourcing/pipelines/trace-processing/subscribers/evaluationTrigger.subscriber.ts:140-191)。护栏本身还带一个可从运维界面翻转的 SYSTEM 开关 ops_es_causality_loop_guard_disabled,出事能立刻回滚而不必重新部署。

⑪ 网关的拦截器链是"能力可缺省"的洋葱。 每个拦截器就是 {Name, Sync, Stream} 三个字段,Build() 从后往前包一遍(services/aigateway/app/pipeline/pipeline.goInterceptorBuild);buildInterceptors() 里某个依赖是 nil 就不挂那一层(services/aigateway/app/app.go:85)。内部服务想跳过全部治理时,还能直接引用裸的 dispatcher 包当库用,不走 HTTP(services/aigateway/dispatcher/dispatcher.goDispatcher.Dispatch)。

⑫ 预测式判断(when 守卫)失败时向"多做一次"倾斜。 订阅者的过滤谓词抛异常时,代码记日志并当作 true 继续派发(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1996-2002subscriberShouldDispatch;契约见 platform/app/src/server/event-sourcing/subscribers/subscriber.types.ts:80):谓词写错最多多跑一次任务,而不是静默吞掉一个副作用。

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

按符号名 grep 定位比按行号更抗上游漂移。路径相对克隆根。

5.1 事件溯源内核

主题文件路径符号名
运行时中枢:建事件存储、建唯一全局队列、注册管线platform/app/src/server/event-sourcing/eventSourcing.tsEventSourcingregistercreateGlobalQueue
管线声明的建造器(类型安全的链式 API)platform/app/src/server/event-sourcing/pipeline/staticBuilder.tsdefinePipelinewithFoldProjectionwithSubscriberwithProcessManager
十九条管线的统一接线platform/app/src/server/event-sourcing/pipelineRegistry.tsPipelineRegistry.registerAll
存事件 + 分发投影的主流程platform/app/src/server/event-sourcing/services/eventSourcingService.tsEventSourcingService.storeEvents
fold / map / 订阅者的分发与错误聚合platform/app/src/server/event-sourcing/projections/projectionRouter.tsProjectionRouter.dispatchregisterSubscriberdispatchToSubscribers
fold 的加载—折叠—落库与乱序重算platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.tsFoldProjectionExecutor.executeexecuteBatch
map 的无状态映射与追加platform/app/src/server/event-sourcing/projections/mapProjectionExecutor.tsMapProjectionExecutor.execute
fold 基类:类型推导出的处理器名、单调时间戳platform/app/src/server/event-sourcing/projections/abstractFoldProjection.tsAbstractFoldProjection.applyinit
fold 状态的 Redis 写穿缓存platform/app/src/server/event-sourcing/projections/redisCachedFoldStore.tsRedisCachedFoldStore
投影重放(改了 fold 逻辑后重算历史)platform/app/src/server/event-sourcing/replay/replayService.tsReplayService
subscriber 与 process manager 的契约platform/app/src/server/event-sourcing/subscribers/subscriber.types.tsplatform/app/src/server/event-sourcing/pipeline/processManagerDefinition.tsSubscriberSpecwhenProcessManagerDefinition

5.2 队列与存储

主题文件路径符号名
同组 FIFO、跨组并行的 Redis 队列platform/app/src/server/event-sourcing/queues/groupQueue/groupQueue.tsGroupQueueProcessor
分层分组键、各类队列的注册platform/app/src/server/event-sourcing/services/queues/queueManager.tsQueueManager.buildGroupKeyinitializeProjectionQueues
ClickHouse 事件存储(含保留期打标)platform/app/src/server/event-sourcing/stores/eventStoreClickHouse.tsEventStoreClickHouse
事件表的读写 SQLplatform/app/src/server/event-sourcing/stores/repositories/eventRepositoryClickHouse.tsinsertEventRecordsgetEventRecords
核心表结构:event_log + 派生表(后续数十个 migration 另加了 suite_runs、dspy_steps、网关预算账本等)platform/app/src/server/clickhouse/migrations/00002_create_schema.sqlevent_logstored_spanstrace_summaries

5.3 trace 与评估

主题文件路径符号名
OTLP 三个入口(traces / logs / metrics)platform/app/src/server/routes/otel.tsotelIngestAuthenforcePlanLimitpeekCustomerTraceIds
REST 采集入口(老协议)platform/app/src/server/routes/collector.tscollectorRESTParamsValidatorSchemaparamsMD5
逐 span 的过滤、去重、派命令platform/app/src/server/app-layer/traces/trace-request-collection.service.tshandleOtlpTraceRequestingestNormalizedSpan
span → 事件(脱敏 / 成本 / token / 甩包)platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.tsRecordSpanCommand.handleRECORD_SPAN_DEDUPLICATION
投影用的瘦身事件与甩包阈值platform/app/src/server/app-layer/traces/lean-for-projection.tsleanForProjectionCOMMAND_INLINE_THRESHOLD
trace 汇总 fold(成本、时延、输入输出、处理上限)platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/traceSummary.foldProjection.tsTraceSummaryFoldProjectionMAX_PROCESSED_SPANS
span 落库 map 投影platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/spanStorage.mapProjection.tsSpanStorageMapProjection
评估触发器(含超大 trace 与因果环护栏)platform/app/src/server/event-sourcing/pipelines/trace-processing/subscribers/evaluationTrigger.subscriber.tscreateEvaluationTriggerSubscribercausalityLoopGuardFired
评估执行命令(前置条件、采样、跑分、落库)platform/app/src/server/event-sourcing/pipelines/evaluation-processing/commands/ExecuteEvaluationCommand
评分器实现库(Python)sdks/python/services/langevals/各 provider 子包

5.4 Go 侧服务

主题文件路径符号名
单二进制多服务入口cmd/service/main.goservices map、ServiceBoot
网关拦截器链的组装services/aigateway/app/app.goApp.buildInterceptors
洋葱式拦截器抽象services/aigateway/app/pipeline/pipeline.goInterceptorBuildPreOnly
各层治理逻辑services/aigateway/app/pipeline/RateLimitPolicyCacheBudgetGuardrailTrace
裸调度器(供内部服务当库用)services/aigateway/dispatcher/dispatcher.goDispatcher.DispatchPassthrough
工作流 DAG 引擎(分层调度,层内并发)services/nlpgo/app/engine/engine.goEngineExecuterunLayer
DAG 的校验、裁剪与分层services/nlpgo/app/engine/planner/planner.goNewlayerizereachableFromEntry
各类节点执行器services/nlpgo/app/engine/blocks/agentblockcodeblockevaluatorblockhttpblock

5.5 仿真

主题文件路径符号名
真实入口:编排仿真执行的 process managerplatform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.tshandleRunQueuedsimulationRunExecutionWake
并发闸门 + 子进程登记表 + 取消标记platform/app/src/server/scenarios/execution/execution-pool.tsScenarioExecutionPoolsubmitdequeueNext
子进程编排:预取数据、spawn、超时与收尾platform/app/src/server/scenarios/scenario.processor.tsexecuteScenarioRun
子进程环境变量组装(密钥注入等)platform/app/src/server/scenarios/execution/child-environment.tsbuildChildProcessEnv
子进程入口:跑 SDK 剧本、flush OTELplatform/app/src/server/scenarios/execution/scenario-child-process.tsexecuteScenarioflushOtelTraces
仿真跑的状态 foldplatform/app/src/server/event-sourcing/pipelines/simulation-processing/projections/simulationRunState.foldProjection.tsSimulationRunStateFoldProjection

引用约定: 本文所有 path:line 均相对克隆根,锚定在 sourceCommit: 7fb95fe 这次提交。行号会随上游漂移,符号名通常不会——定位不到时请用表中的符号名 grep