跳到主要内容

数据截至 (上游 commit 7fb95fe9048f)

事件溯源内核:命令、事件、投影、订阅者

30 秒导读: LangWatch 的服务端把所有状态变更写成不可变事件,再由三种"投影原语"把事件流算成可查询的视图和外部副作用。本章只讲这套通用内核——它长什么样、怎么装配、怎么排队、失败了怎么办。具体业务管线(trace、评估、仿真)分别在 030406

术语变更(2026-08,ADR-098): 这套内核的"投影落库成功之后触发的副作用挂钩"以前叫反应器(reactor),现在正式定名为订阅者(subscriber)——机器没换,名字换了(dev/docs/adr/098-post-event-work-subscribers-and-process-managers.md 写明:fold/map 绑定的订阅者会被编译成与旧反应器完全相同的内部注册、相同的队列路径与去重语义)。同时新增了一个更硬的原语 process manager(带持久状态与精确一次消费的编排器)。本章统一用新词,历史文献里看到 reactor 指的就是 subscriber。


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

一句话定义: 一个内置在 LangWatch Next.js 服务里的事件溯源框架——不直接改状态表,而是先追加一条"发生过什么"的事件,再由下游把事件算成状态。

解决什么问题。 假设一条 trace 的 span 是分批、乱序、可能重复地飞进来的,而你要同时做四件事:更新 trace 汇总、把 span 落到列存、触发评估、给前端推实时更新。

如果四件事各自去改数据库,你会立刻撞上三堵墙:

具体表现
竞态两个 worker 同时改同一条 trace 的汇总,后写覆盖先写
不可复现汇总算错了,原始输入已经没了,没法重算
耦合加一个下游消费者,就要改一次写入路径

事件溯源把这三堵墙一次拆掉:唯一的写入是"追加一条事件",其它全是从事件派生出来的。

它能做什么(内核层面):

  • 定义命令(意图)和它产出的事件,带 Zod schema 校验
  • 把事件流累积成聚合状态(fold),或逐条映射成记录(map)
  • 在状态落库成功之后触发外部副作用(订阅者,subscriber;尽力而为)
  • 用 process manager 承载"丢一次不可接受"的有状态编排(精确一次收件箱 + 事务性发件箱)
  • 保证同一个聚合的事件严格按序处理,不同聚合并行处理
  • 出错自动重试、退避、最终隔离;需要时把历史事件重放一遍重建投影

用起来什么样。 一条完整管线就是一条链式调用,下面是仿真管线的真实骨架(platform/app/src/server/event-sourcing/pipelines/simulation-processing/pipeline.ts:94-152):

// 示意,非源码(结构与真实 pipeline.ts 一致,删去了依赖注入细节)
const pipeline = definePipeline<SimulationProcessingEvent>()
.withName("simulation_processing") // 管线名,全局唯一
.withAggregateType("simulation_run") // 聚合类型,事件按它分区
.withFoldProjection("simulationRunState", stateProjection) // 有序累积状态
.withSubscriber("suiteRunSync", { ... }) // 落库后同步套件状态
.withProcessManager(simulationRunExecutionProcess) // 有状态编排(执行/取消/失速)
.withCommand("queueRun", QueueRunCommand) // 一条命令 = 一种意图
.withCommand("finishRun", FinishRunCommand)
.build();

const registered = eventSourcing.register(pipeline); // 接上 ClickHouse/Redis
await registered.commands.queueRun.send({ tenantId, occurredAt, runId });

一句话直觉: 把事件存储当账本(只能追加、永不涂改),把投影当报表(随时可以按账本重新算出来)。账本是真相,报表只是方便查询的缓存。


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

怎么读这张图: 从上到下是一条命令的完整旅程;事件存储是分水岭,它上面是"写",下面全是"派生"。

commands.queueRun.send(payload)
│ (进队列,不是直接调用)

┌──────────────────┐
│ 命令处理器 │ Zod 校验意图 → 产出事件
│ CommandHandler │
└────────┬─────────┘

┌──────────────────┐
│ 事件存储 │ 只追加;ClickHouse 或内存
│ event_log │
└────────┬─────────┘
│ 投影路由器按事件类型分发
┌────────────┼─────────────┬──────────────┐
▼ ▼ ▼ ▼
┌─────────┐ ┌─────────┐ ┌──────────┐ ┌──────────────┐
│ fold 投影│ │ map 投影 │ │ 全局投影 │ │ process │
│ 有序累积 │ │ 无序并行 │ │ 跨管线 │ │ manager │
└────┬────┘ └────┬────┘ └──────────┘ └──────────────┘
│ 落库成功 │ append 成功
▼ ▼
┌────────────────┐
│ 订阅者 │ 广播 / 触发评估 / 指标回填
└────────────────┘

部件一句话职责:

部件干什么文件(相对克隆根)
Event / TenantId / AggregateType内核的三个基础类型,全部带 Zod 校验platform/app/src/server/event-sourcing/domain/
CommandHandler校验意图,返回要落库的事件数组platform/app/src/server/event-sourcing/commands/command.ts
definePipeline()静态装配 DSL,不碰任何运行时依赖platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts
EventSourcing中央类:持有事件存储、唯一一条全局队列、全局作业注册表platform/app/src/server/event-sourcing/eventSourcing.ts
EventSourcingService存事件 + 把事件交给路由器分发platform/app/src/server/event-sourcing/services/eventSourcingService.ts
ProjectionRouter三种投影原语的执行与订阅者派发platform/app/src/server/event-sourcing/projections/projectionRouter.ts
QueueManager给每种作业类型造一个"门面",往全局队列注册路由元数据platform/app/src/server/event-sourcing/services/queues/queueManager.ts
GroupQueueProcessor分组队列:每组 FIFO、跨组并行,靠 Redis Lua 实现platform/app/src/server/event-sourcing/queues/groupQueue/groupQueue.ts
ReplayServiceevent_log 重放事件,重建 fold / map 投影platform/app/src/server/event-sourcing/replay/replayService.ts

主线走一遍(不进代码):

  1. 业务代码调 commands.xxx.send(payload) —— 这只是入队,不是执行。
  2. worker 进程从队列里取出这条命令,用 Zod 校验,调 handle() 拿到事件数组。
  3. EventSourcingService.storeEvents()(platform/app/src/server/event-sourcing/services/eventSourcingService.ts:308)先把事件写进事件存储(必须成功)。
  4. 写成功后,事件被瘦身(leanForProjection,eventSourcingService.ts:365)再交给投影路由器。
  5. 路由器按事件类型筛选,把事件再次入队分发给每个 fold / map 投影。
  6. fold 落库成功 → 派发订阅者作业;map 追加成功 → 派发挂在 map 上的订阅者作业。

注意第 3 步和第 5 步之间:写事件失败会抛,分发失败只记日志(eventSourcingService.ts:306 的流程注释:"Events are dispatched to all projections via ProjectionRouter - errors are logged but don't fail")。事件存储是唯一的强一致点,投影的补偿手段是重试和重放。


3. 命令层:意图怎么变成事件

这一节讲"写"这一侧:一个 payload 怎么变成一条带完整身份的不可变事件。

