- architecture.md:系统定位与双跑策略、技术栈、总体拓扑(两条单线程管道)、 模块职责、关键决策 D1–D8、数据边界(六辅助表 vs legacy 旧表)、两阶段权威 与就绪度(诚实口径:生产默认不可服务,前置 U05/U07/U09/U13/U15)、 部署与安全姿态、可观测性。 - design.md:状态机与错误分类(含重放白名单)、六表数据模型(含 U18 已知缺口)、 流程 1–4 与泵作业语义、失败/重试统一设计(U08)、不变量 I1–I5 落点、 参数表、测试策略、可观测性、缺口清单对照 ACM2-10 U01–U30。 - README:新增文档导航;关联节按 ACM2-10 T16 修正改为 ACM2 现行入口 (ACM2-3 架构权威 / ACM2-4 脚手架 / ACM2-10 评审),ACMA 仅归档。
13 KiB
msgexchange-v2 设计文档
本文对应仓库当前实现,给出模块级设计语义与依据;架构总览见 architecture.md,权威架构为 Plane ACM2-3,实施计划与逐项验收为 ACM2-10(U01–U30)。文中标注「TODO/未实装」的条目均为已知开放项,不属文档遗漏。
1. 领域模型
1.1 状态机与错误分类
ProcStatus(PROC_STATE.STATE,消息处理侧):
PENDING ──处理成功──▶ SUCCEEDED(终态,回填 CMINMSGS)
│ ──同 identity 已绑定──▶ SKIPPED(终态,lastError=duplicate-of:<id>)
└──失败──▶ FAILED(非终态,attempts+1 + nextAttemptAt 退避)
│ attempts ≥ maxAttempts 或 队头滞留超 head-deadline
▼
DEAD(终态/DLQ,ERROR_CLASS=EXHAUSTED 规范化)
EventStatus(MSG_EVENT.STATE,投递侧):
PENDING ──▶ SENT;失败退避回 PENDING;attempts 耗尽整批/单条 → DEAD(DLQ)
ErrorClass(两侧共用):
MALFORMED 报文非法 → 直接 DEAD,永不重放
CODEC_ERROR 可随 codec 修复 → FAILED 可重放(白名单内)
UNSUPPORTED 能力未实装 → FAILED 可重放(白名单内)
INFRA 基础设施抖动 → FAILED 可重放(白名单内)
EXHAUSTED 重试耗尽(终态规范化类)→ 人工复核后可重放(白名单内)
- 「未实装 ≠ 非法」是通用规则(D4):无 handler、快照 staging 未实装都写
FAILED(UNSUPPORTED),绝不写终态——阶段 2 前接入流量不会把报文变砖。 - 显式重放入口
ReplayService:仅白名单{CODEC_ERROR, UNSUPPORTED, INFRA, EXHAUSTED}可从 FAILED/DEAD 回 PENDING (ATTEMPTS=0、NEXT_ATTEMPT_AT=NULL,errorClass/lastError 保留审计); 含 MALFORMED 的请求对该类静默忽略。运维接口(controller/runbook)属 U11 遗留。
1.2 报文模型(sealed 分派)
MetaFields(sndr, type, styp, seqn, dttm)——实名沿用 legacy META.java。MsgKindsealed:Schd(RESP|DNLD|ADFT)+Flop(29 类 STYP);typeTag产出SCHD-XXX/FLOP-xxx,与HandlerRegistry.keyOf同源(查表键=日志类型,禁止分叉)。- 解码失败二分:
DecodeResult.Err(MALFORMED)→ DEAD;Err(CODEC_ERROR)→ FAILED 退避 (T06/U11)。 - Handler 为纯函数:
decide(flightView, msg) → Decision(flightChanges / msgNotifies / schdPush / outboundIntents / refUpserts),不触碰 Redis/Kafka——副作用全部由泵边界执行。 - Handler 实装:0/32(骨架),翻译属阶段 2/3,逐条对照 ACM2-4 行为基线与 KEEP/FIX 矩阵。
2. 数据模型(六辅助表)
db/migration/V2.0.0__aux_tables.sql(legacy 旧表不在本仓库声明,见 architecture.md §6):
| 表 | 角色 | 关键列/约束 |
|---|---|---|
| PROC_STATE | 消息处理伴生状态(不动 CMINMSGS 旧列) | UK_PROC_IDENTITY(IDENTITY_KEY) 唯一约束=I3 依据;IDX_PROC_HEAD(STATE, CMINMSGS_ID)=队头查询 |
| MSG_EVENT | 统一投递 outbox | EVENT_ID 自增=全序;IDX_EVT_HEAD(TARGET, STATE, EVENT_ID)=每 target 队头 |
| REF_DATA | 21 类参考数据 + SCHD_GEN | VERSION 列支撑流程 4 CAS;SOURCE 区分 ADMINAPI/AODB/PIPELINE |
| REQ_TRACK | 15 类请求状态机 | REGISTERED/SENT/WAITING/DONE/EXPIRED |
| PUMP_JOB | 泵作业队列 | kind:ARCHIVE/HISTORY_SWEEP/PROJECTION_REBUILD |
| FLIGHT_STATE | 阶段 B 权威 | replaceDay 单事务删差集+写新代+版本提升 |
已知 DDL 缺口(U18,未修):REQ_TRACK.COUTMSGS_ID 应 INT→BIGINT;FLIGHT_STATE
主键应含 FDAY;时间列 TIMESTAMP(秒级+会话时区)应 DATETIME(6)/显式 UTC,否则退避/毒丸
判定存在系统性偏移风险。
3. 核心流程设计
3.1 流程 1:收报(InboxService.accept)
事务 1 = insertRaw(CMINMSGS 原文)+ procState.insert(伴生 PENDING 行);
不解析报文、接收层无唯一约束(I3)。响应 = 记录 ID(「已持久化」语义,与现役逐字对拍
后固化,U16)。wakePump() 目前为 TODO 空操作——泵 1s 轮询兜底,唤醒仅为加速。
事务边界随 U05(@Transactional + allopen)补齐。
3.2 流程 2:主泵 tick(Pump.tick)
每 tick 按序判定:
headQueued()取作业、headUnfinished()取最小未完成 CMINMSGS_ID(含 FAILED, 消息间严格保序:队头退避未到期即 sleep 至到期点,后方消息永不越队)。jobBefore(job, head):head 为空或队头 FAILED 时作业先行——当前为入队时间近似, 与「统一 FIFO」注释存在已知偏离(U15 未实装);统一序号列定案后消除。- 队头 FAILED 且退避未到期:
poisoned()判定(attempts≥maxAttempts 或滞留超 head-deadline 10m)→ DEAD(EXHAUSTED) 毒丸升级(Pump 侧;投递侧同语义属 U13,未实装); 否则 sleep 至 nextAttemptAt。 - 正常队头 →
MessageProcessor.processOne: 入口守卫(FAILED 且已 exhausted → DEAD)→rawOf缺失 → DEAD(MALFORMED) → decode(MALFORMED→DEAD / CODEC_ERROR→FAILED)→ identity 首绑 (tryBindIdentity失败 → SKIPPED,I3)→ Schd DNLD →SnapshotFlow→ 其余 →Handler.decide(redis.hgetAllFlightInfo(), msg)→ 阶段 A:Redis 先写(I2 happens-before,TODO redisApply)→ 事务 2: MSG_EVENT 插入 + CMINMSGS 回填 + SUCCEEDED。 - 异常边界(U08):
processOne内 try/catch →ProcFailure.fail(INFRA)(attempts+1、 退避、达上限 DEAD);InterruptedException恢复中断位后上抛;loop 仅 catchException作最后防线,Error任其终止进程(异常必可见)。
幂等键(I3):SNDR|TYPE|STYP|SEQN(Identity.of 唯一入口);「含日边界」可配置且
默认关闭(SEQN 重置作用域 CONFIRM 前,上线后不改幂等键)。
3.3 流程 3:投递(Dispatcher)
- 每 target 严格 FIFO(I1 双层同策略):
headUnsent取队头;队头退避未到期 → 等待不跳过。 - 逐条循环显式排除
KAFKA_SCHD(D3/U06):schd 唯一出口flushSchd——claimBatch(ORDER BY EVENT_ID)→ 按 FLID 分组取 max(EVENT_ID)(FIX:现役 buffer 无去重会重发旧值)→ FLTR JSON 数组一次发出;失败整批 attempts+1 退避、队首未到期不 claim、lastFlush仅成功后推进;达上限整批 DEAD(DLQ)。周期/批上限取参数表 (3s / 500)。 - 轮询间隔取参数表(下限 50ms,N18);无 200ms 硬编码。
- 阶段 B(定案 2/D2):ES 投递成功 → 同线程同步
insertSync删除事件 (deleteOf:Jackson 结构化序列化,refs 可空恒合法 JSON——U14)。
3.4 流程 4:日计划快照(SnapshotFlow)
staging(流式解析+整包校验,TODO 阶段2;未实装→FAILED(UNSUPPORTED))
→ Redis Lua SNAPSHOT_REPLACE(同一 hash 原子「覆盖新代+按代差删」,删除集=旧代flids−新代)
→ putGenIfVersion CAS(version 未变才写;CAS 后重放=version 已达标→no-op 成功)
→ SUCCEEDED
CAS 冲突 → FAILED(INFRA)+退避(串行泵下不应发生→告警语义)
已知缺口(U09,未定案):Lua 与 CAS 分属两存储,非同一事务;崩溃窗口 (Lua 后/CAS 前、CAS 后/SUCCEEDED 前)与幂等重放判据(版本不二次自增)的显式恢复协议 待定案并补测试。
3.5 泵作业(JobExecutor,经 PUMP_JOB 同队列)
- HISTORY_SWEEP(3:30 清场,I4 同步链):判史 → 同步写 ES → 仅删成功集。
占位门禁(U10/T07 修订):ES saveSync 接线前
pickHistory恒空集、删除量恒 0, 禁止「全量可删」fail-open 默认;现役五条判史规则 golden 通过后才允许接线。 - ARCHIVE(3:00):1 天前且仅终态(SUCCEEDED/SKIPPED/DEAD)可迁 CMINMSGS_HST(TODO)。
- PROJECTION_REBUILD(阶段 B 切入时全量重建,TODO activeDays)。
4. 失败与重试统一设计(U08)
| 侧 | 组件 | 迁移语义 |
|---|---|---|
| ProcState(处理/快照) | ProcFailure + FailureScheduler |
attempts+1 → exhausted ? DEAD(EXHAUSTED) : FAILED+nextAttemptAt |
| MsgEvent 逐条(投递) | Dispatcher.retryOrDead |
同上(DLQ 保留行,attempts 审计) |
| MsgEvent 批量(schd) | flushSchd 整批 |
队首未到期不 claim;整批退避;达上限整批 DEAD |
- 退避表
[1s,2s,4s,8s,16s],单档封顶backoff-cap-ms=60s;attempt≤0兜底首档(N28)。 - 时间一律经可注入
java.time.Clock(TimeFactory;测试用 MutableClock,无真实睡眠)。 - loop 兜底 catch 不做状态迁移(迁移已在边界完成),仅防线程静默死亡。
5. 不变量与实现落点
| 不变量 | 语义 | 落点 | 状态 |
|---|---|---|---|
| I1 | 单写者严格 FIFO + HOL 阻塞 + 毒丸升级 | headUnfinished/headUnsent 队头语义、poisoned() |
实装(job 相对队头的全序属 U15) |
| I2 | Redis 先写、后于事件创建(happens-before) | processOne 阶段 A 分支 |
TODO redisApply(流程占位已留) |
| I3 | identity 首绑幂等;接收层无唯一约束;SUCCEEDED 回填 | Identity/tryBindIdentity/backfillOnSuccess |
实装 |
| I4 | 清场仅删 ES 成功集;按代差删 | HistorySweepJob/SNAPSHOT_REPLACE delFields |
门禁实装,ES 接线 TODO |
| I5 | 阶段 A Redis 写仅主泵线程;Delivery 不写 Redis | 单线程拓扑 + Targets.phaseA | 实装(拓扑约束,U26 运行期保护未做) |
6. 配置参数(msgx.*,ACMA-8 参数表初值)
| 键 | 默认 | 说明 |
|---|---|---|
phase |
A | 阶段总开关(A/B 权威切换) |
service-name |
msgexchangeapi | 契约冻结;影子= msgexchangeapi-shadow |
pipeline.poll-interval |
1s | 泵轮询节律(KEEP 现役) |
pipeline.max-attempts |
5 | 处理/投递同值 |
pipeline.backoff-ms / backoff-cap-ms |
1s..16s / 60s | 指数退避表与封顶 |
pipeline.head-deadline |
10m | 队头滞留上界(毒丸升级) |
pipeline.autostart |
false | 生命周期门禁:true 才装配 Pump/Dispatcher 线程(dev+stubs 开) |
schd.flush-period / flush-limit |
3s / 500 | schd 聚合节律与批上限 |
identity.include-day-boundary |
false | 幂等键日边界(CONFIRM 前禁开) |
consistency-check.on-startup / daily-sample-ratio |
true / 0.01 | 一致性哨兵(实装属 U25) |
基础设施键位口径(Micronaut 5.1,U03):datasources.default.*、
flyway.datasources.default.*、kafka.producers.default.*、eureka.client.*;
logback 独立于本文件,环境变量前缀 MSGX_LOGSTASH_*。
7. 测试策略
- 接口驱动 + 假仓储:管道语义全部离线单测(无 DB/Redis/Kafka),时间用
MutableClock。 - 不变量测试:FIFO/HOL、schd 批退避与 DLQ、毒丸升级、重试上限、重放白名单、 聚合最新态、配置绑定、DI 装配冒烟(PipelineSmokeTest 端到端:收报→FAILED(UNSUPPORTED) →重放→schd 聚合发出)。
- 现状 37 测试全绿(
./gradlew test;wrapper 钉 9.6.1 + JDK 25;内网构建设GRADLE_USER_HOME/TMPDIR指向可写目录)。 - U29 门禁(规划):不变量清单化入 CI 红即阻塞;「新增逻辑必伴生不变量测试」入贡献约定。
8. 可观测性设计(U12 已落地部分)
- logstash TCP 经 AsyncAppender(queueSize 4096 / neverBlock / discardingThreshold 0): logstash 不可达丢弃日志而非阻塞业务线程(N31b 顺序约束:先降级再补日志)。
- MDC
traceId= cminmsgsId/eventId(TraceLog.withTrace),覆盖处理/投递日志片段; 收报→投递全链贯穿与 micrometer gauge(队列深度/投递延迟)未实装。 - 生命周期结构化日志:收报 INFO、SUCCEEDED/SKIPPED INFO、FAILED WARN、DEAD/毒丸 ERROR、 flush 批次 INFO、整批 DLQ ERROR——告警暂以 ERROR 日志为落点(DEAD 告警出口属 U13/U25)。
9. 已知缺口与定案待办(对照 ACM2-10)
| 项 | 缺口 | 计划 |
|---|---|---|
| U05 | 仓储接口无实装;@Transactional/allopen 未引入 → 生产 DI 装配失败、事务1/2 未原子 |
阶段 1(Micronaut Data JDBC + allopen + 事务生效回归) |
| U07/U26 | autostart 默认关=有意门禁,但生产无 fail-fast;双实例无运行期防护 |
fail-fast 定案 + 租约/DB 锁拒启 |
| U09 | 快照跨存储恢复协议未定案 | 崩溃窗口清单 + 幂等判据 + 测试 |
| U13 | 投递侧无 createdAt/headDeadline 超时升级、无 DEAD 告警出口 | WP2 |
| U15 | job 与队头消息无统一全序(入队时间近似) | 统一序号列定案 |
| U16 | /cminmsgs/send 无 @Consumes/字符集、无错误路径契约(@ControllerAdvice) |
WP2(legacy 逐字对拍固化) |
| U18 | §2 所列 DDL 缺口 | WP2 |
| U19–U21 | 请求状态机量纲/死分支、identity 绑定静默跳过、REF_DATA 接口/DDL 对齐 | WP2 |
| U22–U24 | 载荷类型收敛、eventSeq 未接线、每报文全量读语义定案 | WP3 |
| U25/U28/U30 | 一致性哨兵实装、README 安全节/入口、索引与杂项 | WP3/4 |
| U27 | 专有材料(SIS md 703KB / XSD 版权头)治理决策 | WP4(ACL 核验先行) |