Staff Engineer 指南
本指南面向需要在不破坏边界的前提下演进 Wow 的工程师。 它描述的是 main 分支中的真实架构,而不是愿景架构。 每个具体事实都链接到建立该事实的代码。 仓库无法证明的属性会明确标记为 未知。
执行摘要
Wow 是围绕聚合级命令执行组织的响应式 CQRS 与事件溯源框架。
wow-api 定义消息信封与公共契约。
wow-core 拥有命令调度、事件溯源、消息处理、投影、Saga、快照与运行时生命周期。
Spring 模块把这些机制适配到依赖注入和应用生命周期。 基础设施模块实现存储与传输契约。 WebFlux 与 OpenAPI 共享运行时路由目录;KSP 则在编译期生成它们依赖的元数据输入。 最重要的一致性边界是把 DomainEventStream 追加到 EventStore。 事件发布和下游处理发生在追加成功之后。 它们并不处在同一个分布式事务中。 框架明确保证聚合级顺序,但没有承诺全局顺序。 重试是选择性的,ack 是显式的,补偿表示重放而不是回滚。
WowRuntime 是已注册运行时组件的唯一生命周期所有者。
它只能启动一次,先关闭准入再排空,并使用有界关闭期限。 安全适配器传播身份相关 Header 和查询标签。 它们本身不能证明服务边界已经完成认证。 仓库具有本地、契约、集成、覆盖率与 JMH 测试层。 这些测试层不能建立生产 SLA 或通用吞吐上限。
唯一核心洞见
Wow 在一个聚合串行通道内把命令转换为不可变、带版本的事件流,先持久化该事件流,再把它扇出给职责独立的消费者。 其余机制都在保护或扩展这条主线。 命令信封携带聚合标识、所有权、租户、请求标识和期望版本。 聚合根决定产生哪些事件载荷。 状态聚合溯源这些事件。 事件存储完成持久追加。 领域事件总线和状态事件总线随后驱动投影、Saga 与快照。
WowRuntime 控制这些处理器何时能够接收工作。
下面的 Python-like 伪代码用于解释边界。 它特意把非事务性的扇出展示出来。
async def execute(command):
lane = lane_for(command.aggregate_id)
async with lane.serialized():
state = await snapshots.load(command.aggregate_id) or new_state()
async for stream in events.load_from(state.next_version):
state.source(stream)
emitted = await aggregate.decide(command, state)
state.source(emitted)
# 命令侧的持久一致性边界。
await events.append(emitted)
# 下游效果属于独立的响应式操作。
await domain_event_bus.send(emitted)
await state_event_bus.send_best_effort(emitted.with_state(state))真实实现会先溯源内存状态再追加,持久化失败时把命令聚合标为过期。 追加契约检测版本冲突和重复请求 ID。 领域事件发送位于聚合处理之后的过滤器中。 状态事件发送更晚,发送错误会被记录后恢复。 来源:命令信封、聚合执行、事件存储契约、追加后发布、尽力发送状态事件。
系统架构
架构按职责分层,而不是按部署拓扑分层。 应用可以在一个进程组合模块,也可以用分布式总线和存储连接多个进程。 仓库没有规定唯一的生产部署拓扑。
所有权表
| 领域 | 所有者 | 委托对象 | 边界证据 |
|---|---|---|---|
| 命令、事件、命名与建模公共契约 | wow-api | 不依赖 core 下层 | API 最小依赖 |
| 命令与事件运行时 | wow-core | 存储和总线接口 | core 依赖 |
| Spring 集成 | wow-spring | core 服务和 Spring 容器 | 模块依赖 |
| 可选 Spring Boot 组合 | wow-spring-boot-starter | feature variants | capabilities |
| HTTP 入口与路由物化 | wow-webflux | RouterSpecs 与处理器 | 模块依赖 |
| 路由契约与 OpenAPI 渲染 | wow-openapi | 元数据、贡献者、Schema 上下文 | RouterSpecs |
| Kafka 传输 | wow-kafka | Reactor Kafka | 模块边界 |
| MongoDB 持久化 | wow-mongo | Mongo 驱动 | 模块边界 |
| Redis 持久化 | wow-redis | Lettuce 与 Redis 脚本 | 模块边界 |
| Elasticsearch 持久化与查询 | wow-elasticsearch | Elasticsearch 客户端 | 模块边界 |
| CoSec 集成 | wow-cosec | WebFlux 请求上下文 | 适配器依赖 |
| 领域测试 DSL | test/wow-test | core 与 JUnit | 测试模块 |
| 后端契约 | test/wow-tck | 存储与调度接口 | TCK 模块 |
依赖方向
core 依赖接口,因此基础设施可以替换。 starter 通过可选 capability 组合组件,不把持久化或传输行为搬进 core。
承重契约
CommandMessage
CommandMessage 是业务命令载荷的运行时信封。
它携带 aggregateId、owner、space、command ID、request ID 与复制语义。 它还携带期望版本以及创建、允许创建、作废等控制标记。 这些字段是框架控制数据,不是领域状态。 来源:CommandMessage。
DomainEvent
DomainEvent 用聚合标识、序号、修订号和事件流位置包装业务事件载荷。
业务事件仍可保持为普通 Kotlin class 或 object。 示例 OrderCreated 只包含业务字段。 来源:DomainEvent、OrderCreated。
DomainEventStream
一个事件流代表一次命令执行产生的事件。 契约规定 command ID 与事件流是一对一关系。 具体事件流非空,并从第一个事件导出聚合与版本元数据。 来源:DomainEventStream。
EventStore
EventStore 拥有追加、请求查询、版本查询和事件流加载契约。
追加契约明确了版本冲突、重复聚合 ID 与重复请求 ID。 它没有定义跨事件发布或投影更新的事务。 来源:EventStore。
SnapshotStore
SnapshotStore 加载和保存状态检查点。
保存规则单调递增:低版本不能原子地覆盖高版本。 接口没有删除或保留期操作。 来源:SnapshotStore。
MessageBus
MessageBus 分离发送和接收。
接收器就绪是契约的一部分,生命周期归 WowRuntime 所有。 本地 sendIfSubscribed 在处理准入成功前保守返回 false。 来源:MessageBus。
RuntimeComponent
组件构造必须无副作用。
prepare、start、quiesce、优雅停止和强制停止是独立阶段。
契约刻意不使用 AutoCloseable,避免任意调用者拥有关闭权。 来源:RuntimeComponent。
领域模型与不变量
框架分离业务载荷、框架信封、命令行为与事件溯源状态。 这种分离支持事件溯源设计,但框架不会阻止命令处理器直接修改状态对象。 应在领域模型中落实该约定:保持 State Setter 私有,让 Command Handler 返回事件,再由 Sourcing Handler 应用事件。 来源:Command Root 构造、封装状态的 Cart 示例。
框架不变量
| Entity | Invariant | Enforced By | Consequence | Source |
|---|---|---|---|---|
CommandMessage | 命令按 named aggregate 与 aggregate ID 路由 | CommandMessage.aggregateId | 每个命令信封包含身份。 | 来源 |
CommandAggregate | 提供期望版本时必须匹配 | SimpleCommandAggregate | 过期写入者在领域调用前失败。 | 来源 |
CommandAggregate | 已提供的 owner 或 space 必须匹配已初始化聚合状态 | SimpleCommandAggregate | 非空值不匹配时在处理前拒绝;空值会跳过该比较。这是一致性检查,不是完整授权机制。 | 来源 |
CommandAggregate | 已删除聚合拒绝普通命令 | SimpleCommandAggregate | 删除状态成为访问保护。 | 来源 |
StateAggregate | 持久化前先溯源内存状态 | SimpleCommandAggregate | 处理期间状态反映已产生事件。 | 来源 |
CommandAggregate | 持久化失败使聚合实例过期 | SimpleCommandAggregate 错误钩子 | 失败实例不能继续充当权威状态。 | 来源 |
Snapshot | 快照版本不能后退 | SnapshotStore | 并发旧保存不能覆盖新状态。 | 来源 |
| 聚合分组 | 一个通道内顺序处理 | AggregateDispatcher | 顺序是分组级而非全局。 | 来源 |
示例订单聚合
示例 Order 展示了推荐分层。
Order 接收命令并返回事件。
OrderState 应用事件并通过 private setter 拥有可变状态。
CreateOrder 校验输入,OrderCreated 是不可变事件载荷。
支付根据金额产生一个或两个有序事件。 状态规则明确:仅 CREATED 可改地址,仅 PAID 可发货,仅 SHIPPED 可收货。 来源:Order 处理器、OrderState、CreateOrder。
aggregateVersion 可以为空;省略时不启用乐观版本前置条件。事件 revision 是语义化版本字符串,默认值为 0.0.1。
命令生命周期
入口
CommandHandlerFunction 提取 body、路径变量和 Header,再委托给 CommandHandler。
CommandHandler 构建命令消息,并选择 SSE 或普通等待行为。
来源:HTTP handler function、command handler。
Gateway
DefaultCommandGateway 在发送前校验消息并检查 request ID。
等待使用绝对超时,不会在每个阶段重新延长期限。 来源:校验与请求检查、等待期限。
调度
CommandDispatcher 从配置的 CommandBus 接收命令交换——来源可以是本地、分布式或两者合并后的本地优先视图——并解析聚合元数据。
它创建聚合专属调度器和 scheduler。 调度器按组处理,使一个聚合通道保持串行。 来源:调度器创建、聚合调度器。
加载与决策
仓储加载快照或创建新状态聚合。 随后从下一个期望版本重放事件流。 命令聚合在调用 handler 前检查版本与删除状态。对已初始化聚合,它还会在对应消息值非空时比较 owner 或 space。 来源:快照加尾部重放、前置条件。
持久化与发布
聚合先溯源产生的事件并追加事件流。 聚合处理完成后,领域事件过滤器才发送事件流。 状态事件过滤器排在领域事件过滤器之后。 它会记录发送失败并恢复。 来源:追加、领域事件发送、状态事件发送。
命令聚合状态
事件、投影、Saga 与快照生命周期
领域事件调度
领域调度器拥有领域事件和状态事件两个子调度器。 Function kind 选择对应子调度器。 同一事件流内通过 concatMap 顺序处理事件。 普通事件处理器的返回值在完成后被丢弃。 来源:组合调度器、事件流内处理、返回值丢弃。
投影
ProjectionDispatcher 同时订阅领域事件与状态事件总线。
它使用事件函数过滤器,因此投影 Publisher 表示完成,不表示新领域事件。 来源:ProjectionDispatcher、ProjectionFunctionFilter。
无状态 Saga
无状态 Saga 是把 handler 结果转换为命令的特殊路径。 新 request ID 从源事件 ID 和结果序号派生。 tenant、space 与 upstream Header 传播到新命令。 这是命令编舞。 它不是分布式事务,也不会自动撤销之前的副作用。 来源:StatelessSagaFunction。
快照
快照是从状态事件派生的检查点。
Starter 的默认快照策略是 ALL。仅当选择 VERSION_OFFSET 时,VersionOffsetSnapshotStrategy 才使用默认的五个版本偏移。
它比较已保存版本,并在需要时保存更新的 SimpleSnapshot。 快照保存不属于事件存储追加事务。 来源:Starter 快照默认值、策略契约、版本偏移策略、快照过滤器。
生命周期对比
| 产物 | 创建者 | 持久边界 | 消费者 | 失败含义 |
|---|---|---|---|---|
| 命令消息 | Gateway 或总线客户端 | 取决于总线 | 命令调度器 | 校验或传输失败 |
| 领域事件流 | 命令聚合 | EventStore.append | 领域事件总线 | 版本、重复或存储失败 |
| 领域事件投递 | 追加后过滤器 | 取决于总线 | 事件处理器、投影、Saga | 重试、处理策略、随后 ack |
| 状态事件 | 领域事件之后的过滤器 | 取决于总线 | 投影和快照调度器 | 即时发送错误被记录并恢复 |
| 快照 | 快照策略 | SnapshotStore.save | 聚合仓储 | 派生检查点可能落后于事件存储 |
| Saga 命令 | 无状态 Saga 结果映射器 | 命令总线及后续事件追加 | 另一个聚合 | 不会隐式回滚源事件 |
来源:发送过滤器、重试过滤器、ack 语义、Saga 映射。
运行时生命周期
WowRuntime 是一次性的生命周期协调器。
状态包括 NEW、STARTING、RUNNING、STOPPING、FORCE_STOPPING 和 STOPPED。 所有组件 prepare 完成后才能进入 start 阶段。 组件意外失败会关闭准入并启动关闭流程。 优雅关闭只有一个所有者和一个全局期限。 顺序是关闭全局准入、quiesce 组件、排空工作、反向停止组件。 超时或优雅停止失败会升级为强制停止。 启动清理只是生命周期回滚。 它不是领域事件或外部副作用回滚。 来源:状态与拓扑、启动与启动清理、关闭所有权、关闭顺序。
组件顺序
组件注册到有序且身份去重的 slot。 prepare 与 start 按注册顺序执行。 优雅停止与强制停止按反向顺序执行。 系统保留第一个失败,同时继续后续清理。 来源:RuntimeComponentGroup。
Spring 桥接
Spring 生命周期桥接保证入口看到已经就绪的 Wow 运行时。 它在入口排空后停止,运行时意外终止时关闭应用上下文。 默认关闭超时为 60 秒,quiet period 为 1 秒。 来源:WowRuntimeLifecycle、WowProperties。
存储架构
存储通过注册表和路由装饰器按聚合选择。 路由器拥有生命周期,并把每个操作委托给选中的后端。 聚合专属映射优先于默认存储。 来源:RoutingEventStore、事件注册表、快照路由、快照注册表。
后端对比
| 后端 | EventStore | SnapshotStore | 重要边界 | 来源 |
|---|---|---|---|---|
| In-memory | 是 | 是 | 用于开发和测试;持久性仅限进程 | InMemoryEventStore、InMemorySnapshotStore |
| MongoDB | 是 | 是 | 有序加载;直接写或可选批量追加 | MongoEventStore、MongoSnapshotStore |
| Redis | 是 | 是 | Lua 追加检查冲突;不支持按事件时间加载 | RedisEventStore、RedisSnapshotStore |
| Elasticsearch | 是 | 是 | refresh 与可选批处理影响可见性和延迟 | ElasticsearchEventStore、ElasticsearchSnapshotStore |
批处理
MongoDB 和 Elasticsearch 批处理默认选择性启用。 默认关闭,因为不满批次会增加最多 maxDelay 的延迟。 默认参数包括批次 128、pending 4096、lane 1、延迟 1ms。 这些是配置默认值,不是所有工作负载的最优测量值。 来源:Mongo 事件选项、Elasticsearch 事件选项。 pending 队列有界,耗尽时可用类型化过载错误拒绝准入。 这是显式背压,而不是无界静默缓存。 来源:Mongo 批量追加器。
消息架构
聚合级顺序
AggregateDispatcher 把消息映射为 group key。
每个组使用 publishOn 加 concatMap 顺序处理。 不同组可以并行。 默认 lane 数为 64 * available processors,可用系统属性覆盖。 这不是全局顺序保证。 来源:分组处理、并行度。
Local-first 行为
Local-first 同时准备本地投递副本与分布式副本。 只有本地准入成功后,分布式副本才标为已在本地处理。 本地投递出错时,分布式路径仍可用。 被过滤的分布式副本会 ack。 这是准入感知优化。 它不能证明集群级 exactly-once。 来源:LocalFirstMessageBus。
Kafka
Kafka 发送在 sender result 返回时完成。 接收使用 consumer group,重试 receive stream,并顺序解码记录。 Kafka key 是 aggregate ID 字符串。 topic converter 提供聚合与函数路由上下文。 来源:发送与接收、订阅、key 与序列化。 默认 Kafka receiver policy 使用 prefetch 1、maximum deferred ack 1、重试 3 次、延迟 10 秒。 这些是配置默认值,不是吞吐保证。 来源:KafkaReceiverPolicy。
Ack 语义
finallyAck 在成功后 ack。
发生错误时,它先 ack 再重新抛出错误。 因此 handler 最终失败本身不意味着 broker 会重新投递。 依赖重放前必须理解重试与补偿策略。 来源:ExchangeAck。
失败处理
默认重试过滤器最多重试三次,backoff 为两秒。 只有标记为 recoverable 的异常才会重试。 它排在聚合、事件函数与快照处理过滤器之前。 默认事件处理错误策略在过滤器策略结束后记录并恢复。 来源:RetryableFilter、事件自动配置。 补偿会重新加载已持久化事件、添加 compensation target 并重发。 状态事件补偿在重发前通过事件溯源重建状态。 两条路径都不会撤销原始事件存储追加。 来源:领域事件补偿、状态事件补偿。
失败模式表
| 失败 | 即时行为 | 持久事实 | Staff Engineer 动作 |
|---|---|---|---|
| 期望版本不匹配 | handler 调用前拒绝 | 现有事件流 | 按乐观并发冲突处理 |
| 重复 request ID | EventStore 契约拒绝重复 | 第一次接受的事件流 | 客户端重试必须保留 request ID |
| EventStore 追加失败 | 聚合实例过期 | 是否提交由后端结果决定 | 重用状态前重新加载并核对后端 |
| 领域事件发送失败 | 命令事件流可能已经持久化 | EventStore 仍是事实来源 | 按异常分类使用重试或显式补偿 |
| 状态事件发送失败 | 记录错误并恢复 | 事件流仍持久 | 监控延迟,必要时执行状态事件补偿 |
| 投影 handler 失败 | 选择性重试,随后执行 ack/error 策略 | 投影可能落后 | handler 幂等并定义重放 runbook |
| Saga 命令失败 | 源事件已提交 | 无自动回滚 | 显式建模业务补偿命令 |
| 运行时启动失败 | 清理已启动组件 | 不意味着领域回滚 | 同时检查首个失败与清理失败 |
| 优雅关闭超时 | 升级为强制停止 | 在途结果可能未知 | 用 request ID 与事件存储对账 |
元数据、生成代码、路由与 OpenAPI
这些是相关但不同的流水线,不是一个生成步骤。
编译期 KSP 元数据
MetadataSymbolProcessor 扫描 bounded context 与 aggregate root。
它合并结果并把元数据资源写成 JSON。
AggregatesMetadataResolver 另行生成调用 aggregateMetadata<Command, State>() 的 Kotlin 访问器,该函数会调用运行时聚合元数据 parser。
这是两条不同的运行时输入,不是一条发现链:MetadataSearcher 加载 JSON 资源,生成访问器则通过 aggregateMetadata() 调用 AggregateMetadataParser。 来源:元数据资源生成、聚合访问器生成、资源搜索、运行时 parser。
运行时路由目录
RouterSpecs 对 route contributor 排序。
它读取运行时 MetadataSearcher、过滤禁用的聚合路由并构建经校验的 RouteCatalog。 目录会拒绝重复 route key 和路径变量不匹配。 来源:路由收集、目录校验。
运行时 WebFlux 物化
RouterFunctionBuilder 遍历路由目录。
它把每个契约物化为 predicate 与 handler function。 Spring Boot 用 RouterSpecs 和 handler registrar 创建该 router。 来源:RouterFunctionBuilder、WebFlux 自动配置。
运行时 OpenAPI 渲染
同一目录被渲染为 OpenAPI 3.1 path 与 component。 Springdoc customizer 把生成目录合并进应用 OpenAPI 对象。 来源:OpenAPI 渲染、OpenAPI 自动配置。
反射边界
不要把 Wow 描述成零反射框架。 core 把 Kotlin reflection 声明为 API 依赖。 元数据 parser 文档明确包含反射分析。 测试 DSL 也反射泛型参数。 KSP 减少部分发现与注册样板,但不能证明运行时零反射。 来源:core reflection 依赖、元数据 parser 契约、AggregateSpec 反射。
安全与信任边界
请求上下文
WebFlux 从 path 和 Header 提取 tenant、owner、space、aggregate ID 与 local-first 提示。 提取不等于认证。 部署必须决定哪些 Header 可来自不受信客户端,哪些必须由可信边缘覆盖。 来源:AggregateRequest。
CoSec 适配器
CoSec extractor 把 request ID 和 space ID Header 写入命令 builder。 其他 CoSec 适配器传播 app 与 device ID。 该模块仅依赖 WebFlux,本身没有建立 authenticator。 来源:builder extractor、消息传播、模块依赖。
聚合授权前置条件
对已初始化聚合,命令处理仅在对应消息值非空时比较 owner 或 space 与已加载聚合状态。 读侧 owner precondition 可以拒绝 owner aggregate 访问。 路由元数据控制 owner path 是 NEVER、ALWAYS 或由 AGGREGATE_ID 决定。 来源:命令检查、owner precondition、路由所有权。
这些条件检查会在上下文已提供时保护聚合一致性;它们不负责认证调用方,也不能替代端点授权。
查询 ABAC
AbacQueryFilter 把 principal tag 转换为查询条件。
principal tag 解析是抽象方法,必须由集成实现。 空 tag 集解析为 Condition.all()。 因此仅存在该过滤器不能证明查询已认证或受限。 来源:AbacQueryFilter。
安全检查清单
- 信任 Wow 身份 Header 前先终止外部认证。
- 在边缘剥离客户端提供的内部 Header。
- 把 tenant、owner、space 绑定到已认证 principal。
- 为 ABAC 提供具体 principal-tag resolver。
- 显式测试空 tag 行为。
- 把 local-first Header 当作内部路由提示。
- 验证补偿端点具有运维授权。
- 验证 metadata 与 BI script 端点符合暴露策略。
- 对外发布前审计生成的 OpenAPI。
- 存储凭证和签名材料不得进入仓库。
代码建立了以上提取与过滤点。 具体生产身份提供方、边缘策略与 secret store 在仓库中 未知。
性能模型
结构性热路径
写路径包括请求解码、校验、request ID 检查、总线准入、lane 调度、快照加载、尾部重放、领域调用、事件序列化、存储追加、发布与可选等待协调。 主导成本取决于工作负载和部署。 仓库不能证明唯一的通用瓶颈。
显式边界与调节项
| 调节项 | 代码默认值 | 限制内容 | 不能证明 |
|---|---|---|---|
| Dispatcher lanes | 64 * processors | 进程内分组并行度 | 最优 CPU 或存储并发 |
| Kafka prefetch | 1 | Receiver demand | 端到端吞吐 |
| Kafka deferred ack | 1 | 未完成 deferred ack | 投递保证 |
| Kafka retry | 3,延迟 10s | receive-stream 重试 | 最终 ack 后 handler 重放 |
| Batch max size | 128 | 一次可选存储批次 | 工作负载最优批次 |
| Batch max pending | 4096 | pending 队列容量 | 饱和时安全内存或延迟 |
| Batch lanes | 1 | coordinator lane 数 | 通用最优顺序策略 |
| Batch max delay | 1ms | 不满批次等待 | 端到端延迟 |
| Runtime timeout | 60s | 默认关闭期限 | 业务操作期限 |
来源:消息并行度、Kafka receiver policy、Mongo batch 选项、运行时默认值。
Benchmark 证据
benchmark 模块包含组件、端到端、WebFlux、MongoDB、Redis 与 Elasticsearch fixture。 它使用 JMH,并依赖 example、test、mock 与基础设施模块。 来源:benchmark 依赖、JMH 版本。 模拟 I/O benchmark 研究 I/O 延迟与 scheduler handoff。 批量 E2E benchmark 按命令归一化结果。 并发 benchmark 明确说重复 key 顺序由功能测试覆盖,而非吞吐 benchmark。 来源:模拟 I/O benchmark、批量 E2E benchmark、coordinator benchmark 范围。
README 压测样例
README 报告了一次示例应用的两分钟压力测试。 它列出特定操作与等待计划的平均和峰值 TPS。 这些数字是链接部署条件下的历史样例。 它们不是 SLA、容量计划或组件性能上限。 来源:README 样例。
性能决策规则
使用可复现工作负载。 固定代码修订和环境。 同时测量存储、broker、CPU、分配与 scheduler 行为。 分离组件筛选与端到端确认。 改变并发时重新验证顺序与过载行为。 不要根据一次 quick benchmark 修改默认值。 没有测量就不要把 EventStore 结果迁移到 SnapshotStore。 在部署专属实验给出结果前,生产容量与尾延迟目标均为 未知。
测试策略
测试层
| 测试层 | 目的 | 证据 |
|---|---|---|
| Domain spec | Given/when/expect 行为 | AggregateSpec |
| Saga spec | 隔离验证产生的命令 | SagaSpec |
| EventStore TCK | 追加、加载、冲突、重复、并发 | EventStoreSpec |
| SnapshotStore TCK | 加载、单调保存、并发 | SnapshotStoreSpec |
| 后端契约实现 | 在真实适配器运行 TCK | Mongo、Redis、Elasticsearch |
| Integration CI | 服务与聚合集成任务 | workflow |
| Static analysis | Detekt | workflow |
| Coverage | 库模块启用 Jacoco;要求阈值时由具体模块配置 | 根 Jacoco 配置、示例 80% 规则 |
| Benchmark | JMH 回归与诊断 | benchmark 模块 |
Domain test 风格
DSL 通过 JUnit dynamic test 暴露 Given、When、Expect 阶段。
AggregateSpec 的泛型命令聚合类型发现使用反射。
示例领域模块可执行 80% Jacoco 下限。 来源:AggregateSpec factory、示例覆盖率。
变更到测试映射
| 变更 | 最小聚焦验证 | 更宽门禁 |
|---|---|---|
| 命令校验或 handler | 成功与拒绝的 Aggregate spec | 领域模块 check |
| 事件溯源规则 | 完整历史与快照尾部重放 | 存储 TCK 与集成测试 |
| EventStore 适配器 | 冲突、重复请求、顺序、并发 | 适配器 check 与集成 workflow |
| SnapshotStore 适配器 | 并发单调保存 | Snapshot TCK 与适配器 check |
| 调度并发 | 同 key 顺序、跨 key 并行、quiesce | core test 与 benchmark 诊断 |
| 运行时生命周期 | prepare barrier、反向清理、超时、取消 | :wow-core:test |
| Route contributor | 目录校验与路由快照 | OpenAPI 与 WebFlux test |
| 元数据 KSP | 生成资源与访问器 golden output | compiler check |
| 安全过滤器 | 已认证、未认证、空 tag、伪造 Header | WebFlux 集成测试 |
| 性能默认值 | 多 fork 组件与 E2E 对比 | 部署代表性压测 |
绿灯测试只证明 fixture 与 assertion 覆盖的内容。 除非条件进入测试,否则不能证明真实 provider 取消、生产授权、迁移安全或部署 SLA。
架构决策
仓库中没有可引用的 ADR 记录这些机制的历史替代方案或原始动机。 因此下表有意将替代方案写为“未声明”;“理由”是对当前代码行为的架构解释,不是历史设计意图的证据。
| 决策 | 考虑过的替代方案 | 理由 | 来源 |
|---|---|---|---|
| 分离业务载荷与框架信封 | 未声明 | 当前信封把路由、身份、所有权和版本控制放在业务载荷之外;这是对现有契约的解释。 | CommandMessage |
| 先持久化再发布领域事件 | 未声明 | 当前过滤器顺序保证 EventStore 追加先完成,下游因此允许延迟或重放。 | 过滤器顺序 |
| 把快照视为派生检查点 | 未声明 | 当前策略在事件处理后保存已溯源状态,不替代事件历史。 | 快照策略 |
| 按聚合派生分组串行处理 | 未声明 | 当前分组 concatMap 保持同组顺序,同时允许不同分组独立推进。 | AggregateDispatcher |
| 使用准入感知 local-first 投递 | 未声明 | 实现会尝试本地投递并保留带标记的分布式副本,以更复杂的 copy 与 ack 语义换取减少 broker 往返。 | LocalFirstMessageBus |
| 由唯一运行时独占生命周期所有权 | 未声明 | 当前契约分离 prepare、start、quiesce、graceful stop 和 force stop,使就绪与清理由一个协调者拥有。 | RuntimeComponent |
| 按聚合路由存储 | 未声明 | 当前注册表允许聚合专属存储,同时保留默认后端。 | 存储注册表 |
| 通过 starter capability 组合适配器 | 未声明 | Feature variant 使基础设施模块可选,同时让 variant resolution 成为发布兼容面。 | starter capabilities |
| 共享一个经校验的路由目录 | 未声明 | 当前目录同时供路由与 OpenAPI 物化使用,降低两类输出之间的契约漂移。 | RouterSpecs |
| 把补偿建模为显式重放操作 | 未声明 | 当前 compensator 重新加载持久事件并向目标重发,不撤销原始追加。 | DomainEventCompensator |
依赖理由
版本目录固定 Kotlin 2.4.10、KSP 2.3.11、Spring Boot 4.1.1、JUnit 6.1.3、Testcontainers 2.0.5 与 JMH 1.37。 来源:version catalog。
引用证据中没有 ADR 或迁移记录说明这些依赖替换过什么。 因此 替代了什么 一列统一写“未声明”,不虚构历史。
| 依赖 | 用途 | 替代了什么 | 来源 |
|---|---|---|---|
| Kotlin 与 KSP | Kotlin 实现框架,KSP 在编译期生成元数据资源与类型化访问器。 | 未声明 | 版本目录、编译器依赖 |
| Spring Boot | 提供自动配置、生命周期集成、WebFlux 组合与 feature variant。 | 未声明 | starter features |
| Reactor | 提供命令、事件、重试、顺序与排空路径使用的非阻塞 Publisher 模型。 | 未声明 | core 依赖 |
| Jackson | 序列化命令、事件、状态与元数据表示。 | 未声明 | core 依赖、消息序列化器 |
| Reactor Kafka | 实现分布式 Kafka 消息总线适配器。 | 未声明 | Kafka 模块 |
| MongoDB reactive driver | 实现 MongoDB 事件、快照与查询持久化。 | 未声明 | MongoDB 模块 |
| Spring Data Redis 与 Lettuce | 实现 Redis 事件与快照持久化及 Redis 传输集成。 | 未声明 | Redis 模块 |
| Spring Data Elasticsearch | 实现 Elasticsearch 事件、快照与查询适配器。 | 未声明 | Elasticsearch 模块 |
| Swagger/OpenAPI libraries | 把运行时路由目录建模并渲染为 OpenAPI 契约。 | 未声明 | OpenAPI 模块 |
| JUnit 与 Testcontainers | 提供动态领域测试、后端 TCK fixture 与外部服务集成测试。 | 未声明 | TCK 依赖 |
已知技术债
仓库没有在可引用的 ADR 或 Issue 中把以下缺口标记为技术债。 定性风险等级是根据当前影响给出的评审优先级,不是维护者承诺。
| 问题 | 风险等级 | 受影响文件 | 来源 |
|---|---|---|---|
| Redis 无法实现公共的按事件时间加载能力,因此该后端不支持时间范围重放。 | 中 | EventStore.kt、RedisEventStore.kt | 契约、Redis 实现 |
| 公共 SnapshotStore 契约缺少删除和保留能力,生命周期政策只能由后端运维或额外应用契约补充。 | 中 | SnapshotStore.kt、选定快照后端与部署政策 | SnapshotStore |
显式框架边界与有意约束
以下行为是代码确认的边界。 没有 ADR、Issue 或维护者决策声明改造意图时,不应直接称其为技术债。
| 约束 | 工程影响 | 来源 |
|---|---|---|
| 状态事件发送错误会在即时过滤器边界记录并恢复。 | 快照和状态事件消费者可能落后;具体总线持久性和重放政策必须补齐运营闭环。 | SendStateEventFilter |
| 认证和 principal-tag 解析由集成负责。 | 只有 Header 提取与 ABAC hook 不能证明访问已经认证或受限。 | CoSec 提取、ABAC 空 tag |
| 普通事件处理器返回值没有发布语义。 | 事件结果需要转成命令时,应使用 stateless saga 映射。 | 事件函数过滤器、Saga 映射器 |
WowRuntime 与 Spring bridge 都是一次性的。 | 嵌入代码必须替换运行时,而不是重启已停止实例。 | 一次性启动、Spring 生命周期状态 |
| KSP 不会消除聚合运行时反射。 | AOT、启动时间或反射削减工作必须测量真实 parser 和 invocation 路径。 | 生成访问器、运行时 parser |
需要部署证据的未知项
- 生产认证提供方未知。
- 可信代理与 Header 清洗策略未知。
- 每个聚合的生产 EventStore 与 SnapshotStore 选择未知。
- broker 复制与保留策略未知。
- 灾难恢复 RPO 与 RTO 未知。
- 投影重放 runbook 未知。
- 补偿端点授权策略未知。
- 可接受状态事件延迟未知。
- 生产命令延迟 SLO 未知。
- 任一部署的安全最大并发未知。
- 各存储后端的容量上限未知。
- 引用契约没有建立事件载荷 Schema 迁移策略。
- 事件流与快照保留策略未知。
- 强制停止部分完成后的运维响应未知。
- 客户端网络重试是否保持 request ID 未知。
这些不一定是框架缺陷。 它们是生产设计必须补充的输入。
Staff Engineer 变更协议
设计前
- 明确拥有行为的 aggregate、bounded context 与模块。
- 明确持久事实是 EventStore、snapshot、projection 还是外部系统。
- 追踪路由使用的信封字段和元数据。
- 找到拥有准入和关闭权的 RuntimeComponent。
- 说明变更影响单聚合通道还是跨聚合协调。
- 列出精确的 retry、ack 与 compensation 行为。
- 区分可信与不可信 Header。
- 判断 KSP 产物、运行时路由目录是否变化。
- 定义持久事件和公共路由的向后兼容性。
- 行为变化时优先先写失败模式测试。
实现期间
- 公共契约留在
wow-api。 - 运行时行为留在
wow-core。 - Spring wiring 留在
wow-spring*。 - 传输与存储细节留在适配器模块。
- 保持 Reactor 路径非阻塞。
- 保持聚合级顺序。
- 不要意外扩大 ack 语义。
- 不要在没有重放路径时隐藏发送失败。
- 不要把手改生成文件作为主要修复。
- 让路由和 OpenAPI 继续由同一个目录驱动。
合并前
- 先运行最窄模块测试。
- 运行相关存储或调度器 TCK。
- 对修改的 Kotlin 运行 static analysis。
- 路由变化时渲染并比较 OpenAPI。
- 注解变化时检查生成元数据。
- 生命周期变化时测试启动、优雅停止、强制停止。
- 并发变化时测试同 key 顺序和跨 key 并行。
- 分别测试 recoverable 与 unrecoverable 错误。
- 验证受影响处理器的补偿幂等性。
- 记录仍然存在的部署假设。
推荐阅读顺序
- 从
CommandMessage理解控制信封。 - 阅读
DomainEvent,分离载荷与元数据。 - 阅读
DomainEventStream,理解命令到事件流关系。 - 阅读
SimpleCommandAggregate,定位一致性边界。 - 阅读
EventStore,理解持久化契约。 - 阅读
EventSourcingStateAggregateRepository,理解快照加尾部重放。 - 阅读两个发布过滤器,理解追加后边界。
- 阅读
AggregateDispatcher,理解顺序。 - 阅读
LocalFirstMessageBus,理解准入感知投递。 - 修改失败策略前阅读
ExchangeAck。 - 阅读
RetryableFilter,理解重试分类。 - 阅读
DomainEventCompensator,理解重放语义。 - 阅读
StatelessSagaFunction,理解跨聚合编舞。 - 阅读
VersionOffsetSnapshotStrategy,理解快照时机。 - 修改生命周期前阅读
RuntimeComponent。 - 阅读
WowRuntime,理解启动所有权。 - 继续阅读关闭所有权。
- 阅读
RuntimeComponentGroup,理解顺序与清理。 - 阅读
MetadataSymbolProcessor,理解编译期元数据。 - 阅读
RouterSpecs,理解运行时路由与 OpenAPI 组装。 - 阅读
RouterFunctionBuilder,理解 HTTP 物化。 - 框架边界明确后再阅读订单示例。
评审启发式规则
拒绝把 snapshot 当成事实来源的变更。 没有显式新一致性模型时,拒绝在 EventStore 追加前发布。 拒绝在命令、事件、投影、Saga 或存储响应式路径引入阻塞 I/O。 没有 broker、ack、handler 幂等与重放证据时,拒绝 exactly-once 结论。 代码只提供 retry 或 compensation replay 时,拒绝自动回滚结论。 仍存在反射依赖和 parser 时,拒绝零反射结论。 拒绝仅根据 README 样例或一次 quick JMH 修改性能默认值。 拒绝仅根据 Header 提取作出授权结论。 持久事件 Schema 变化必须提供显式迁移路径。 新增 RuntimeComponent 必须有生命周期测试。 新增队列或 batch coordinator 必须有过载测试。 路由元数据变化必须检查路由目录与 OpenAPI。
术语表
Aggregate lane:从 aggregate ID 派生的串行处理组。 Command envelope:CommandMessage 加路由、身份、所有权与版本控制数据。 Domain event payload:描述事实的应用级不可变对象。 Domain event stream:一次命令执行产生的非空有序事件集合。 Event sourcing:在可选快照之后应用持久事件流重建状态。 State event:用已溯源聚合状态装饰的领域事件流。 Snapshot:用于减少重放工作的派生状态检查点。 Projection:更新读模型的事件消费者。 Stateless saga:其结果被转换为新命令的事件函数。 Compensation:针对目标函数显式重放持久领域事件或重建状态事件。 Local-first:保留分布式投递路径,同时尝试本地准入。 Quiesce:停止接收新工作,但允许已准入工作排空。 Force stop:优雅完成不再可行后的尽力关闭。 Route catalog:由 WebFlux 路由与 OpenAPI 物化共享的已校验运行时契约。 Generated metadata:KSP 生成由 MetadataSearcher 加载的 JSON 资源,以及调用运行时聚合元数据 parser 的独立访问器。
最终心智模型
从一个聚合和一个命令开始。 沿命令信封进入一个串行通道。 从快照加事件尾部重建状态。 让聚合产生事件,而不是直接修改存储。 把一个事件流追加为命令的持久结果。 把之后的总线、投影、Saga 与快照效果视为职责独立的异步边界。 让 WowRuntime 决定这些所有者何时可以接收和完成工作。 对分类为临时的失败使用 retry。 需要下游重新处理时,对持久事件执行显式 compensation replay。 对安全、投递与性能保证只采用证据,不采用标签。