3.1 事件的形状

事件是内核里唯一的不可变事实,字段全部由 EventSchema 约束(platform/app/src/server/event-sourcing/domain/types.ts:44-67):

字段含义为什么重要
id纯 KSUID可按字典序排序 = 天然的写入顺序游标
aggregateId / aggregateType这条事件属于哪个聚合、什么类型数据库按 tenantId + aggregateType 分区
tenantId租户品牌类型,见 §7.3
createdAt写入时刻(ms)存储层排序、重放游标
occurredAt业务发生时刻(ms)fold 的排序基准、乱序检测基准
type事件类型串投影路由的唯一依据
versionYYYY-MM-DD 日历日期事件数据 schema 的版本
idempotencyKey可选存储层去重

createdAtoccurredAt 必须分清,这是整个内核最容易踩的一处:队列里作业的排序分数、乱序重折叠判定都只看 occurredAt(见 §4.1),与写入时刻无关。

事件类型走一套四段分类法,写在 domain/types.ts:41-42 的注释里:<来源>.<领域>.<聚合类型>.<具体标识>,例如 lw.obs.trace.span_received。聚合类型是第三段,且不是事件类型——AggregateTypeSchema 是一个从 schemaTypeIdentifiers 的常量数组生成的枚举(platform/app/src/server/event-sourcing/domain/aggregateType.ts:18),想加新聚合必须先改那张表。

3.2 命令 = 意图,不是事实

Command 的形状很薄:tenantId + aggregateId + type + data + 可选 metadata(platform/app/src/server/event-sourcing/commands/command.ts)。

处理器接口只有一个必需方法(commands/command.ts:69-97):

// 真实接口,commands/command.ts:85-97
handle(command: TCommand): CommandHandlerResult<EventType>;
cleanupAfterStore?(command: TCommand): Promise<void>;

cleanupAfterStore 是个值得记的钩子:它在事件写入持久化成功之后才被调用,专门放"绝不能在落库前跑"的尽力而为副作用(注释里举的例子是删除临时的 S3 spool)。注释还点破了一个隐藏陷阱——原始命令要作为参数传进去,不能存成实例字段,因为处理器实例在并行队列作业之间是共享的。

3.3 静态契约:类的静态成员就是配置

框架不要求你实现某个基类,而是从类的静态属性上读配置(platform/app/src/server/event-sourcing/commands/commandHandlerClass.ts:14-46):

静态成员必需作用
schema命令类型 + payload 的 Zod 校验器
getAggregateId(payload)决定事件落到哪个聚合
getGroupKey(payload)覆盖队列分组键,把"按聚合串行"换成更细的并行粒度
getSpanAttributes(payload)可观测性属性
dispatcherName覆盖注册时的名字

CommandHandlerClass 的类型本身是"静态契约 ∩ 零参构造签名"(commandHandlerClass.ts:96-101),所以框架能直接 new handlerClass()(platform/app/src/server/event-sourcing/services/queues/queueManager.ts:540 的注释也点明预构造实例优先于它)。需要构造函数注入依赖时,改用 withCommandInstance() 传一个预构造实例(platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts:503)。

3.4 defineCommand:一条命令 = 一条事件的语法糖

绝大多数命令只做一件事:把 payload 原样变成一条事件。defineCommand() 把这件事压成了配置(platform/app/src/server/event-sourcing/commands/defineCommand.ts:49-128):

// 示意,非源码(参数名与 defineCommand.ts 一致)
export const StartSuiteRunCommand = defineCommand({
commandType: "lw.suite_run.start",
eventType: "lw.suite_run.started",
eventVersion: "2026-03-01",
aggregateType: "suite_run",
schema: suiteRunStartedEventDataSchema, // 事件数据 schema 才是真相
aggregateId: (d) => d.batchRunId,
idempotencyKey: (d) => `${d.tenantId}:${d.batchRunId}:${d.idempotencyKey}`,
});

关键设计在信封(envelope):命令 payload = 事件数据 schema 合并三个信封字段 tenantId / occurredAt / idempotencyKey(platform/app/src/server/event-sourcing/commands/commandEnvelope.ts:20withCommandEnvelope)。handle() 再把信封剥掉,只把纯业务数据写进事件(defineCommand.ts:105stripEnvelope)。

这条约定的收益很实在:事件数据 schema 是唯一真相,命令 schema 自动派生,两者永远不会漂移。

3.5 派发端:命令看起来只是个 async 函数

register() 返回的对象上,commands.xxx 是队列门面。mapCommands() 再把它们拍平成普通异步函数(platform/app/src/server/event-sourcing/mapCommands.ts:17),于是调用方完全看不见队列的存在。

组合根就是靠它把十九条管线的命令暴露成一组扁平 API(platform/app/src/server/event-sourcing/pipelineRegistry.ts:456registerAll(),里面 19 次 eventSourcing.register(...))。


4. 三种投影原语 + 两种事后原语

这一节是内核的核心:同样一条事件,消费方式不同,保证完全不同。

先看总表——挑哪个原语,基本由这几行决定:

fold 投影map 投影订阅者process manager
有无累积状态无(读 fold/map 的结果)有(每聚合持久状态)
顺序保证同聚合严格 FIFO每键有序
生命周期getapplystoremapappend上游落库成功后才触发事件 → 收件箱 → 纯 handler 演进状态
纯度要求apply 必须纯函数map 必须纯函数就是干副作用的handler 纯;副作用走发件箱
失败后果整条事件重试该条重试独立重试,不影响已落库状态发件箱租约重试,有次数上限
定义platform/app/src/server/event-sourcing/projections/foldProjection.types.ts:33platform/app/src/server/event-sourcing/projections/mapProjection.types.ts:29platform/app/src/server/event-sourcing/subscribers/subscriber.types.ts:60platform/app/src/server/event-sourcing/pipeline/processManagerDefinition.ts

4.1 fold 投影:有状态、有序

要解决的小问题: 一条 trace 的汇总(span 数、总耗时、模型列表)必须随每条 span 事件增量更新,而且不能因为两条事件并行处理而算错。

思路: 把它写成一个左折叠——init() 给初始状态,apply(state, event) 是纯函数,存储只负责读写。框架保证同一聚合的事件一条一条来。

# 示意,非源码:fold 的本质
state = store.get(aggregate_id) or projection.init() # 1. 取
state = projection.apply(state, event) # 2. 纯函数算
store.store(state) # 3. 存

真实执行器就是这三步,加上一段乱序保护(platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.ts:355 起,FoldProjectionExecutor.execute)。

乱序重折叠(re-fold)是这里最硬的一段。 状态里有一个 LastEventOccurredAt 字段记录"我见过的最大 occurredAt"。执行器比较新事件与它:

// 真实逻辑,foldProjectionExecutor.ts:450-462(节选)
eventOccurredAt < prevLastOccurred // 严格小于才算乱序;相等不触发

判定为乱序时,执行器用 eventLoader 把该聚合的历史事件拉回来,从 init() 重放一遍。相等触发重折叠——源码注释说得很直白:同一逻辑时刻的两条事件(比如 SDK 同时发 snapshot 和 finished),到达顺序就是正确的裁决者(foldProjectionExecutor.ts:450-455)。

