1. 页面概览
1.1 是什么
「流事件化规则」页面(页面内标题为 流事件化规则)是 LightFoundry 的 CDC 事件 → 规则条件求值 → 执行动作 的自动化闭环管理界面,V5 Stage 3(B3-4)交付。页面描述把整条链路浓缩为:「CDC 事件 → 规则条件求值({{change_type}}/table/pk_value/payload)→ 执行动作(webhook/通知/同步/画像/写回)。」
在数据集成场景里,数据库的表变化(插入、更新、删除)会被 CDC(Change Data Capture)检测器捕获并发布到平台事件总线(platform/cdc),每条事件带有 change_type(insert/update/delete/snapshot)、table_name、pk_value、payload 等信息。业务上常常需要「当订单表插入一条新记录时立刻通知主管」「当某张表的金额超过阈值时触发一次数据同步」「当订单状态变为已发货时把状态写回另一套系统」这类实时响应。本页面就是配置这些响应规则的地方。
一条规则 = 绑定一个 CDC watcher(数据源)+ 一个条件表达式 + 一个动作。规则启用后即订阅该 watcher 的全部事件;事件到达时后端对条件求值,命中则执行动作并落一条投递记录,未命中则不产生任何记录。页面同时提供「测试触发」——构造一条模拟事件直接走完整处理流水线(不经过事件总线订阅),以及「投递历史」抽屉——查看每次动作执行的 event_id、状态与错误。
1.2 核心价值
| 维度 | 说明 |
|---|---|
| 事件化响应 | 把数据库变更变成可编程的事件流,按规则自动触发动作,无需人工轮询 |
| 条件精确筛选 | 条件表达式支持 {{var}} 语法与 == != > < && || ! 运算,可精确圈定「什么事件才响应」 |
| 动作类型丰富 | 五种动作:webhook(HTTP POST JSON)、notification(审计+通知)、sync_run(触发同步引擎)、quality_run(触发质量画像)、writeback(Action 写回流水线) |
| 实时测试 | 「测试触发」构造模拟事件直接走处理流水线,不落订阅即可验证规则是否生效 |
| 投递可追溯 | 每次动作执行落一条投递记录(event_id/状态/错误),命中即留痕,未命中不产生记录 |
| 幂等防重 | (rule_id, event_id) 复合唯一索引保证同事件不重复执行动作,重放安全 |
1.3 一句话总结
2. 访问入口
2.1 路由与菜单
- 路由:
/foundry/stream - 路由名称:
FoundryStream - 路由 meta:
title: 流事件规则,requiresAuth: true,挂在父路由/foundry(FoundryLayout)下 - 菜单位置:Foundry 左侧边栏「流事件规则」(FoundryLayout 菜单项)
- 前端源码:
action/web/src/views/StreamPage.vue - API 客户端:
action/web/src/api/streamApi.js
2.2 认证与权限
- 页面路由挂
requiresAuth: true,未登录访问重定向到/login。 - 数据请求走
streamApi实例:baseURL/api/v1,请求拦截器自动附带Authorization: Bearer <aip_token>;响应拦截器 401 时清理aip_token/aip_username并跳转/login。 streamApi用unwrap()解析统一响应体{code, data}:code !== 0时抛出body.error || body.message,页面在 alert 中展示。- 后端
/stream/*与/cdc/watchers均挂 protected 语义,未带有效令牌返回 401。 - Watcher 列表(
GET /cdc/watchers)不可用时(如 CDC 未接线)前端不阻断规则管理,仅 watchers 下拉为空。
2.3 端口与 API 前缀
- Foundry 后端端口:18081。
- API 前缀:
/api/v1(streamApi 的 baseURL)。 - 完整请求示例:
GET /api/v1/stream/rules、POST /api/v1/stream/rules/:id/test。
3. 界面布局
页面为单栏布局,从上到下依次为工具栏、规则编辑表单(条件显示)、规则列表、测试触发面板(条件显示)、投递历史抽屉(条件显示):
┌────────────────────────────────────────────────────────────────┐
│ 流事件化规则(页头 + 描述) │
│ [alert 操作结果提示条(可关闭)] │
│ 工具栏:[新建规则] [刷新] │
├────────────────────────────────────────────────────────────────┤
│ [新建/编辑规则表单,v-if form.show] │
│ 规则名称 * | Watcher(CDC 检测源)* │
│ 条件表达式(可空=恒真;{{var}} 语法) │
│ 动作类型 * | 启用(保存后立即订阅该 watcher 事件) │
│ 动作配置 action_config(JSON)* [保存] [取消] │
├────────────────────────────────────────────────────────────────┤
│ 流规则(N) 表格:名称 / Watcher / 条件 / 动作 / 状态 / 触发次数 / │
│ 最近触发 / 操作(测试|投递|编辑|停用|启用|删除) │
├────────────────────────────────────────────────────────────────┤
│ [测试触发:规则名,v-if testRule] │
│ change_type / table_name / pk_value / payload(JSON 摘要文本) │
│ [触发测试] 命中= 状态= (skipped/error) │
├────────────────────────────────────────────────────────────────┤
│ [投递历史:规则名(N),v-if deliveryRule] │
│ 表格:event_id / 状态 / 错误 / 时间 [关闭] │
└────────────────────────────────────────────────────────────────┘
各板块职责:
- 工具栏:「新建规则」打开编辑表单;「刷新」同时重载规则列表与 watcher 列表。
- 新建/编辑规则表单:录入规则名称、绑定 watcher、条件表达式、动作类型与动作配置 JSON、启用开关;「保存」提交到后端并自动同步订阅。
- 流规则列表:全部规则的汇总表,行内可直接「测试」「投递」「编辑」「停用/启用」「删除」。
- 测试触发面板:构造模拟事件(change_type/table_name/pk_value/payload),走完整处理流水线验证规则;展示命中、状态、跳过原因或错误。
- 投递历史抽屉:某条规则的全部投递记录,含 event_id、状态(success/failed)、错误、时间。
4. 交互元素详解
4.1 工具栏
| 元素 | 含义 | 操作效果 | 触发后端调用 |
|---|---|---|---|
| 「新建规则」按钮 | 打开新建表单 | 重置表单(watcher 取列表第一个的 id,action_type=webhook,enabled=true),标题显示「新建规则」 | 无(仅打开表单) |
| 「刷新」按钮 | 重载数据 | 同时执行 fetchRules 与 fetchWatchers | GET /stream/rules、GET /cdc/watchers |
4.2 新建/编辑规则表单
| 元素 | 含义 | 必填与默认值 | 操作效果 | 触发后端调用 |
|---|---|---|---|---|
| 规则名称 * | 规则的可读名称 | 必填,placeholder 如 订单插入→通知主管 | 校验后保存 | 写入 name |
| Watcher(CDC 检测源)* | 绑定哪个 watcher | 必选下拉,选项文案 name(table_name,datasource_id) | 决定订阅哪个事件源 | 写入 watcher_id |
| 条件表达式 | 求值条件 | 可空(空=恒真);{{var}} 语法,支持 == != > < && || !,placeholder 如 {{change_type}} == 'insert' && {{table}} == 'orders' | 命中才执行动作 | 写入 condition |
| 动作类型 * | 动作分型 | 必选下拉:webhook(HTTP POST JSON)/ notification(审计+通知)/ sync_run(触发同步引擎)/ quality_run(触发质量画像)/ writeback(Action 写回流水线),默认 webhook | 切换时若 action_config 为空自动填充模板 | 写入 action_type |
| 启用 | 是否启用规则 | checkbox,默认勾选 | 勾选表示「保存后立即订阅该 watcher 事件」 | 写入 enabled |
| 动作配置 action_config(JSON)* | 按动作类型的配置 | 必填 JSON 文本区,按类型自动填充 placeholder 与模板 | 前端先 JSON.parse 校验,非法则提示「action_config 不是合法 JSON」不提交 | 写入 action_config |
| 「保存」按钮 | 提交表单 | — | 创建提示「规则已创建(已订阅事件源)」/ 更新提示「规则已更新(订阅已同步)」 | POST /stream/rules 或 PUT /stream/rules/:id |
| 「取消」按钮 | 关闭表单 | — | 仅隐藏表单,不提交 | 无 |
4.3 规则列表
| 元素 | 含义 | 操作效果 | 触发后端调用 |
|---|---|---|---|
| 「测试」按钮 | 打开测试触发面板 | 预填 change_type=insert、table_name 空,显示「测试触发:规则名」 | 无(仅打开面板) |
| 「投递」按钮 | 打开投递历史抽屉 | 加载最近 100 条投递记录 | GET /stream/rules/:id/deliveries?limit=100 |
| 「编辑」按钮 | 打开编辑表单 | 回填规则字段,标题「编辑规则」 | 无(打开表单) |
| 「停用/启用」按钮 | 切换 enabled | 提示「规则已停用(订阅已退订)」或「规则已启用(已订阅事件源)」 | PUT /stream/rules/:id(仅传 {enabled: !enabled}) |
| 「删除」按钮 | 删除规则 | window.confirm 二次确认「确认删除规则「name」?投递历史保留。」后删除 | DELETE /stream/rules/:id |
4.4 测试触发面板
| 元素 | 含义 | 默认值 | 操作效果 | 触发后端调用 |
|---|---|---|---|---|
| change_type 下拉 | 事件类型 | insert | 可选 insert/update/delete/snapshot | 测试事件字段 |
| table_name 输入 | 表名 | 空 | 为空时后端测试自动用 watcher_id 填充 | 测试事件字段 |
| pk_value 输入 | 主键值 | 空,placeholder 如 1001 | 模拟事件主键 | 测试事件字段 |
| payload 输入 | JSON 摘要文本 | 可空,placeholder 如 {"count":150,"watermark":"2026-08-29"} | 事件的 payload | 测试事件字段 |
| 「触发测试」按钮 | 执行模拟触发 | 触发中置灰显示「触发中...」 | 展示 命中= 状态= ,skipped 或 error 附注 | POST /stream/rules/:id/test |
| 测试结果提示 | 结果摘要 | — | test-success(绿)或 test-failed(红) | — |
4.5 投递历史抽屉
| 元素 | 含义 | 操作效果 |
|---|---|---|
| event_id | 事件唯一 ID | 等宽字体展示 |
| 状态 | success/failed | success 显示绿色徽标,其余红色 |
| 错误 | 失败原因 | 超出宽度省略,悬停 title 查看完整 |
| 时间 | 投递创建时间 | YYYY-MM-DD HH:MM:SS 格式 |
| 「关闭」链接按钮 | 关闭抽屉 | 清空 deliveryRule |
| 空态提示 | 暂无投递记录 | 显示「暂无投递记录(条件未命中不会产生记录)」 |
5. 后端关联
5.1 API 客户端
- 文件:
action/web/src/api/streamApi.js - baseURL:
/api/v1;超时 30000ms;Content-Type: application/json - 拦截器:请求自动附带
Authorization: Bearer <aip_token>;响应 401 清除令牌并跳登录 unwrap():解析{code, data},code !== 0抛body.error || body.message- 导出函数:
listStreamRules / getStreamRule / createStreamRule / updateStreamRule / deleteStreamRule / testStreamRule / listStreamDeliveries / listWatchers,以及常量ACTION_TYPE_LABELS
5.2 端点表
| 方法 | 路径 | 请求体/参数 | 说明 |
|---|---|---|---|
| GET | /stream/rules | — | 规则列表(created_at 倒序) |
| POST | /stream/rules | CreateRuleRequest | 创建规则(创建后同步订阅) |
| GET | /stream/rules/:id | — | 规则详情 |
| PUT | /stream/rules/:id | UpdateRuleRequest | 更新规则(指针语义,未提供保留原值) |
| DELETE | /stream/rules/:id | — | 删除规则(硬删;投递历史保留) |
| POST | /stream/rules/:id/test | TestRuleRequest | 测试触发(构造模拟事件,不落订阅) |
| GET | /stream/rules/:id/deliveries | ?limit= | 投递历史(默认 100,上限 500) |
| GET | /cdc/watchers | — | CDC watcher 列表(规则绑定目标来源) |
5.3 响应结构
规则对象(rules 数组元素,{code:0, data:{rules:[...]}}):
{ "id": "...", "name": "订单插入→通知主管", "watcher_id": "w_orders",
"condition": "{{change_type}} == 'insert' && {{table}} == 'orders'",
"action_type": "webhook", "action_config": "{\"url\":\"http://host:8080/hook\"}",
"enabled": true, "last_fired_at": "2026-08-30T...", "fire_count": 3,
"created_at": "...", "updated_at": "..." }
投递记录(deliveries 数组元素):
{ "id": "...", "rule_id": "...", "event_id": "evt_xxx", "status": "success",
"error": "", "created_at": "2026-08-30T..." }
测试触发结果(data 直接为 TestResult):
{ "rule_id": "...", "event_id": "evt_auto_xxx", "matched": true,
"status": "success", "error": "", "skipped": "" }
5.4 关联模块表
| 后端包 | 职责 |
|---|---|
products/foundry/stream/rest.go | REST 端点、TestRuleRequest 定义、统一响应 |
products/foundry/stream/models.go | 表结构:fst_rules / fst_deliveries;动作类型与投递状态常量;action_config 分型结构 |
products/foundry/stream/service.go | 规则 CRUD、订阅生命周期(Start/Stop/SyncSubscription)、事件处理流水线、五种动作执行、TestRule、ListDeliveries |
products/foundry/stream/webhook_validate.go | webhook 目标 SSRF 校验(fail-closed) |
platform/cdc | 事件总线与 watcher(Subscribe/Publish 语义,只调用不修改) |
products/foundry/workflow | 条件表达式(EvalCondition)与 {{var}} 模板(ResolveTemplate) |
products/foundry/writepath | writeback 动作的 VALIDATE_AND_EXECUTE 五步流水线 |
5.5 关键机制
- 订阅生命周期:
Start(ctx)对每条 enabled 规则执行bus.Subscribe(watcherID, handler);handler 在事件到达时重查规则(防陈旧配置/已删除/已停用)再处理;CRUD 内部自动同步订阅;SyncSubscription支持外部直改库后手动重载。 - 条件求值:
eventVars把 payload 字段顶层展开,并保留payload / change_type / table / watcher_id / watcher_name / datasource_id / pk_value / event_id / occurred_at;EvalCondition求值失败记日志不报错,未命中不产生投递记录。 - 幂等防重入:先落一条 status=failed、error="pending" 的投递记录占位,(rule_id, event_id) 复合唯一索引冲突时跳过(幂等重放);动作执行后回写终态 success/failed。
- 动作语义:webhook 为 HTTP POST JSON(超时 10s、失败重试 1 次、请求前 SSRF 校验 fail-closed);notification 为审计落点 + 通知渠道(未注入降级为审计+投递记录不报错);sync_run / quality_run 依赖注入接口(未注入记录不可用);writeback 经 writepath Executor VALIDATE_AND_EXECUTE 走全五步流水线,幂等键
stream:<rule>:<event>。 - 触发计数:命中并执行后规则表
fire_count + 1、last_fired_at更新(best-effort,失败不阻断投递结果)。
6. 核心流程详解
6.1 主流程:建规则 → 订阅 → 事件到达 → 求值 → 执行 → 落投递
- 新建规则:点「新建规则」,填写规则名称、选择 watcher(需先在 CDC watcher 数据源存在,watchers 由 cdc REST 创建)、写条件表达式、选动作类型并填 action_config JSON(可先切换动作类型自动填充模板)、勾选启用,点「保存」。
- 订阅生效:保存后前端提示「规则已创建(已订阅事件源)」;后端对 enabled 规则
bus.Subscribe(watcherID, handler),此后该 watcher 的全部事件都会进入规则。 - 事件处理:事件到达 → handler 重查规则(确保仍启用)→
EvalCondition求值 → 命中则幂等 claim 占位 → 执行动作 → 回写投递终态 → 更新规则触发计数。 - 验证与追溯:用「测试」构造模拟事件直接走处理流水线(不落订阅)验证条件与动作;用「投递」抽屉查看每次执行的 event_id、状态与错误。
6.2 条件表达式要点
- 变量引用:
{{change_type}}、{{table}}、{{pk_value}}、{{payload}},payload 内字段顶层展开可直接{{count}}。 - 运算符:
== != > < && || !(注意 HTML 中需写><&&)。 - 示例:
{{change_type}} == 'insert' && {{table}} == 'orders';条件留空表示恒真(每条事件都触发)。 - 求值失败:后端记 warn 日志并跳过该事件(不产生投递记录),不会误触发动作。
6.3 分支流程:测试触发
- 点规则行的「测试」打开面板(预填 change_type=insert)。
- 修改 change_type、table_name、pk_value、payload。
- 点「触发测试」:前端组装
{watcher_id, table_name: table_name || watcher_id, change_type, pk_value, payload}提交。 - 结果展示:
命中=true 状态=success(绿)/命中=false 状态=(条件未命中,未触发动作)/命中=true 状态=failed 错误=xxx(红)/ 幂等重放时 skipped 提示。 - 测试会真实执行动作并落投递记录(webhook 会真的发起 HTTP 请求),测试后规则 fire_count 会增加。
6.4 投递记录语义
- 只有条件命中才产生投递记录;未命中无记录(抽屉空态文案明确说明)。
- status 只有 success / failed 两种终态;执行中的占位记录 error="pending" 理论上只在崩溃残留时可见。
- 同规则同事件重复触发(幂等重放)被跳过,不产生第二条投递记录,也不重复执行动作。
7. 权限与安全
- 认证:全部端点经 authMiddleware(JWT,aip_token),未授权 401。
- webhook SSRF 防护:所有 webhook 触发路径(订阅触发与 TestRule 模拟触发)均经
validateWebhookTarget校验目标 URL 与 method,拒绝内网/私网/云元数据等目标,fail-closed。 - 动作依赖注入:sync_run / quality_run / writeback 依赖后端注入的接口,未注入时记录不可用(failed + 明确错误),不会静默失败。
- 写操作防护:删除规则用
window.confirm二次确认且提示「投递历史保留」;停用/启用直接提交。 - 幂等安全:(rule_id, event_id) 唯一索引防止重复执行动作,重放/重复投递安全。
8. 常见问题与排错
问题一:保存规则时提示「action_config 不是合法 JSON」并拒绝提交
现象:点保存后 alert 报 JSON 非法,表单不提交。
原因:动作配置文本区输入了非法 JSON(多余逗号、单引号、注释等)。
排查步骤:检查文本区内容是否与所选动作类型的模板结构一致(可先切换动作类型让前端填充模板再修改);用 JSON 校验工具验证;webhook 需 {"url":"http://..."},writeback 需 {"action_name":"...","object_type_id":N,"params":{}},sync_run 需 {"target_id":"..."},quality_run 需 {"profile_id":"..."}。
问题二:Watcher 下拉为空,无法新建规则
现象:新建表单的 watcher 下拉没有选项,规则列表空态提示「请先创建 CDC watcher 再新建规则」。
原因:GET /cdc/watchers 返回空或不可用(CDC 未接线/未创建 watcher)。
排查步骤:确认后端 CDC watcher 是否已创建(通过 cdc REST 接口);检查浏览器 Network 里 /api/v1/cdc/watchers 的返回;确认后端已启动并注入 cdc 事件总线。
问题三:触发测试显示「命中=false(条件未命中,未触发动作)」
现象:测试结果 matched=false 且 skipped 提示条件未命中。
原因:条件表达式与模拟事件不匹配(change_type、table、payload 字段对不上)。
排查步骤:核对条件中引用的变量是否都在测试面板中构造(如条件用 {{count}} 则 payload 必须包含 count 字段);确认 {{table}} 与 table_name 一致;若条件留空仍不命中,检查后端日志的求值警告。
问题四:触发测试命中但状态 failed,投递记录错误列有内容
现象:matched=true 但 status=failed,投递抽屉 error 列有内容。
原因:动作执行失败——webhook 目标不可达/非 2xx、sync_run 未注入、writeback 参数不对等。
排查步骤:打开「投递」抽屉查看 error 全文;webhook 检查 url 可达性与目标 SSRF 校验;sync_run/quality_run 确认后端注入对应服务;writeback 确认 action_name 存在且 object_type_id 正确。
问题五:规则启用了但实际事件到达时不触发
现象:规则显示启用,但真实 CDC 事件到达后没有任何投递记录。
原因:watcher 未发布事件、订阅未建立(bus 为 nil)、或事件到达时规则已被停用。
排查步骤:确认 watcher 的数据源确实有变更产生;查看后端日志「流事件规则订阅完成」的规则数;确认规则状态为「启用」;可先用「测试触发」验证规则本身可用。
9. 已知缺陷与边界
| 边界 | 说明 |
|---|---|
| 条件语法 | 仅支持 == != > < && || ! 与 {{var}},不支持函数调用/正则/大小写折叠(不同于函数页表达式语法) |
| watcher 依赖 | 规则必须绑定已存在的 CDC watcher;watcher 列表不可用时新建表单无法选源 |
| 动作依赖注入 | sync_run / quality_run / writeback 依赖后端注入接口,未装配时动作恒失败并记录 |
| notification 降级 | 通知渠道(B3-6)未注入时 notification 动作降级为审计+投递记录,不报错 |
| webhook 重试 | 仅网络/超时错误重试 1 次;非 2xx 响应不重试(立即失败) |
| 投递上限 | 投递历史默认返回 100 条、上限 500 条;前端固定请求 limit=100 |
| 删除语义 | 删除规则为硬删,投递历史保留但不随规则展示 |
| 幂等范围 | 幂等键为 (rule_id, event_id),不同规则对同一事件会各自执行一次动作 |