数据规则完整执行 M3B 阶段验收
日期:2026-07-23
分支:codex/data-rule-execution-m3a-m5
范围:Task 5–6(跨源 Polars 执行、产物交接、规则运行证据)
验收结论
M3B 通过。
平台已形成可复现的跨源数据规则执行链:受治理的已发布规则计划从
PostgreSQL、MySQL 和 MinIO 读取真实输入,在受操作系统资源限制的隔离
Polars worker 中执行,生成不可变 Parquet 产物;后续规则节点仅通过服务端
登记的 opaque artifact reference 消费上游结果。每个节点执行均绑定精确的
部署、环境、数据流版本、组件、规则版本、计划摘要和 correlation ID,并记录
成功、失败、取消或提交结果未知等证据。
M3B 不负责把编译计划发布为可执行计划。生产 Runner 仍只消费
published 规则与计划;发布门禁由 M4 Task 7 负责。
完成能力
跨源批处理与资源隔离
- 采用闭合 JSON Polars plan,不执行模型生成的 Python、pickle、callable
或 Polars 内部序列化计划。
- 支持 normalization、filter、lookup join、assert/reject/quarantine、
deduplicate、aggregate 和受治理 mask policy。
- Parquet 在解压前校验 compressed size、row count、row-group
uncompressed size、digest、schema、TTL 和 store ownership。
- Polars 计划在 spawn 子进程中执行,使用 RSS 监控、
RLIMIT_AS、
timeout 和安全 kill;主 Runner 不设置不可恢复的全局资源上限。
- Decimal precision/scale、timestamptz timezone、nullable 和 typed cast
在编译与运行时共同校验。
不可变产物交接
- MinIO key 仅由服务端生成,格式限定在
rules/<correlation-id>/<artifact-id>.parquet。
- PostgreSQL handoff 采用
pending -> ready/failed 状态机,并保存
binding hash、digest、schema hash、row count、expiry 和 artifact kind。
- 同一 correlation/binding/kind 的相同结果幂等复用;不同结果 fail closed。
- 上传、数据库提交结果未知和进程崩溃均有双向、限量、带 grace period 的
reconciliation;临时存储错误不会破坏 ready 结果。
- SQL 同源交接使用服务端登记的 opaque staging receipt,不接受调用方拼接
的可预测引用。
执行证据与违规样本
rule_runs 在 adapter 执行前创建,终态保存 counts、timings、
commit outcome、最小公开结果和证据摘要。
- 任务令牌签名 exact deployment/environment/dataflow version/node/JTI;
请求 body 不能声明证据身份。
- evidence key 跨不同签名 token 保持节点级幂等。
- running execution 使用 attempt、lease、heartbeat 和过期收敛;两个崩溃
窗口不会永久返回 202。
- 同 JTI 响应丢失可从 durable ledger/evidence 重放终态;过期 token
仅允许 exact terminal replay,不能启动新执行。
- 违规样本最多 100 行,在隔离 worker 与 evidence boundary 双重全字段脱敏,
不进入应用日志。
- violation sample 和 SQL staging receipt 均有 claim lease、过期接管、
confirmed-missing 与 transient-store-error 分类以及 bounded cleanup。
编排交接
- Kestra 单前驱节点使用 bracket-safe 的上游 output expression。
- 受治理规则节点不接收 raw row JSON 或无关 workflow parameters。
- 当前规则运行时只支持单一主输入,因此 governed fan-in 在编译期 fail closed,
不生成无法执行的
input_artifacts 伪契约。
迁移与兼容性
- 历史 migration 140、150、160、170、180 保持不可变。
- 新增 migration 190 后,本地 PostgreSQL 当前 head 为
20260723_190。
- 已验证旧 140 到当前 head、150 异常历史样本预检、170 到 180、
180 遗留 cleanup claim 到 190 的真实 PostgreSQL 升级。
- 190 升级前必须停止 revision-180 cleanup worker;历史非空 claim 会被
回填为已到期并可通过 CAS 接管,避免永久锁死。
验收证据
阶段聚焦测试
命令覆盖:
- Polars compiler、artifact store、isolated worker;
- artifact handoff、rule evidence、Runner API、task ledger;
- Kestra compiler;
- 历史 migration 升级;
- 真实 PostgreSQL + MySQL + MinIO 两节点执行。
结果:
118 passed, 3 skipped in 67.33s
真实链路验证:
- 第一节点从 PostgreSQL 主数据和 MySQL lookup 读取输入;
- isolated Polars worker 执行 normalize/join/assert/deduplicate;
- 输出写入 MinIO 并登记 ready handoff;
- 违规行以全字段脱敏 Parquet 样本保存;
- 第二节点消费第一节点的 exact artifact ref 并产生新产物;
- 两条 rule run 证据与同 correlation 对齐;
- 同 token/不同 token retry 不重复执行;
- failed、cancelled、unknown commit、lease expiry 和 response-loss
均按可信状态收敛。
全量 Python 回归
606 passed, 28 skipped, 59 subtests passed in 11.11s
Docker 状态
健康:
- backend
- frontend
- runner
- PostgreSQL
- source PostgreSQL
- source MySQL
- MinIO
- Neo4j
- n8n
已知待办:
- Kestra 当前为
Exited (1),日志显示 PostgreSQL broken pipe/closed
connection 后未捕获队列线程异常。该问题属于 M5 Task 9 的专用数据库、
连接稳定性、restart policy、health dependency 和 30 分钟 soak 验收,
不属于 M3B Runner/产物/证据验收范围。
开源依赖与许可证
运行依赖已固定并记录在 CycloneDX 1.5 SBOM:
- Polars 1.42.1 / polars-runtime-32 1.42.1:MIT
- PyArrow 21.0.0:Apache-2.0
- MinIO Python client 7.2.10:Apache-2.0
- psutil 5.9.8:BSD-3-Clause
SBOM 使用固定的官方 CycloneDX、SPDX 与 JSF schema 进行离线校验。
进入 M4 的条件
以下 M3B 输入已冻结,可进入 M4:
RuleVersion + ExecutionPlan + SchemaSnapshot + DatasetBinding 的 canonical
attestation;
- compiled-only plan persistence 与 published-only Runner boundary;
- Polars/SQL backend result contract;
- immutable artifact/staging handoff;
- exact deployment execution identity;
- trusted compile/runtime evidence primitives;
- replay、lease、cleanup 和 reconciliation contracts。
M4 必须在这些边界之上实现可信 generation receipt、compile/test evidence
门禁、计划 compiled -> tested -> published 状态迁移,以及 Data Standard /
Data Flow 的统一资产选择体验,不能回退到 inline rule 或调用方自报 evidence。