eventLoader 不需要投影自己提供:注册时框架自动接上事件存储(platform/app/src/server/event-sourcing/services/eventSourcingService.ts:126-131)。没有 eventLoadercanRefold 返回 false,只记警告("cannot re-fold",foldProjectionExecutor.ts:53-65);另外 fold 还可以声明 refoldOnOutOfOrder: false 整体关掉重折叠(foldProjectionExecutor.ts:128 注释,trace 汇总等热 fold 现在靠它避开全量重算)。

AbstractFoldProjection 把时间戳管起来了。 继承它之后(platform/app/src/server/event-sourcing/projections/abstractFoldProjection.ts:101-234):

  • eventTypes 从 Zod schema 数组自动派生,不用手写
  • 事件类型 lw.suite_run.started 自动映射到方法名 handleSuiteRunStarted;漏写一个方法是编译错误(类型 FoldEventHandlers,:37-53),运行时构建 dispatch map 时再兜一层(:174-189)
  • initState() 的返回类型禁止包含时间戳字段(时间戳键清单在 :28-35),init() 自己补上
  • apply() 自动维护单调递增的 UpdatedAt:Math.max(Date.now(), prev + 1)(:233)

那个 +1 不是洁癖:同一毫秒内处理两条事件时,它保证 ClickHouse 里产生两行可区分的版本。

批合并(coalescing)。 当某个组积压时,执行器可以一次读、折叠 N 条、只写一次,把 O(n²) 的读写摊成 O(n)(platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.ts:519executeBatch)。默认上限 500 条(2026-08 从 100 提到 500,platform/app/src/server/event-sourcing/projections/projectionRouter.ts:71DEFAULT_FOLD_COALESCE_MAX_BATCH,理由写在 :55-70 的注释:积压组的派发周期减少 5 倍),投影可以设 coalesceMaxBatch: 1 关掉。

它之所以安全,是因为纯左折叠的最终结果与逐条应用完全一致——注释还强调"提高上限只改吞吐,永不改正确性"(projectionRouter.ts:64-70)。

4.2 map 投影:无状态、并行

要解决的小问题: 把每条 span 事件规范化成一行记录,追加到列存里。这件事既不需要历史,也不需要顺序。

生命周期只有两步(platform/app/src/server/event-sourcing/projections/mapProjectionExecutor.ts:32 起):

// 真实实现,精简自 mapProjectionExecutor.ts
const record = projection.map(event);
if (record === null) return null; // 返回 null = 显式跳过
await projection.store.append(record, context);

map 返回 null 是一等公民的"跳过"语义(platform/app/src/server/event-sourcing/projections/mapProjection.types.ts:40)。存储接口只有 append 一个方法——这正是"无序"的类型级体现:没有 get,就无从依赖历史。

4.3 订阅者:落库成功之后的副作用

要解决的小问题: trace 汇总更新完之后,要把它推给前端、触发评估、回填指标。这些都会失败、都会重试、都不该阻塞投影本身。

关键保证: 订阅者在投影落库成功之后才被派发。fold 路径见 platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1739-1799——先 foldExecutor.execute 拿到 foldState(在 :1731 的事务闭包里),成功了才 dispatchSubscribersAfterStore(:1791)。fold 失败 → 抛异常 → 订阅者永远不会触发,于是订阅者看到的状态一定是已经持久化的。

订阅者既可以挂在 fold 上,也可以挂在 map 上:声明形如 .withSubscriber("名字", { fold: "traceSummary", when, handler }),fold / map 字段指明挂在哪个投影(platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts:326-345,重名在 :288-293ConfigurationError)。map 执行成功时同样派发挂在它上面的订阅者(projectionRouter.ts:727 在 append 成功的闭包里执行)。真实例子:trace 管线的 spanStorageBroadcast 挂在 spanStorage map 上(platform/app/src/server/event-sourcing/pipelines/trace-processing/pipeline.ts:303-313);跨管线全局投影里的 billingMeterDispatch 挂在 orgBillableEventsMeter map 上(platform/app/src/server/event-sourcing/projections/global/billingMeterDispatch.subscriber.ts:37 的注释)。

订阅者的配置项都在 SubscriberDispatchOptions 里(platform/app/src/server/event-sourcing/subscribers/subscriber.types.ts:33-50):

选项作用
delay派发前延迟多少毫秒——配合去重就是"防抖"
makeJobId去重键;有它才启用去重
ttl去重窗口(毫秒)
deduplication完整的 GroupQueue 去重契约(withSubscriber 会把两条路指向同一个键函数,防止漂移)
runIn限定只在某些进程角色跑
groupKeyFn覆盖分组键
killSwitch / disabled开关

delay + makeJobId 这个组合是内核里最实用的一招:高频事件(一条 trace 的几百个 span)会把同一个订阅者作业反复压成同一个,最后只真正跑一次。

还有一个纯函数谓词 when(声明侧)/ shouldDispatch(内部契约),在入队之前决定要不要派发。它的错误策略值得抄:失败即放行——谓词抛异常会被记日志然后当作 true(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1996-2002subscriberShouldDispatch;契约注释在 subscriber.types.ts:71-88),因为"多跑一次冗余作业"远比"静默丢掉一个副作用"便宜。ADR-026 定下的这条契约在改名后原样保留(ADR-098 明说 shouldReact 契约平移为 when)。

批合并模式下,订阅者派发按 jobId 折叠:去重键带事件 id 的(如 customEvaluationSync 要逐条读 span)仍每条事件各派一次;聚合级键的(广播、告警)合并成一个任务,免掉 N 次 serialize+gzip+blob 往返(projectionRouter.ts:1957-1969collapseByJobId 注释)。

4.4 process manager:丢一次不可接受的编排

订阅者是尽力而为的:作业暂存上队列之后才保证重试,进程死在"fold 落库与暂存之间"的那条缝里,这次副作用就静默丢了(ADR-098 把这一点写成退役旧词汇的理由之一)。stake-sensitive(赌注敏感)的活——比如"看到 queued 事件就把这条仿真 run 编排进执行池、处理取消、检测失速"——交给 process manager(.withProcessManager(...),platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts:421):

  • 每个聚合键一行持久状态,纯事件 handler 演进状态、产出意图(intent);
  • 意图经事务性发件箱派发,带租约、重试、次数上限;
  • 精确一次收件箱(exactly-once inbox)去重消费;
  • wake 调度:状态可以声明"到点叫醒我"(比如失速检测的宽限期)。

仿真执行就是它的标杆用例(platform/app/src/server/event-sourcing/pipelines/simulation-processing/process-manager/simulationRunExecution.process.ts:206handleRunQueued:318simulationRunExecutionWake),细节见 06。全库现有十三个 process manager 分布在九条管线上(ADR-098 的决策记录)。


5. 装配:从静态 DSL 到运行时

这一节讲那条链式调用背后到底发生了什么。

5.1 分阶段类型:错误在编译期就挡住

definePipeline() 返回的不是一个大对象,而是分阶段接力的类型(platform/app/src/server/event-sourcing/pipeline/staticBuilder.ts):

definePipeline<E>() StaticPipelineBuilder
│ .withName(name)

StaticPipelineBuilderWithName
│ .withAggregateType(type)

StaticPipelineBuilderWithNameAndType
│ .withFoldProjection / .withMapProjection
│ .withSubscriber / .withProcessManager / .withEventSubscriber
│ .withCommand / .withCommandInstance / .withFeatureFlagService
▼ .build()
StaticPipelineDefinition

前两个阶段的 build() 直接抛错(staticBuilder.ts:76-77:100-101),所以"忘了写名字"这类失误不会跑到运行时。

泛型参数会随每次 withXxxProjection 累积投影名字面量类型,withSubscriberfold 字段被约束成已注册的投影名(staticBuilder.ts:315-322)——挂到不存在的投影上是类型错误,运行时还会再抛一次 ConfigurationError(:379registerProjectionSubscriber)。重名同样在构建期抛(:190-192)。

build() 只是把 Map 和数组打包成一个纯数据对象,外加一个零运行时开销commandRegistry 空对象——它唯一的作用是让下游能从类型上推断出命令名和 payload(:611 的注释)。

5.2 register():接上真实基础设施

EventSourcing.register() 是静态定义与运行时的接缝(platform/app/src/server/event-sourcing/eventSourcing.ts:240-315):

  1. 事件溯源被禁用或没有事件存储 → 返回 DisabledPipeline,命令被静默丢弃(:253-269)
  2. buildServiceOptions() 把 Map 拍平成数组(:735)——注释特意提醒不能用展开语法,因为 eventTypes 这类 getter 活在原型上,{...obj} 会丢掉
  3. 构造 EventSourcingPipeline,它内部构造 EventSourcingService(platform/app/src/server/event-sourcing/runtimePipeline.ts:13)
  4. 从服务里取出命令队列,挂成 .commands 字段

5.3 一条全局队列,不是每种作业一条

这是内核最反直觉、也最省资源的一个决定:整个进程只有一条 Redis 队列,名字是 makeQueueName("event-sourcing/jobs")(platform/app/src/server/event-sourcing/eventSourcing.ts:529)。

那不同作业怎么区分?靠路由元数据 + 全局作业注册表:

QueueManager.createFacade() ──▶ 给每种作业造一个门面
│ 门面在 send 时注入三个字段:
│ __pipelineName / __jobType / __jobName

globalJobRegistry["管线:类型:名字"] = { process, groupKeyFn, scoreFn, ... }

│ 全局队列的每个回调都先 lookupEntry(payload) 找到条目
└── groupKey / score / process / processBatch / spanAttributes

代码上就是 createFacade()(platform/app/src/server/event-sourcing/services/queues/queueManager.ts:223)注册条目、注入元数据,lookupEntry()(eventSourcing.ts:391)反查条目、剥掉元数据。注册表里找不到就跳过并记警告——这正是滚动发布时旧作业的处理方式。

分组键是分层的,格式 ${tenantId}/${jobPath}/${domainKey}(platform/app/src/server/event-sourcing/services/queues/queueManager.ts:175-191buildGroupKey):

作业类型jobPath默认 domainKey
map 投影map/<名字><聚合类型>:<聚合ID>
fold 投影fold/<名字><聚合类型>:<聚合ID>
命令command/<名字><聚合类型>:<聚合ID 或 groupKey>
订阅者<fold|map>/<父投影>/reactor/<名字><聚合类型>:<聚合ID>

注意订阅者的 jobPath 里仍保留 reactor/——ADR-098 改名时刻意不动队列拓扑,改名前的已暂存作业才不会找不到家(queueManager.ts:799 拼这个路径;platform/app/src/server/event-sourcing/ARCHITECTURE.md:90 也注明 "the reactor segment" 是历史保留)。

租户 ID 放在第一段不是为了好看——分组队列的 Lua 脚本靠"第一个 / 之前的片段"识别租户,用来做每租户在途上限和暂停(见 §6.7)。


6. 排队机制(本章硬核)

前面所有的"有序""并行""重试",最终都落在这一个组件上:GroupQueueProcessor(platform/app/src/server/event-sourcing/queues/groupQueue/groupQueue.ts:303)。

上游自带的 platform/app/src/server/event-sourcing/ARCHITECTURE.md:143 现在也明确写着:"Not BullMQ"——分组队列基于 Redis 原语 + 自研 Lua,进程内并发限流用 fastq。真实实现:一套 Lua 暂存脚本(queues/groupQueue/scripts.ts,两千多行)+ fastq 处理槽位。

6.1 它到底要解决什么

一句话:同一个聚合的作业必须一条一条来,不同聚合可以随便并行。

朴素方案是分布式锁,但那会让每条 span 都去抢一次锁。分组队列换了个思路——每个组在 Redis 里有一把"活跃锁",拿到锁才能出队,处理完才还锁。锁本身就是 FIFO 的实现,不是额外开销。

6.2 Redis 里放了什么

类型装什么
<队列>:gq:readyZSET组 ID → 最早待派发时刻
<队列>:gq:group:<组>:jobsZSET暂存作业 ID → 待派发时刻
<队列>:gq:group:<组>:dataHASH暂存作业 ID → 作业信封
<队列>:gq:group:<组>:activeSTRING该组正在处理的作业 ID(带 TTL,就是那把锁)
<队列>:gq:group:<组>:errorHASH重试耗尽时的错误信息
<队列>:gq:blockedSET被隔离的组
<队列>:gq:dedup:<键>STRING去重键 → 暂存作业 ID(带 TTL)
<队列>:gq:signalLIST唤醒派发循环的信号
<队列>:gq:blob:<id>STRING卸载出去的大 payload(见 §6.6)

键名构造见 GroupStagingScripts 构造函数(platform/app/src/server/event-sourcing/queues/groupQueue/scripts.ts:1959 起)。队列前缀还是 Redis Cluster 的 hash tag,同队列所有键落同一个 slot,Lua 才能原子地摸多把键(queues/groupQueue/ARCHITECTURE.md:36)。

6.3 一条作业的生命周期

怎么读这张图: 左边是生产者,中间是派发循环,右边是 worker;方括号是 Lua 脚本。

send() 派发循环(BRPOP 唤醒) fastq worker
│ │ │
▼ ▼ ▼
[STAGE / STAGE_BATCH] ──▶ ready 有序集 ──▶ [DISPATCH_BATCH] ──▶ 拿到组的 active 锁 ──▶ 处理
│ 扫描:到期 + 未阻塞 │
│ + 无 active 的组 ├─成功──▶ [COMPLETE] 还锁 + 重排
└─去重命中则原地 squash ├─失败──▶ [RETRY_RESTAGE] 退避重排
└─耗尽──▶ [RESTAGE_AND_BLOCK] 组隔离

入队(send/sendBatch,groupQueue.ts:679:808): 算出组 ID、暂存作业 ID、派发时刻 score + delay,把 payload 编码成信封,调 scripts.stage() / scripts.stageBatch()。单发与批量共用同一套暂存不变量。

派发(GroupQueueDispatcher,platform/app/src/server/event-sourcing/queues/groupQueue/dispatcher.ts:15): 一个 BRPOP 阻塞循环,拿到信号就连续派发直到派不动。它的 BRPOP 超时会按最早到期时间收紧(dispatcher.ts:111-120nextWakeTimeoutSec),这样延迟作业不用干等一个固定轮询周期。

出队的核心不变量DISPATCH_BATCH_LUA 里(scripts.ts:909 起;单条派发已并入批量脚本):扫描时只有 active不存在且不在 blocked 集里的组才会被取队首(:892-893),取到之后立刻 SET activeKey ... EX activeTtlSec(:1084),并把该组在 ready 集里的分数改成一个未来时刻(:1086)以抑制重复派发。

完成(COMPLETE_LUA,scripts.ts:1244): 校验 active 键还是自己的作业 ID(:1276,不是就说明已被接管),删锁(:1281),然后按该组还剩没剩作业决定重排 ready 分数还是把组移出 ready。

崩溃恢复现在还有"死亡计数"。 worker 认领组时会 recordClaim,连续确认死亡(认领后心跳消失)达到阈值就触发毒丸守卫、隔离该组(groupQueue.ts:973 的错误文案:"Poison guard: N confirmed worker deaths while this group was in flight")。

6.4 主 Lua 脚本,各守一个不变量

脚本位置(scripts.ts)守住的不变量
STAGE_LUA:617组处于 active 或 blocked 时,不许改 ready 分数(:662,那时分数归别人管)
STAGE_BATCH_LUA:754批量入队复用同一套不变量
DISPATCH_BATCH_LUA:909一个组同时最多一个作业在飞
DRAIN_GROUP_LUA:1185批合并时排空到期兄弟作业,原子地摘出 jobs+data
COMPLETE_LUA:1244只有锁的主人能还锁
REFRESH_LUA:1345心跳只续自己的锁;组已被隔离就不再塞回 ready
RESTAGE_AND_BLOCK_LUA:1390隔离组:出 ready、进 blocked、清 TTL(等人工处理)
RETRY_RESTAGE_LUA:1596退避期间继续持锁,保住 FIFO

同一个文件还有会被拼进主脚本的 helper:租户在途计数(TENANT_ACTIVE_HELPER_LUA,:102)、路由元数据(gqRoutingMeta,:395)等;blob 生命周期另有独立脚本(blobDeleteLua.tsblobGraceLua.tsblobSweepLua.ts)。

用 Lua 而不是 MULTI 的理由很直接:这些操作都要先读后判再写(读 active 键决定要不要改 ready、读 ZRANK 决定 squash 还是新增),事务做不到,脚本可以。

6.5 去重:原地压扁,不是丢弃

STAGE_LUA 的去重逻辑(scripts.ts:683-708)分两种情况:

  • 去重键指向的作业还在暂存层原地 squash:复用旧的作业 ID,按 extend 决定要不要顺延派发时刻,按 replace 决定用新 payload 还是保留旧的,待处理计数净变化为零
  • 去重键指向的作业已经派发了 → 键是陈旧的,删掉,把新作业当全新的入队(:708 注释:把被顶掉的 payload 作为 orphaned value 返回,迟到的再触发也回收不了已派发的作业)。

第二条是个真正的 TOCTOU 修复:内存队列版本也做了同样的事(platform/app/src/server/event-sourcing/queues/memory.ts:220,注释直接写着 "same TOCTOU fix as GroupQueue Lua")。

默认去重窗口 200ms(groupQueue.ts:185DEFAULT_DEDUPLICATION_TTL_MS)。门面层会给去重 ID 加上 管线/类型/名字 命名空间前缀防跨管线撞车(platform/app/src/server/event-sourcing/services/queues/queueManager.ts:252-261)。

squash 会丢掉一个 payload,如果它带着卸载出去的 blob,那个 blob 就可能成孤儿。所以 Lua 把被顶掉的值作为 orphanedValue 返回(scripts.ts:683-703),调用方按需回收/转移租约;blob 另有 TTL 与清扫兜底(见下节)。

6.6 大 payload 不走 Lua:两级信封 + 分层 blob

信封现在有两代格式(platform/app/src/server/event-sourcing/queues/groupQueue/jobEnvelope.ts:142-176 的头注释):

格式阈值体放哪
内联≤ 4 KiB信封内(1 KiB 以上才尝试 gzip+base64,且必须真的变小,:155)
GQ1> 32 KiB独立 randomUUID() Redis blob 键(:160),7 天兜底 TTL(queues/groupQueue/blobConstants.ts:14)
GQ2> 4 KiB 且配置了 tiered store内容寻址、按租户命名空间的分层存储:先 Redis、后对象存储(S3);相同字节只存一份(ADR-029,jobEnvelope.ts:168-176;4 天兜底 TTL + 3 天租约,blobConstants.ts:26-34)

GQ2 的内容寻址去重正对事件溯源的扇出场景:同一条事件要喂给 N 个投影/订阅者,以前是 N 份大 payload 各自暂存,现在 N 个信封指向同一份 blob。blob 的读写由客户端直连完成(redisJobBlobStore.tstieredBlobStore.ts),Lua 脚本只碰小字符串。

头里冗余存了 __pipelineName/__jobType/__jobName 三个路由字段,这样 Lua 侧的暂停判定可以只切头、不碰体(scripts.ts:395gqRoutingMeta)。

信封写开关 GROUP_QUEUE_ENVELOPE_WRITES_ENABLED 仍是分阶段灰度:读侧永远认识信封,写侧要等全队列的消费者都升级(jobEnvelope.ts:172-181)。

6.7 心跳与崩溃恢复

active 锁带 300 秒 TTL(dispatcher.ts:157-167GROUP_QUEUE_CONFIG.activeTtlSec),纯粹是崩溃兜底。作业跑着的时候,心跳会按 TTL 的三分之一周期去续期(同段注释)。

心跳脚本还顺手做两件事(REFRESH_LUA,scripts.ts:1345):把该组暂存兄弟作业的安全网 TTL 也续上(:1355-1367),以及只续自己的锁(:1360-1362)——别人的锁碰不得。

每租户在途计数用的是 ZSET 而不是计数器,TENANT_ACTIVE_HELPER_LUA 的注释(scripts.ts:95-135)解释得非常好:标量计数器在 worker 非优雅死亡时只会往上漏、永不自愈;改成 ZSET 后,每个在途槽位携带与 activeKey 心跳同步的过期分数,心跳一停,分数过期,槽位自动不再计数。注释里保留了 2026-05 的 ElastiCache 事故作为设计动机。

6.8 失败阶梯

失败分三级处理(groupQueue.ts:1398-1451 一带):

失败 ──▶ 可重试 且 attempt < 25 ──▶ RETRY_RESTAGE:退避后重排
│ (退避 500ms 起翻倍,封顶 600s)

├──▶ 不可重试(CRITICAL 分类) ──┐
└──▶ 重试次数耗尽 ──┴──▶ RESTAGE_AND_BLOCK:组隔离 + 存错误

25 次尝试的累计等待约 2 小时 27 分(platform/app/src/server/event-sourcing/queues/shared.ts:10-23JOB_RETRY_CONFIG 与注释)——注释说这是为了扛过一次 ClickHouse 滚动重启,而失败的作业一直躺在 Redis 有序集里,长退避不丢数据,只是把运维负担换成了自动恢复

重试的实现细节值得注意:它不是在 worker 槽位里 sleep,而是重新入队一个未来分数的作业,立刻释放并发槽位(shared.ts:4-5 注释)。同时 active 锁的 TTL 被设成退避时长,所以退避期间该组仍然被锁着,FIFO 不会被后面的事件插队(RETRY_RESTAGE_LUA)。另外失败错误现在可以自带最小退避下限(queues/dispatchError.ts:21-23),调度器的指数退避把它当地板。

批合并遇到失败还有一层补偿:被排空的兄弟作业会按原分数重新入队(groupQueue.ts:1082-1097drainGroupReady 返回 drainedSiblings,失败时逐个还回去),因为批量折叠只在最后落一次库,失败意味着这些兄弟的工作全白做了。

6.9 没有 Redis 怎么办

EventSourcing.createGlobalQueue()(eventSourcing.ts:528)在没有 Redis 时换成 EventSourcedQueueProcessorMemory

它不是简化版的分组队列,而是一个语义不同的东西:send() 返回的 Promise 直到该作业处理完才 resolve(platform/app/src/server/event-sourcing/queues/memory.ts:166-174),所以调用方 await 一下就等于同步执行。它保留了去重 squash 语义(:220),但没有分组、没有 FIFO、没有持久化、没有重试。够测试和本地开发用,生产不行。


7. 三件设计取舍

7.1 为什么不需要 checkpoint 存储

传统事件溯源系统都有一张 checkpoint 表,记录"每个投影处理到第几号事件了"。这套内核没有,理由分三条:

原语为什么不需要 checkpoint
fold持久化的状态本身就是 checkpoint——它的 LastEventOccurredAt 字段(platform/app/src/server/event-sourcing/projections/abstractFoldProjection.ts:34-35 的时间戳键)直接告诉系统进度到哪了
map无状态逐条追加,没有"进度"这个概念
订阅者落库之后才触发,失败独立重试,靠 makeJobId 去重

再加上顺序由队列而不是序号保证:分组队列的 active 锁已经实现了每聚合 FIFO,不需要序号追踪器。

代价也很明确:追加型存储必须能容忍重复。fold 的 store 失败会导致整条事件重试,重试时重新读状态、重新应用——所以存储要么幂等要么用 upsert 语义。这一条写在 platform/app/src/server/event-sourcing/README.md:296 的"Fold store failures"里,是这套设计付出的真实成本。

7.2 processRole:进程分工

配置项有四个值(platform/app/src/server/app-layer/config.ts:3):

含义
web只派发命令、入队事件,不跑消费者
worker跑全部消费者(生产 worker 部署)
all单进程模式,web 进程内嵌 worker(仅开发,WORKERS_IN_PROCESS=1)
migration直接调 processCommand(),排除订阅者

真正决定"跑不跑消费者"的是 roleRunsWorkers()(config.ts:16-18):

// 真实代码,config.ts:16-18
export function roleRunsWorkers(role: ProcessRole | undefined): boolean {
return role === "worker" || role === "all";
}

全局队列的 consumerEnabled 就取它(eventSourcing.ts:657)。processRole 完全没设(undefined)时消费者不启动——但现在有显式的 "all" 值覆盖开发单进程场景,注释写明它永不用于生产(config.ts:9-14)。

有意思的是注册路径是所有进程都走的:EventSourcingService 构造时无条件初始化所有队列门面,注释写明理由——共享队列的 worker 必须认识每一种作业类型,才能派发它捡到的任何作业(platform/app/src/server/event-sourcing/services/eventSourcingService.ts:262)。所以 web 进程也要完整注册,只是不消费。

订阅者还有一层进程级过滤 runIn,由 roleSatisfiesRunIn 判定——没写 runIn 就到处跑;"all" 角色扮演所有身份,否则进程内嵌模式下 runIn: ["worker"] 的订阅者会全军覆没(config.ts:24-40 的注释)。

7.3 租户隔离在存储层怎么强制

这套内核把租户隔离做成了四层防御,一层一层往下收:

① 类型层:TenantId 是品牌类型,普通 string 赋不进去
platform/app/src/server/event-sourcing/domain/tenantId.ts:6-12 .brand<"TenantId">()

② 构造层:createTenantId() 校验失败抛 SecurityError(不是普通 Error)
domain/tenantId.ts(import 自 services/errorHandling)

③ 操作层:每个存储操作入口先 validateTenantId(context, "操作名")
platform/app/src/server/event-sourcing/stores/abstractEventStore.ts:340

④ 数据层:批内每条事件都要与上下文租户一致,才允许写
stores/abstractEventStore.ts:532 validateEventTenant(event, context, i)

第 ④ 层是最容易被忽略的一环:storeEvents() 逐条校验批内每个事件的租户和聚合类型都与上下文匹配,任何一条不符就整批拒绝(abstractEventStore.ts:532 附近)。

隔离还渗进了键的构造:

  • fold 缓存键是 fold:<前缀>:<租户>:<聚合ID>(platform/app/src/server/event-sourcing/projections/redisCachedFoldStore.ts:350)
  • 队列分组键第一段就是租户(queueManager.tsbuildGroupKey)
  • 重放的内存累加器用 <租户>::<投影键> 做复合键,注释明说是"ensures tenant isolation"(platform/app/src/server/event-sourcing/replay/replayExecutor.ts:50-54)
  • 批量写入时,只有租户和保留策略都一致的条目才走原生批插,否则退化成逐条写(platform/app/src/server/event-sourcing/projections/repositoryFoldStore.ts:70-90,注释:"a mixed batch must fall back to")

fold 存储的缓存层还有一个漂亮的顺序约定(redisCachedFoldStore.ts:93 的 write-through 注释):先写 ClickHouse(失败就抛),再写 Redis 缓存(失败只记日志)。这个顺序不用事务就保证了正确性——CH 失败则缓存没动、事件会重试;CH 成功而缓存失败,下次读回落到 CH。


8. 重放:用事件重建投影

为什么需要。 fold 的 apply 改了逻辑、状态字段加了一个、某次事故写坏了投影——这时候事件账本还在,报表可以重算。

入口与通道。 ReplayService(platform/app/src/server/event-sourcing/replay/replayService.ts:50)提供 discover()(找聚合)与 replay()(:90):fold 与 map 投影走同一个批量引擎(runFoldMapReplay,每批每个聚合只拉一次事件,逐投影同时重算),.withProjection() 声明的关系型状态投影走独立通道(replayStateLane,:129:先 pause + drain 该投影的队列,再从 canonical 事件确定性重建 Postgres 行)。重放累加器按租户×投影复合键聚合,内存只与批内聚合数成正比(replay/replayExecutor.ts:50-65FoldAccumulator;map 侧另有 MapAccumulator,:171)。

怎么和在线处理共存是重放最难的部分:重放在读历史事件的同时,新事件还在往里灌。方案是一组 Redis 标记:

做什么
① mark给这批聚合打上 pending 标记(platform/app/src/server/event-sourcing/replay/replayMarkers.ts:46markPendingBatch)
② 暂停该投影的队列消费、排空在飞作业重放引擎负责
③ cutoff定出分界事件,标记改成 <时间戳>:<事件ID>
④ replay分页拉事件,喂给累加器,批量落库
⑤ unmark删标记、恢复消费(replayService.ts:187cleanup)

在线侧看到标记会怎么做(platform/app/src/server/event-sourcing/projections/replayMarkerCheck.ts:30-41 的契约):

标记状态决定
无标记process —— 正常处理(单次 HGET,约 0.1ms)
pendingReplayDeferralError —— 走队列退避重试
事件 ≤ 分界点skip —— 重放会处理它
事件 > 分界点ReplayDeferralError —— 等重放做完再说

ReplayDeferralError 继承 RecoverableError(replayMarkerCheck.ts:14;基类在 platform/app/src/server/event-sourcing/services/errorHandling.ts:82),所以它天然落进分组队列的退避重试路径,不需要任何特殊处理。分界点格式 {timestamp}:{eventId},比较刻意对齐 ClickHouse 的 ORDER BY EventTimestamp ASC, EventId ASC(replayMarkerCheck.ts:51-56 的注释),这样 Redis 侧的判定和 CH 侧的查询边界完全一致。标记有 7 天 TTL 兜底(platform/app/src/server/event-sourcing/replay/replayConstants.ts:9),防止一次被放弃的重放永久卡住在线处理。

一个容易漏的一致性细节: 重放在 apply 之前也调 leanForProjection(platform/app/src/server/event-sourcing/replay/replayExecutor.ts:99,注释:"once at materialization, matching the live dispatch"),和在线分发用的是同一个函数(eventSourcingService.ts:365)——让重放和在线产出逐字节相同的投影状态。

在线的乱序重折叠也是一种"小重放",它有一个专门的成本优化(platform/app/src/server/event-sourcing/stores/rehydrationWindow.ts)。只有时间局部的聚合类型(TIME_LOCAL_AGGREGATE_TYPES,:21)才允许把 event_log 扫描下界卡在触发事件前 45 天(:43REHYDRATION_WINDOW_DAYS,理由注释在 :39),让 ClickHouse 剪掉更老的分区;其它聚合类型返回 undefined = 永不设界(:55-59)。

那份"允许清单"是这个优化的正确性契约:会跨越任意时间累积事件的聚合永远无界扫描;窗口选 45 天也是刻意保守,因为"太小的代价是静默丢事件 = 投影损坏",所以宁可偏大。


9. 巧妙之处(可以直接借鉴的)

① 状态即 checkpoint。 只要投影状态里带一个"我见过的最大 occurredAt",就同时拿到了进度追踪和乱序检测两个能力,一张 checkpoint 表都不用建(platform/app/src/server/event-sourcing/projections/abstractFoldProjection.ts:233 的单调 UpdatedAt 也是同一哲学的延伸)。

② 事件数据 schema 是唯一真相,命令 schema 自动派生。 信封合并 + 剥离两个纯函数(platform/app/src/server/event-sourcing/commands/commandEnvelope.ts:20:30)就消灭了"命令 schema 和事件 schema 漂移"这类 bug。

③ 事件类型 → 方法名的编译期强制。 handle${DotSnakeToPascal<StripPrefix<T>>} 这个类型体操(platform/app/src/server/event-sourcing/projections/abstractFoldProjection.ts:37-53)让"新增一种事件却忘了写处理方法"变成编译错误,而不是运行时静默跳过。

④ 谓词失败即放行。 when/shouldDispatch 抛异常时当作 true(platform/app/src/server/event-sourcing/projections/projectionRouter.ts:1996-2002)。判定逻辑的 bug 只会造成一次冗余作业,不会静默吞掉一个副作用——错误方向选对了。

⑤ 在途计数用带过期分数的 ZSET,而不是计数器。 计数器在进程非优雅死亡时只漏不还,ZSET 让每个槽位与心跳同生共死,几分钟内自愈,不需要人工重置(platform/app/src/server/event-sourcing/queues/groupQueue/scripts.ts:95-135)。这是从一次真实事故里长出来的设计。

⑥ 大 payload 绕开 Lua,而且内容寻址共享。 4 KiB 以上的体走 blob、由客户端直接读写,Lua 脚本永远只处理小字符串(platform/app/src/server/event-sourcing/queues/groupQueue/jobEnvelope.ts:142-176);GQ2 更进一步,相同字节只存一份、按租户命名空间分层到 S3(queues/groupQueue/blobConstants.ts)。同时头里冗余存路由字段,让 Lua 能"只切头不碰体"。

⑦ 一个全局队列 + 路由注册表。__pipelineName/__jobType/__jobName 三个字段代替 N 条队列(platform/app/src/server/event-sourcing/services/queues/queueManager.ts:223platform/app/src/server/event-sourcing/eventSourcing.ts:391),连接数、指标基数、运维面全部塌缩成一份。

⑧ 改名不动物理拓扑。 ADR-098 把 reactor 改叫 subscriber,但队列 jobPath 里的 reactor/ 段原样保留——已暂存的作业不该因为一次词汇清理而找不到回家的路(queueManager.ts:799)。词汇是给人看的,拓扑是给数据走的。


10. 边界与局限

它刻意不做的:

  • 不做快照(snapshot)。 没有"每 N 个事件存一次全量状态"的机制,重建只有"增量 apply"和"全量重放"两档。
  • 不做严格的一次性投递。 追加型存储必须容忍重复,这是 platform/app/src/server/event-sourcing/README.md:296 明列的陷阱。(订阅者是尽力而为;要一次性保证,升级用 process manager。)
  • 不做跨聚合事务。 分组队列只保证组内有序,跨聚合无任何顺序保证。
  • map 投影不做排序。 类型层面就没给 get,想按顺序处理只能用 fold。

已知会崩/会退化的地方:

场景后果依据
生产环境没配 ClickHouse事件溯源直接禁用,命令被静默丢弃platform/app/src/server/event-sourcing/eventSourcing.ts:253-269
processRole 没设消费者不启动,作业只堆积不处理platform/app/src/server/app-layer/config.ts:16-18eventSourcing.ts:657
乱序事件但投影没有 eventLoader(且未关 re-fold)只记警告,状态保持错的继续往下走platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.ts:53-65
订阅者内联执行失败(队列缺失时的回退)状态已落库但副作用丢失platform/app/src/server/event-sourcing/projections/projectionRouter.ts:2293,日志里明写了这一点
某组重试 25 次仍失败整组被隔离,需要人工解除platform/app/src/server/event-sourcing/queues/groupQueue/scripts.ts:1390queues/shared.ts:19-23
信封写开关(GROUP_QUEUE_ENVELOPE_WRITES_ENABLED)提前打开旧版本 pod 会丢弃带前缀的信封作业platform/app/src/server/event-sourcing/queues/groupQueue/jobEnvelope.ts:172-181

11. 和本组其它章的关系

本章只讲通用内核。registerAll() 实际接上的是十九条管线(platform/app/src/server/event-sourcing/pipelineRegistry.ts:456),与本文各章相关的几条分散在别的章:

管线想看什么去哪章
trace_processingtrace 汇总 fold + span 存储 map 的具体实现03-trace-processing.md
evaluation_processing评估状态机怎么用 fold + 订阅者串起来04-evaluation.md
experiment_run_processing离线实验跑的逐条结果与汇总04-evaluation.md
simulation_processing仿真运行的生命周期投影与执行 process manager06-simulations.md
suite_run_processing一批仿真(套件)的聚合状态06-simulations.md
其余(authz-grants、automations、billing-reporting、coding-agent-processing、gateway-spend-processing、langy-* 等)同一内核在别的域上的复用,无专章platform/app/src/server/event-sourcing/pipelines/

命令是怎么被生产出来的(一条 span 怎么被接住、变成 recordSpan),见 01-ingestion.md

AI 网关(05-ai-gateway.md)是独立的 Go 数据面,不走这套 TypeScript 内核;它只在下游通过网关消费处理管线把用量回灌进事件体系。


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

所有路径相对克隆根;符号名比行号抗漂移,优先用符号 grep。

主题文件符号
事件/投影核心类型platform/app/src/server/event-sourcing/domain/types.tsEventSchemaEvent
租户品牌类型platform/app/src/server/event-sourcing/domain/tenantId.tsTenantIdSchemacreateTenantId
聚合类型枚举platform/app/src/server/event-sourcing/domain/aggregateType.tsAggregateTypeSchema
命令与处理器接口platform/app/src/server/event-sourcing/commands/command.tsCommandCommandHandler
命令静态契约platform/app/src/server/event-sourcing/commands/commandHandlerClass.tsCommandHandlerClassStaticCommandHandlerClass
命令信封platform/app/src/server/event-sourcing/commands/commandEnvelope.tswithCommandEnvelopestripEnvelope
单事件命令糖platform/app/src/server/event-sourcing/commands/defineCommand.tsdefineCommand
命令派发扁平化platform/app/src/server/event-sourcing/mapCommands.tsmapCommands
链式装配 DSLplatform/app/src/server/event-sourcing/pipeline/staticBuilder.tsdefinePipelinewithSubscriberwithProcessManager
订阅者契约platform/app/src/server/event-sourcing/subscribers/subscriber.types.tsSubscriberDispatchDefinitionSubscriberDispatchOptionsshouldDispatch
process manager 契约platform/app/src/server/event-sourcing/pipeline/processManagerDefinition.tsProcessManagerDefinitionTriggerContextWakeHandler
中央运行时platform/app/src/server/event-sourcing/eventSourcing.tsEventSourcingregistercreateGlobalQueuelookupEntry
静态定义→运行时platform/app/src/server/event-sourcing/runtimePipeline.tsEventSourcingPipeline
组合根(十九条管线)platform/app/src/server/event-sourcing/pipelineRegistry.tsPipelineRegistryregisterAll
存事件 + 分发platform/app/src/server/event-sourcing/services/eventSourcingService.tsEventSourcingService.storeEventsregisterJob
队列门面与路由元数据platform/app/src/server/event-sourcing/services/queues/queueManager.tsQueueManager.createFacadebuildGroupKey
fold 定义与存储接口platform/app/src/server/event-sourcing/projections/foldProjection.types.tsFoldProjectionDefinitionFoldProjectionStore
fold 执行器platform/app/src/server/event-sourcing/projections/foldProjectionExecutor.tsFoldProjectionExecutor.executeexecuteBatchcanRefold
fold 基类platform/app/src/server/event-sourcing/projections/abstractFoldProjection.tsAbstractFoldProjectionFoldEventHandlers
map 定义与执行platform/app/src/server/event-sourcing/projections/mapProjection.types.tsmapProjectionExecutor.tsMapProjectionDefinitionAppendStoreMapProjectionExecutor
投影路由platform/app/src/server/event-sourcing/projections/projectionRouter.tsProjectionRouterregisterSubscriberdispatchSubscribersAfterStorecollapseByJobId
跨管线全局投影platform/app/src/server/event-sourcing/projections/projectionRegistry.tsProjectionRegistry
fold 缓存/仓储存储platform/app/src/server/event-sourcing/projections/redisCachedFoldStore.tsrepositoryFoldStore.tsRedisCachedFoldStoreRepositoryFoldStore
分组队列platform/app/src/server/event-sourcing/queues/groupQueue/groupQueue.tsGroupQueueProcessor
派发循环platform/app/src/server/event-sourcing/queues/groupQueue/dispatcher.tsGroupQueueDispatchernextWakeTimeoutSecGROUP_QUEUE_CONFIG
Lua 暂存脚本platform/app/src/server/event-sourcing/queues/groupQueue/scripts.tsSTAGE_LUADISPATCH_BATCH_LUACOMPLETE_LUAGroupStagingScripts
作业信封与 blobplatform/app/src/server/event-sourcing/queues/groupQueue/jobEnvelope.tsblobConstants.tsredisJobBlobStore.tstieredBlobStore.tsENVELOPE_PREFIX_V2GQ1_BLOB_BACKSTOP_TTL_SECONDSRedisJobBlobStoreTieredBlobStore
重试配置platform/app/src/server/event-sourcing/queues/shared.tsJOB_RETRY_CONFIGgetBackoffMs
失败错误分类platform/app/src/server/event-sourcing/queues/dispatchError.ts最小退避下限契约
内存队列回退platform/app/src/server/event-sourcing/queues/memory.tsEventSourcedQueueProcessorMemory
事件存储抽象与实现platform/app/src/server/event-sourcing/stores/abstractEventStore.tseventStoreClickHouse.tsAbstractEventStore.storeEventsEventStoreClickHouse
重建扫描窗口platform/app/src/server/event-sourcing/stores/rehydrationWindow.tsTIME_LOCAL_AGGREGATE_TYPESrehydrationLowerBoundMs
重放编排platform/app/src/server/event-sourcing/replay/replayService.tsReplayService.replaycleanup
重放折叠累加platform/app/src/server/event-sourcing/replay/replayExecutor.tsFoldAccumulatorMapAccumulator
重放标记platform/app/src/server/event-sourcing/replay/replayMarkers.tsreplayConstants.tsmarkPendingBatchisAtOrBeforeCutoffMarker
在线侧标记检查platform/app/src/server/event-sourcing/projections/replayMarkerCheck.tsRedisReplayMarkerCheckerReplayDeferralError
租户校验platform/app/src/server/event-sourcing/utils/event.utils.tsEventUtils.validateTenantIdcreateEvent
进程角色platform/app/src/server/app-layer/config.tsProcessRoleroleRunsWorkersroleSatisfiesRunIn
Reactor→Subscriber 退役决策dev/docs/adr/098-post-event-work-subscribers-and-process-managers.mdADR-098