diff --git a/IMPLEMENTATION.md b/IMPLEMENTATION.md
index a29b1a5..d9aa6da 100644
--- a/IMPLEMENTATION.md
+++ b/IMPLEMENTATION.md
@@ -10,8 +10,9 @@
- `frontend`:Vue3 + TypeScript 前端。
- `infra`:PostGIS、Redis、Kafka、MinIO、Prometheus、各服务 Docker Compose 编排。
- `docs/IMPLEMENTATION_COMPLIANCE.md`:设计说明书到实现的映射关系。
+- `docs/full-chain-visualization-design.md`:全链路示例可视化展示设计方案。
-当前 AI 服务保留确定性推理适配器,接口、模型组、返回结构均按真实模型接入方式设计。后续接入 YOLO、Mask R-CNN、SAM、ONNX Runtime、TensorRT、Open3D/PCL、GDAL 等模型或算法时,替换对应 FastAPI 服务内的推理适配逻辑即可。
+当前 AI 服务保留确定性推理适配器,接口、模型组、返回结构均按真实模型接入方式设计。全链路示例已实现运行实例、阶段事件、SSE 实时推送和前端运行看板。后续接入 YOLO、Mask R-CNN、SAM、ONNX Runtime、TensorRT、Open3D/PCL、GDAL 等模型或算法时,替换对应 FastAPI 服务内的推理适配逻辑即可。
## 2. 目录结构
@@ -53,6 +54,11 @@ POST /api/v1/exemptions
GET /api/v1/models
GET /api/v1/stats/overview
POST /api/v1/demo/run-full-chain
+POST /api/v1/demo/runs
+GET /api/v1/demo/runs/{run_id}
+GET /api/v1/demo/runs/{run_id}/events
+GET /api/v1/demo/runs/{run_id}/events/history
+POST /api/v1/demo/runs/{run_id}/cancel
```
## 5. 验证
@@ -61,7 +67,7 @@ POST /api/v1/demo/run-full-chain
.\scripts\smoke_test.ps1
```
-冒烟脚本默认通过前端 Nginx 反向代理访问 API,即 `http://localhost:8088`。如需直连后端,可设置 `RAIL_API_BASE_URL=http://localhost:8080`。
+冒烟脚本默认通过前端 Nginx 反向代理访问 API,即 `http://localhost:8088`。脚本会启动新的可视化全链路运行实例,轮询运行快照直到完成,并校验工单接口与前端入口。如需直连后端,可设置 `RAIL_API_BASE_URL=http://localhost:8080`。
离线静态校验:
@@ -105,5 +111,6 @@ python -m pytest
1. 启动完整技术栈。
2. 打开 Web 前端。
3. 点击“运行全链路示例”。
-4. 查看任务、AI 结果、告警、工单、模型版本。
-5. 执行 `scripts/smoke_test.ps1` 验证 API 返回。
+4. 查看全链路运行看板中的阶段时间轴、GIS态势、阶段详情、事件流和完成汇总。
+5. 查看任务、AI 结果、告警、工单、模型版本。
+6. 执行 `scripts/smoke_test.ps1` 验证 API 返回。
diff --git a/docs/IMPLEMENTATION_COMPLIANCE.md b/docs/IMPLEMENTATION_COMPLIANCE.md
index 630fc0a..c908621 100644
--- a/docs/IMPLEMENTATION_COMPLIANCE.md
+++ b/docs/IMPLEMENTATION_COMPLIANCE.md
@@ -30,6 +30,7 @@
| 工单闭环 | `/api/v1/workorders/callbacks/status` |
| 基线库与豁免区 | `baseline_objects`、`exemption_areas` 表结构和豁免接口 |
| 数据标注与模型迭代 | `model_versions` 表和模型版本接口,样本库可继续接 MinIO 路径 |
+| 全链路可视化 | `demo_runs`、`demo_run_events`、SSE 事件流、前端运行看板 |
| 系统运维 | Actuator、Prometheus、Docker Compose |
## 3. 接口落实情况
@@ -46,6 +47,11 @@
| 工单回调 | `POST /api/v1/workorders/callbacks/status` |
| 豁免区配置 | `POST /api/v1/exemptions` |
| 模型版本 | `GET /api/v1/models` |
+| 全链路运行实例 | `POST /api/v1/demo/runs` |
+| 全链路运行快照 | `GET /api/v1/demo/runs/{run_id}` |
+| 全链路事件流 | `GET /api/v1/demo/runs/{run_id}/events` |
+| 全链路历史事件 | `GET /api/v1/demo/runs/{run_id}/events/history` |
+| 取消全链路运行 | `POST /api/v1/demo/runs/{run_id}/cancel` |
## 4. 部署方式
@@ -71,8 +77,10 @@
验证内容:
- 后端健康检查
-- 全链路示例
+- 全链路可视化运行实例
+- 阶段事件与运行快照
- 统计接口
+- 工单接口
- 前端页面
离线环境下,如果无法访问 Docker Hub 或 Maven/npm/PyPI 仓库,可先执行静态校验:
@@ -87,7 +95,26 @@
- 前端 JSON 配置
- Docker Compose 配置合法性
-## 5.1 离线镜像要求
+## 5.1 全链路可视化实现
+
+运行“全链路示例”时,前端不再等待同步接口结束后一次性刷新,而是创建运行实例并订阅 SSE 事件流。后端按阶段发布事件,前端实时更新阶段时间轴、GIS 态势、当前阶段详情、事件流和完成汇总。
+
+已实现阶段包括:
+
+- 运行初始化
+- 巡检任务创建
+- 航线任务下发
+- 多源数据接入
+- 数据预处理
+- AI 推理分析
+- GIS 规则判定
+- 告警生成
+- 工单派发
+- 工单闭环
+- 样本回流
+- 完成汇总
+
+## 5.2 离线镜像要求
完整部署默认使用已验证可用的华为云 SWR Docker Hub 代理 `swr.cn-north-4.myhuaweicloud.com/ddn-k8s/docker.io`,并配置了阿里 PyPI、Maven 多源兜底、npmmirror npm 源。需要以下基础镜像可用:
diff --git a/docs/full-chain-visualization-design.md b/docs/full-chain-visualization-design.md
new file mode 100644
index 0000000..ce4c145
--- /dev/null
+++ b/docs/full-chain-visualization-design.md
@@ -0,0 +1,826 @@
+# 全链路示例可视化展示设计方案
+
+## 1. 设计目标
+
+运行全链路示例时,系统需要把“任务创建、无人机巡检、数据接入、预处理、AI 推理、GIS 规则判定、告警生成、工单派发、闭环处置、样本回流”完整展示为可观察、可追踪、可复盘的过程。展示方式不应只在执行结束后刷新统计数据,而应在执行过程中持续呈现每个阶段的状态、输入、输出、耗时、异常和业务结果。
+
+核心目标如下:
+
+- 每个环节都有明确的可视化承载区域,不出现空白占位。
+- 每个环节都有状态:未开始、执行中、成功、失败、跳过、重试中。
+- 每个环节都有业务实体关联:任务、资源、分析任务、AI 结果、告警、工单、样本。
+- 每个环节都有可理解的过程反馈:进度、关键指标、事件日志、空间位置、数据流向。
+- 全链路示例执行完成后,可保留本次运行记录,支持回看与问题定位。
+
+## 2. 当前问题与改造方向
+
+当前 `POST /api/v1/demo/run-full-chain` 为同步接口,后端一次性完成任务创建、资源接入、分析、告警和工单生成,前端只在接口返回后刷新列表与统计。因此用户无法看到中间过程,也无法判断耗时发生在哪个环节。
+
+改造方向:
+
+- 将“全链路示例”从同步动作升级为“示例运行实例”。
+- 后端按阶段执行,并在每个阶段发布事件。
+- 前端打开“全链路运行看板”,通过 SSE 或 WebSocket 实时订阅阶段事件。
+- 每个阶段事件携带可视化负载,前端按事件类型更新进度条、地图、卡片、表格、日志和工单流。
+- 执行结束后,保留运行快照,支持回放。
+
+## 3. 展示总览
+
+全链路运行看板采用一屏式工作台布局:
+
+```text
+┌────────────────────────────────────────────────────────────────────┐
+│ 顶部:运行控制区 │
+│ 运行编号 / 当前阶段 / 总进度 / 开始时间 / 耗时 / 重新运行 / 取消 │
+├────────────────────────────────────────────────────────────────────┤
+│ 阶段时间轴 │
+│ 任务创建 → 航线下发 → 数据接入 → 预处理 → AI推理 → GIS规则 → 告警 → 工单 → 样本闭环 │
+├───────────────────────────────┬────────────────────────────────────┤
+│ 左侧:GIS空间态势 │ 右侧:当前阶段详情 │
+│ 线路、航线、资源点、告警点 │ 阶段说明、输入输出、指标、结果 │
+│ 防护区、防洪区、工单位置 │ │
+├───────────────────────────────┴────────────────────────────────────┤
+│ 底部:多标签详情区 │
+│ 事件流 / 数据资源 / AI结果 / 规则命中 / 告警工单 / 样本回流 │
+└────────────────────────────────────────────────────────────────────┘
+```
+
+页面行为:
+
+- 点击“运行全链路示例”后,按钮进入运行态,并打开运行看板。
+- 阶段时间轴从左到右推进,当前阶段高亮并展示进度。
+- GIS 区域持续叠加航线、数据采集点、隐患点、规则区域和工单点位。
+- 右侧详情区根据当前阶段动态切换展示内容。
+- 底部事件流持续追加结构化事件,支持按阶段筛选。
+- 执行完成后,顶部显示总耗时、生成结果和闭环状态。
+
+## 4. 全链路阶段设计
+
+| 阶段 | 阶段编码 | 主要动作 | 可视化展示 | 阶段输出 |
+| --- | --- | --- | --- | --- |
+| 运行初始化 | `run.initializing` | 创建运行实例、锁定场景参数、初始化进度 | 运行编号、场景清单、线路范围、阶段时间轴置灰 | `run_id`、场景集、线路信息 |
+| 巡检任务创建 | `task.created` | 创建巡检任务,写入任务表,发布任务事件 | 任务卡片、线路里程范围、优先级、任务状态 | `task_id` |
+| 航线任务下发 | `uav.dispatching` | 模拟航线下发、任务状态回调 | GIS 航线动画、起降点、航线里程、无人机状态 | 航线状态、任务回调记录 |
+| 多源数据接入 | `resource.ingesting` | 接入可见光、红外、TIF/点云等资源 | 资源上传卡片、文件类型图标、接入进度、采集点 | `resource_ids` |
+| 数据预处理 | `preprocess.running` | 格式解析、时间同步、坐标归一、影像/点云预处理 | 处理队列、预处理步骤勾选、质量检查指标 | 标准化资源、质量报告 |
+| AI 推理分析 | `ai.inferencing` | 调用视觉推理和点云分析服务 | 模型矩阵、推理进度、识别结果缩略卡、置信度分布 | `ai_result_ids` |
+| GIS 规则判定 | `rule.evaluating` | 空间规则、业务阈值、豁免区、基线库匹配 | 地图规则区域、高亮命中点、规则命中列表 | 规则决策、抑制结果 |
+| 告警生成 | `alarm.generating` | 生成有效告警、分类分级、证据关联 | 告警级别分布、告警卡片、证据链、空间定位 | `alarm_ids` |
+| 工单派发 | `workorder.dispatching` | 自动创建工单、派发复核人员 | 工单流转泳道、派发状态、负责人、截止时间 | `workorder_ids` |
+| 工单闭环 | `workorder.closing` | 模拟处置反馈、关闭告警、更新工单 | 闭环进度、已处置/待处置数量、操作记录 | 工单状态、告警状态 |
+| 样本回流 | `sample.feedback` | 生成样本回流记录、模型迭代候选 | 样本队列、标注状态、模型版本影响指标 | 样本记录、模型迭代待办 |
+| 完成汇总 | `run.completed` | 汇总运行结果、耗时、数量、异常 | 总览指标、阶段耗时条、最终地图、结果摘要 | 运行快照 |
+
+## 5. 阶段可视化细节
+
+### 5.1 运行初始化
+
+展示内容:
+
+- 运行编号:`run_id`
+- 示例场景:违法开挖、塔吊、人员入侵、防洪点精密监测、接触网异物、杆上设备发热
+- 线路范围:`line-demo K123+000 - K130+000`
+- 预计阶段数:12
+
+界面元素:
+
+- 顶部运行状态条
+- 阶段时间轴
+- 全链路参数摘要
+
+### 5.2 巡检任务创建
+
+展示内容:
+
+- 任务 ID、线路、里程范围、优先级、任务状态
+- 任务创建事件:`inspection.task.created`
+- 任务写库结果
+
+界面元素:
+
+- 任务卡片从“待创建”切换为“已创建”
+- 时间轴第一个节点点亮
+- 事件流追加任务创建日志
+
+### 5.3 航线任务下发
+
+展示内容:
+
+- 航线 ID、起降点、里程范围、预计飞行时间
+- 无人机任务状态:已下发、执行中、返航、完成
+- 与巡检管理平台的状态回调
+
+界面元素:
+
+- GIS 地图绘制航线轨迹
+- 航线进度沿线路移动
+- 状态徽标显示“任务下发中 / 飞行中 / 已完成”
+
+### 5.4 多源数据接入
+
+展示内容:
+
+- 可见光图片、红外 TIF、防洪 TIF/点云资源
+- 每类资源的接入状态、大小、采集时间、空间位置
+- 存储路径:MinIO/S3 URL
+
+界面元素:
+
+- 资源卡片网格
+- 上传/接入进度条
+- 地图上出现采集点
+- 底部数据资源表追加行
+
+### 5.5 数据预处理
+
+展示内容:
+
+- 图像解码、抽帧、尺寸归一、去噪、质量检查
+- 红外温度矩阵解析、温度阈值提取
+- TIF/点云坐标解析、坐标系统一、体积变化指标提取
+
+界面元素:
+
+- 预处理步骤清单逐项打勾
+- 质量指标:清晰度、坐标有效性、时间同步状态、资源可读性
+- 异常时展示失败原因和可重试入口
+
+### 5.6 AI 推理分析
+
+展示内容:
+
+- 视觉模型:目标检测、分割、变化检测
+- 红外模型:设备发热识别
+- 点云/TIF 模型:防洪点体积变化分析
+- 每个模型的状态、耗时、结果数量、平均置信度
+
+界面元素:
+
+- 模型推理矩阵:
+
+```text
+┌────────────────────┬──────────┬──────────┬────────────┐
+│ 模型组 │ 状态 │ 结果数 │ 平均置信度 │
+├────────────────────┼──────────┼──────────┼────────────┤
+│ vision-detector │ 已完成 │ 12 │ 0.90 │
+│ thermal-analyzer │ 已完成 │ 3 │ 0.90 │
+│ pointcloud-analyzer│ 已完成 │ 3 │ 0.93 │
+└────────────────────┴──────────┴──────────┴────────────┘
+```
+
+- 识别结果卡片:场景、类别、置信度、模型版本、资源证据
+- 置信度分布条形图
+
+### 5.7 GIS 规则判定
+
+展示内容:
+
+- 线路 100 米范围内识别
+- 防护网内人员入侵
+- 接触网/线缆 ROI 命中
+- 防洪点体积变化阈值命中
+- 设备温度阈值命中
+- 豁免区与基线库过滤情况
+
+界面元素:
+
+- GIS 区域叠加规则缓冲区
+- 命中点位高亮闪烁
+- 规则命中列表展示每条规则的输入、阈值、结果
+
+规则命中展示示例:
+
+```text
+塔吊 / 高风险
+- 距线路:72.5m
+- 阈值:<= 100m
+- 置信度:0.91
+- 规则:线路100米范围内 + 置信度达标
+```
+
+### 5.8 告警生成
+
+展示内容:
+
+- 告警 ID、场景、类别、等级、置信度、状态
+- 告警证据:资源路径、AI 结果、规则命中
+- 告警分级统计:严重、高、中、低
+
+界面元素:
+
+- 告警卡片从 AI 结果区流入告警中心
+- GIS 告警点位按等级着色
+- 告警中心表格实时追加行
+
+### 5.9 工单派发
+
+展示内容:
+
+- 工单 ID、关联告警、处置人、派发状态、更新时间
+- 自动派发结果
+- 工单与告警的关联关系
+
+界面元素:
+
+- 工单流转泳道:
+
+```text
+告警生成 → 工单创建 → 派发复核 → 现场处置 → 复核关闭
+```
+
+- 工单列表实时增加
+- 待处置数量增加
+- 地图点位从“告警点”增加“工单状态标识”
+
+### 5.10 工单闭环
+
+展示内容:
+
+- 工单状态从 `created/processing` 变更为 `closed`
+- 告警状态同步关闭
+- 处置结果、处置说明、操作人、更新时间
+
+界面元素:
+
+- 工单卡片状态变更为“已闭环”
+- 统计卡片从待处置转入已闭环
+- 事件流追加 `workorder.closed`
+- 地图点位状态颜色变化
+
+### 5.11 样本回流
+
+展示内容:
+
+- 高价值样本数量
+- 待标注样本数量
+- 低置信度/误报/复核关闭样本归档
+- 模型版本影响范围
+
+界面元素:
+
+- 样本队列卡片
+- 模型版本旁展示“新增样本数”
+- 完成汇总中展示样本闭环指标
+
+### 5.12 完成汇总
+
+展示内容:
+
+- 总耗时
+- 任务、资源、分析任务、AI 结果、告警、工单、样本数量
+- 各阶段耗时
+- 成功/失败/重试次数
+
+界面元素:
+
+- 总览指标卡
+- 阶段耗时条
+- 最终 GIS 态势
+- 运行结果摘要
+
+## 6. 事件模型设计
+
+全链路可视化应以事件驱动。后端每完成一个阶段动作或子动作,就发布一条事件。前端只关心事件流和运行快照,不直接推断后端执行过程。
+
+### 6.1 运行实例
+
+建议新增 `demo_runs` 表:
+
+```sql
+create table demo_runs (
+ id varchar(64) primary key,
+ status varchar(32) not null,
+ current_step varchar(64),
+ progress int not null default 0,
+ task_id varchar(64),
+ summary jsonb not null default '{}'::jsonb,
+ started_at timestamptz not null,
+ completed_at timestamptz,
+ error_message text
+);
+```
+
+### 6.2 阶段事件
+
+建议新增 `demo_run_events` 表:
+
+```sql
+create table demo_run_events (
+ id varchar(64) primary key,
+ run_id varchar(64) not null references demo_runs(id),
+ step_key varchar(64) not null,
+ step_name varchar(128) not null,
+ status varchar(32) not null,
+ progress int not null,
+ title varchar(256) not null,
+ message text,
+ entity_type varchar(64),
+ entity_id varchar(64),
+ metrics jsonb not null default '{}'::jsonb,
+ visual_payload jsonb not null default '{}'::jsonb,
+ created_at timestamptz not null
+);
+```
+
+### 6.3 事件结构
+
+```json
+{
+ "event_id": "event-001",
+ "run_id": "run-001",
+ "step_key": "ai.inferencing",
+ "step_name": "AI推理分析",
+ "status": "running",
+ "progress": 58,
+ "title": "视觉模型推理完成",
+ "message": "vision-detector 生成 12 条识别结果",
+ "entity_type": "analysis_job",
+ "entity_id": "job-001",
+ "metrics": {
+ "result_count": 12,
+ "avg_confidence": 0.9,
+ "elapsed_ms": 1260
+ },
+ "visual_payload": {
+ "type": "inference_matrix",
+ "model_group": "vision-detector",
+ "status": "completed",
+ "items": [
+ {
+ "scene": "塔吊",
+ "category": "塔吊",
+ "confidence": 0.91,
+ "mileage": "K123+456"
+ }
+ ]
+ },
+ "created_at": "2026-07-21T16:20:00+08:00"
+}
+```
+
+## 7. 接口设计
+
+### 7.1 创建运行实例
+
+```http
+POST /api/v1/demo/runs
+```
+
+响应:
+
+```json
+{
+ "run_id": "run-xxx",
+ "status": "running",
+ "event_stream_url": "/api/v1/demo/runs/run-xxx/events"
+}
+```
+
+### 7.2 查询运行快照
+
+```http
+GET /api/v1/demo/runs/{run_id}
+```
+
+返回当前运行状态、阶段列表、最新统计、关联业务实体。
+
+### 7.3 订阅运行事件
+
+```http
+GET /api/v1/demo/runs/{run_id}/events
+Accept: text/event-stream
+```
+
+采用 SSE 的理由:
+
+- 当前场景主要是服务端向前端推送事件,SSE 足够稳定。
+- 浏览器原生支持 `EventSource`。
+- 相比 WebSocket,实现成本更低,便于后续通过 Nginx 代理。
+- 断线后可通过 `Last-Event-ID` 或快照接口恢复。
+
+### 7.4 查询历史事件
+
+```http
+GET /api/v1/demo/runs/{run_id}/events/history
+```
+
+用于页面刷新后的恢复和运行记录回放。
+
+### 7.5 取消运行
+
+```http
+POST /api/v1/demo/runs/{run_id}/cancel
+```
+
+用于长耗时阶段或异常情况下终止示例运行。
+
+## 8. 后端实现设计
+
+### 8.1 新增服务
+
+建议新增以下后端模块:
+
+```text
+platform/backend/src/main/java/com/ai/trackwalker/demo/
+ DemoRunController.java
+ DemoRunService.java
+ DemoRunEventPublisher.java
+ DemoRunEvent.java
+ DemoStep.java
+```
+
+职责:
+
+- `DemoRunController`:提供运行创建、查询、事件订阅、取消接口。
+- `DemoRunService`:负责编排全链路阶段执行。
+- `DemoRunEventPublisher`:负责事件入库、Kafka 发布、SSE 推送。
+- `DemoStep`:定义阶段编码、名称、权重、默认进度。
+
+### 8.2 编排逻辑
+
+运行实例创建后,后端按顺序执行阶段:
+
+```text
+createRun
+ -> emit initializing
+ -> createTask
+ -> dispatchUav
+ -> ingestResources
+ -> preprocessResources
+ -> createAnalysisJob
+ -> runAiInference
+ -> evaluateRules
+ -> generateAlarms
+ -> dispatchWorkorders
+ -> closeWorkorders
+ -> feedbackSamples
+ -> completeRun
+```
+
+每个阶段遵循统一模板:
+
+```text
+emit step.started
+执行业务动作
+emit step.progress
+写入阶段产物
+emit step.completed
+```
+
+异常处理:
+
+- 阶段失败时,写入 `step.failed` 事件。
+- `demo_runs.status` 更新为 `failed`。
+- 前端展示失败阶段、错误详情和重试入口。
+- 可恢复阶段支持从最近成功阶段继续执行。
+
+### 8.3 与现有服务的关系
+
+现有 `PlatformService.runDemo()` 可以保留为兼容接口,但新的可视化运行应调用 `DemoRunService`。
+
+建议逐步将现有同步逻辑拆分为可复用方法:
+
+- `createTask`
+- `completeResources`
+- `createAnalysisJob`
+- `runAnalysis`
+- `listAlarms`
+- `listWorkOrders`
+- `workOrderCallback`
+
+可视化运行编排只负责阶段组织和事件发布,不重复实现业务逻辑。
+
+## 9. 前端实现设计
+
+### 9.1 页面结构
+
+建议新增全链路运行组件:
+
+```text
+frontend/src/components/demo-run/
+ FullChainRunConsole.vue
+ RunHeader.vue
+ RunStageTimeline.vue
+ RunGisPanel.vue
+ RunCurrentStagePanel.vue
+ RunEventStream.vue
+ RunResourcePanel.vue
+ RunInferencePanel.vue
+ RunRulePanel.vue
+ RunAlarmPanel.vue
+ RunWorkorderPanel.vue
+ RunSummaryPanel.vue
+```
+
+### 9.2 状态管理
+
+建议新增组合式逻辑:
+
+```text
+frontend/src/composables/useDemoRun.ts
+```
+
+职责:
+
+- 创建运行实例。
+- 建立 SSE 连接。
+- 维护运行状态、阶段状态、事件列表。
+- 将事件转换为 GIS marker、资源卡片、AI 结果、工单列表。
+- 断线后通过快照接口恢复。
+
+### 9.3 阶段时间轴
+
+阶段状态样式:
+
+| 状态 | 样式 |
+| --- | --- |
+| 未开始 | 灰色节点 |
+| 执行中 | 蓝色节点 + 动态进度环 |
+| 成功 | 绿色节点 + 对勾 |
+| 失败 | 红色节点 + 错误标记 |
+| 跳过 | 浅灰节点 + 跳过标识 |
+| 重试中 | 黄色节点 + 旋转状态 |
+
+### 9.4 GIS 面板
+
+GIS 面板需要从事件中持续更新以下图层:
+
+- 线路中心线
+- 无人机航线
+- 防护区
+- 重点防洪点
+- 资源采集点
+- AI 识别点
+- 规则命中点
+- 告警点
+- 工单点
+
+图层状态:
+
+- 数据接入阶段:显示资源采集点。
+- AI 推理阶段:显示 AI 识别点。
+- GIS 规则阶段:显示规则缓冲区和命中点。
+- 告警阶段:按告警等级着色。
+- 工单阶段:点位增加工单状态标识。
+- 闭环阶段:已闭环点位颜色变为绿色或打勾。
+
+### 9.5 当前阶段详情
+
+右侧详情面板根据当前阶段动态展示:
+
+- 任务阶段:任务 ID、线路、里程、优先级。
+- 数据阶段:资源类型、数量、采集时间、存储路径。
+- 预处理阶段:步骤清单、质量指标。
+- AI 阶段:模型组、版本、结果数、置信度。
+- 规则阶段:规则命中原因、阈值、输入值。
+- 告警阶段:告警等级、证据链。
+- 工单阶段:派发人、状态、处置结果。
+- 样本阶段:样本数量、标注状态、模型版本。
+
+### 9.6 事件流
+
+事件流用于保留完整操作记录:
+
+```text
+[16:20:01] 运行初始化完成
+[16:20:02] 巡检任务 task-xxx 创建成功
+[16:20:04] 航线任务已下发
+[16:20:06] 接入可见光资源 visible-001.jpg
+[16:20:08] vision-detector 生成 12 条结果
+[16:20:09] 塔吊命中线路100米范围规则
+[16:20:10] 生成高风险告警 alarm-xxx
+[16:20:11] 工单 wo-xxx 已派发
+[16:20:14] 工单 wo-xxx 已闭环
+```
+
+支持筛选:
+
+- 全部
+- 当前阶段
+- 异常
+- 告警
+- 工单
+- 模型
+
+## 10. 前后端交互流程
+
+```mermaid
+sequenceDiagram
+ participant UI as 前端运行看板
+ participant API as Spring Boot API
+ participant Demo as DemoRunService
+ participant DB as PostgreSQL/PostGIS
+ participant AI as FastAPI AI服务
+ participant Kafka as Kafka事件总线
+
+ UI->>API: POST /api/v1/demo/runs
+ API->>DB: 创建 demo_runs
+ API-->>UI: 返回 run_id 和事件流地址
+ UI->>API: GET /api/v1/demo/runs/{run_id}/events
+ API->>Demo: 异步执行全链路阶段
+ Demo->>DB: 创建巡检任务
+ Demo->>Kafka: 发布 task.created
+ API-->>UI: SSE 推送 task.created
+ Demo->>DB: 写入资源
+ API-->>UI: SSE 推送 resource.ingesting
+ Demo->>AI: 调用视觉/点云推理
+ AI-->>Demo: 返回 AI 结果
+ Demo->>DB: 写入 AI 结果、告警、工单
+ API-->>UI: SSE 推送 AI、规则、告警、工单事件
+ Demo->>DB: 更新 demo_runs completed
+ API-->>UI: SSE 推送 run.completed
+```
+
+## 11. 运行快照结构
+
+前端刷新或断线恢复时,使用运行快照重建页面。
+
+```json
+{
+ "run_id": "run-xxx",
+ "status": "running",
+ "progress": 68,
+ "current_step": "rule.evaluating",
+ "started_at": "2026-07-21T16:20:00+08:00",
+ "steps": [
+ {
+ "key": "task.created",
+ "name": "巡检任务创建",
+ "status": "completed",
+ "progress": 100,
+ "elapsed_ms": 320
+ }
+ ],
+ "entities": {
+ "task_id": "task-xxx",
+ "resource_ids": ["res-001", "res-002"],
+ "analysis_job_id": "job-xxx",
+ "alarm_ids": ["alarm-001"],
+ "workorder_ids": ["wo-001"]
+ },
+ "metrics": {
+ "tasks": 1,
+ "resources": 3,
+ "ai_results": 18,
+ "alarms": 18,
+ "workorders": 18,
+ "closed_workorders": 0
+ }
+}
+```
+
+## 12. 视觉设计规范
+
+页面应保持巡检平台的业务工作台气质,避免营销式大图和装饰性内容。视觉重点放在信息密度、状态可读性和空间关系上。
+
+设计要求:
+
+- 颜色含义固定:蓝色表示执行中,绿色表示完成,橙色表示待处理,红色表示严重/失败,灰色表示未开始。
+- 同一业务实体在不同区域保持同一 ID 和同一颜色标识。
+- 地图点位数量较多时自动聚合,悬浮展示详情。
+- 阶段详情区域不滚动整页,内部列表独立滚动。
+- 重要事件进入事件流顶部,并短暂高亮。
+- 失败事件必须显式显示失败原因和所在阶段。
+- 全链路完成后,顶部运行态切换为完成态,不再显示加载动画。
+
+## 13. 测试方案
+
+### 13.1 后端测试
+
+单元测试:
+
+- 阶段定义完整性测试。
+- 事件发布顺序测试。
+- 事件 payload schema 测试。
+- 失败阶段状态更新测试。
+
+集成测试:
+
+- `POST /api/v1/demo/runs` 能创建运行实例。
+- SSE 能收到 `run.initializing` 到 `run.completed` 的完整事件。
+- 运行完成后任务、资源、AI 结果、告警、工单数量正确。
+- 断线后通过快照接口能恢复运行状态。
+
+### 13.2 前端测试
+
+组件测试:
+
+- 阶段时间轴状态渲染。
+- GIS marker 渲染。
+- 事件流追加和筛选。
+- 当前阶段详情切换。
+- 工单闭环状态更新。
+
+浏览器验证:
+
+- 点击“运行全链路示例”后,运行看板打开。
+- 每个阶段至少出现一次可视化变化。
+- 运行过程中按钮禁用并显示运行状态。
+- 运行完成后统计值与后端返回一致。
+- 失败事件能显示错误态。
+
+### 13.3 冒烟测试
+
+现有 `scripts/smoke_test.ps1` 需要扩展:
+
+- 创建 demo run。
+- 等待 SSE 或轮询快照直到 completed。
+- 校验阶段数量。
+- 校验关键阶段均有事件。
+- 校验最终任务、资源、AI 结果、告警、工单数据。
+
+## 14. 分阶段实施计划
+
+### 第一阶段:事件与运行实例
+
+- 新增 `demo_runs`、`demo_run_events` 表。
+- 新增 `DemoRunService` 和事件发布器。
+- 新增运行创建、快照查询、历史事件接口。
+- 保留原同步接口兼容现有调用。
+
+交付结果:
+
+- 后端能创建运行实例并记录阶段事件。
+- 不要求前端实时展示,先可通过接口查看事件。
+
+### 第二阶段:SSE 实时推送
+
+- 新增 `/api/v1/demo/runs/{run_id}/events`。
+- 后端阶段执行时实时推送事件。
+- 支持运行结束、失败、取消事件。
+- 支持断线恢复。
+
+交付结果:
+
+- 前端可实时收到阶段事件。
+- 事件流不依赖手动刷新。
+
+### 第三阶段:前端运行看板
+
+- 新增全链路运行看板组件。
+- 接入运行创建和 SSE。
+- 实现阶段时间轴、运行状态、事件流。
+- 原“运行全链路示例”按钮改为打开运行看板并启动运行。
+
+交付结果:
+
+- 每个阶段有可见状态。
+- 能看到事件逐步推进。
+
+### 第四阶段:阶段专属可视化
+
+- GIS 图层动态更新。
+- 数据资源卡片动态追加。
+- AI 推理矩阵动态更新。
+- 规则命中面板动态展示。
+- 告警与工单流转动态展示。
+- 样本回流和完成汇总展示。
+
+交付结果:
+
+- 每个业务环节都有独立可视化表现。
+- 执行过程具备演进感和可解释性。
+
+### 第五阶段:测试与稳定性
+
+- 扩展后端测试和冒烟脚本。
+- 增加前端浏览器验证。
+- 增加失败、重试、取消、断线恢复测试。
+- 优化大数据量下的事件裁剪和地图点聚合。
+
+交付结果:
+
+- 全链路可视化稳定可用。
+- 页面刷新和异常场景可恢复。
+
+## 15. 验收标准
+
+功能验收:
+
+- 点击“运行全链路示例”后,能看到运行实例编号和总进度。
+- 12 个阶段全部在时间轴中展示。
+- 每个阶段执行时,页面都有对应的可视化变化。
+- 执行过程中事件流持续追加事件。
+- GIS 区域能展示航线、资源点、AI 识别点、告警点和工单状态。
+- AI 推理阶段能展示模型组、结果数量、置信度。
+- GIS 规则阶段能展示规则命中原因。
+- 告警阶段能展示告警等级和证据。
+- 工单阶段能展示派发、处置、闭环状态。
+- 运行完成后展示汇总指标和阶段耗时。
+
+稳定性验收:
+
+- 后端异常时前端显示失败阶段和失败原因。
+- 页面刷新后能恢复当前运行状态。
+- SSE 断开后能自动重连或通过快照恢复。
+- 重复点击运行不会产生不可控的并发状态。
+- 运行历史可查询,便于复盘。
+
+## 16. 后续扩展
+
+- 接入真实无人机航迹点,替换当前示意航线。
+- 接入真实影像缩略图、红外热力图、点云剖面图。
+- 接入 MapLibre 地图底图,实现真实地理坐标渲染。
+- 接入 WebSocket,用于双向控制无人机任务、暂停和重试阶段。
+- 接入 Grafana 或 Prometheus 指标,展示阶段耗时和服务延迟。
+- 接入模型评估指标,展示样本回流对模型版本的影响。
diff --git a/frontend/nginx.conf b/frontend/nginx.conf
index d50c2db..5f03205 100644
--- a/frontend/nginx.conf
+++ b/frontend/nginx.conf
@@ -13,6 +13,9 @@ server {
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
+ proxy_buffering off;
+ proxy_cache off;
+ proxy_read_timeout 3600s;
}
location /actuator/ {
diff --git a/frontend/src/App.vue b/frontend/src/App.vue
index b219f20..5e45da9 100644
--- a/frontend/src/App.vue
+++ b/frontend/src/App.vue
@@ -5,8 +5,12 @@
铁路无人机智能巡检平台
任务接入、AI识别、GIS规则、告警工单、样本闭环
- 运行全链路示例
+
+
+
+
+
+
+ 全链路运行看板
+ {{ activeRunId || "等待运行" }} · {{ currentRunTitle }}
+
+
+
+
+
+
+
+ {{ stageIndex(stage.key) }}
+ {{ stage.name }}
+
+
+
+
+
+
+
+
+ {{ routeProgress }}%
+
+
+
+
{{ marker.index }}
+
+ {{ marker.scene }}
+ {{ marker.mileage }} / {{ marker.distance || "线路邻近" }}
+ {{ severityLabel(marker.severity) }} {{ statusLabel(marker.status) }}
+
+
+
+
+ 线路中心线
+ 防护区
+ 重点防洪点
+ 阶段点位
+
+
+
+
+
+
+
+
+
事件流
+
+
+ {{ event.step_name }}
+ {{ event.title }}
+
+
+
+
+
+ {{ item.value }}
+ {{ item.label }}
+
+
+
+
+
+
巡检任务
@@ -155,9 +305,20 @@
diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts
index 0885a79..7242c12 100644
--- a/frontend/src/services/api.ts
+++ b/frontend/src/services/api.ts
@@ -45,3 +45,22 @@ export async function runDemo() {
const { data } = await client.post("/demo/run-full-chain", {});
return data.data;
}
+
+export async function startDemoRun() {
+ const { data } = await client.post("/demo/runs", {});
+ return data.data;
+}
+
+export async function demoRunSnapshot(runId: string) {
+ const { data } = await client.get(`/demo/runs/${runId}`);
+ return data.data;
+}
+
+export async function cancelDemoRun(runId: string) {
+ const { data } = await client.post(`/demo/runs/${runId}/cancel`, {});
+ return data.data;
+}
+
+export function demoEventsUrl(runId: string) {
+ return `/api/v1/demo/runs/${runId}/events`;
+}
diff --git a/frontend/src/styles.css b/frontend/src/styles.css
index b1ae79f..34a681c 100644
--- a/frontend/src/styles.css
+++ b/frontend/src/styles.css
@@ -30,6 +30,13 @@ body {
color: #dce8f5;
}
+.header-actions {
+ display: inline-flex;
+ align-items: center;
+ gap: 10px;
+ flex-wrap: wrap;
+}
+
.stats {
display: grid;
grid-template-columns: repeat(auto-fit, minmax(130px, 1fr));
@@ -55,6 +62,330 @@ body {
margin-bottom: 18px;
}
+.run-console {
+ margin-bottom: 18px;
+}
+
+.run-header {
+ display: grid;
+ grid-template-columns: minmax(260px, .55fr) minmax(320px, 1fr);
+ gap: 18px;
+ align-items: center;
+ padding-bottom: 16px;
+ border-bottom: 1px solid #e5edf5;
+}
+
+.run-header strong,
+.run-header span {
+ display: block;
+}
+
+.run-header strong {
+ color: #17365d;
+ font-size: 22px;
+}
+
+.run-header span {
+ margin-top: 6px;
+ color: #64748b;
+ font-size: 13px;
+}
+
+.run-timeline {
+ display: grid;
+ grid-template-columns: repeat(auto-fit, minmax(112px, 1fr));
+ gap: 10px;
+ margin: 16px 0;
+}
+
+.run-stage {
+ min-height: 68px;
+ padding: 10px;
+ border: 1px solid #d8e3ee;
+ border-radius: 8px;
+ background: #f8fafc;
+ color: #64748b;
+ display: grid;
+ grid-template-columns: 30px minmax(0, 1fr);
+ gap: 8px;
+ align-items: center;
+}
+
+.run-stage i {
+ width: 28px;
+ height: 28px;
+ border-radius: 999px;
+ display: grid;
+ place-items: center;
+ background: #e2e8f0;
+ color: #475569;
+ font-size: 12px;
+ font-style: normal;
+ font-weight: 800;
+}
+
+.run-stage span {
+ min-width: 0;
+ font-size: 13px;
+ font-weight: 700;
+ line-height: 1.25;
+}
+
+.stage-running {
+ border-color: rgba(37, 99, 235, .5);
+ background: #eff6ff;
+ color: #1d4ed8;
+}
+
+.stage-running i {
+ background: #2563eb;
+ color: #fff;
+}
+
+.stage-completed {
+ border-color: rgba(22, 163, 74, .42);
+ background: #f0fdf4;
+ color: #166534;
+}
+
+.stage-completed i {
+ background: #16a34a;
+ color: #fff;
+}
+
+.stage-failed {
+ border-color: rgba(220, 38, 38, .45);
+ background: #fef2f2;
+ color: #b91c1c;
+}
+
+.stage-failed i {
+ background: #dc2626;
+ color: #fff;
+}
+
+.stage-cancelled {
+ border-color: rgba(245, 158, 11, .45);
+ background: #fffbeb;
+ color: #92400e;
+}
+
+.stage-cancelled i {
+ background: #f59e0b;
+ color: #fff;
+}
+
+.run-board {
+ display: grid;
+ grid-template-columns: minmax(0, 1.15fr) minmax(360px, .85fr);
+ gap: 18px;
+ align-items: stretch;
+}
+
+.run-map {
+ position: relative;
+ min-height: 430px;
+ overflow: hidden;
+ border: 1px solid #d8e3ee;
+ border-radius: 8px;
+ background:
+ radial-gradient(circle at 18% 72%, rgba(34, 197, 94, .12), transparent 24%),
+ radial-gradient(circle at 78% 32%, rgba(14, 165, 233, .14), transparent 22%),
+ linear-gradient(90deg, rgba(203, 213, 225, .22) 1px, transparent 1px),
+ linear-gradient(0deg, rgba(203, 213, 225, .22) 1px, transparent 1px),
+ linear-gradient(135deg, #f8fbff 0%, #edf7f1 100%);
+ background-size: 100% 100%, 100% 100%, 64px 64px, 64px 64px, 100% 100%;
+}
+
+.stage-detail {
+ min-width: 0;
+ border: 1px solid #d8e3ee;
+ border-radius: 8px;
+ background: #fff;
+ padding: 14px;
+ overflow: hidden;
+}
+
+.stage-title {
+ display: flex;
+ justify-content: space-between;
+ align-items: center;
+ gap: 12px;
+ margin-bottom: 12px;
+}
+
+.stage-title strong {
+ color: #17365d;
+ font-size: 18px;
+}
+
+.metric-grid {
+ display: grid;
+ grid-template-columns: repeat(auto-fit, minmax(120px, 1fr));
+ gap: 10px;
+ margin-bottom: 14px;
+}
+
+.metric-grid div {
+ min-height: 64px;
+ padding: 10px;
+ border: 1px solid #e2e8f0;
+ border-radius: 8px;
+ background: #f8fafc;
+}
+
+.metric-grid strong,
+.metric-grid span {
+ display: block;
+}
+
+.metric-grid strong {
+ color: #17365d;
+ font-size: 18px;
+ word-break: break-word;
+}
+
+.metric-grid span {
+ margin-top: 6px;
+ color: #64748b;
+ font-size: 12px;
+}
+
+.visual-block {
+ margin-top: 12px;
+}
+
+.visual-block h3 {
+ margin: 0 0 8px;
+ color: #334155;
+ font-size: 14px;
+}
+
+.visual-row {
+ display: grid;
+ grid-template-columns: minmax(110px, .7fr) minmax(0, 1fr);
+ gap: 10px;
+ min-height: 34px;
+ padding: 8px 0;
+ border-bottom: 1px solid #eef2f7;
+}
+
+.visual-row span,
+.visual-row em {
+ min-width: 0;
+}
+
+.visual-row span {
+ color: #17365d;
+ font-weight: 700;
+}
+
+.visual-row em {
+ color: #64748b;
+ font-style: normal;
+ word-break: break-word;
+}
+
+.run-bottom {
+ display: grid;
+ grid-template-columns: minmax(0, 1.1fr) minmax(340px, .9fr);
+ gap: 18px;
+ margin-top: 18px;
+}
+
+.event-stream,
+.run-summary {
+ border: 1px solid #d8e3ee;
+ border-radius: 8px;
+ background: #fff;
+}
+
+.event-stream {
+ max-height: 300px;
+ overflow: auto;
+}
+
+.stream-title {
+ position: sticky;
+ top: 0;
+ z-index: 1;
+ padding: 12px 14px;
+ border-bottom: 1px solid #e5edf5;
+ background: #fff;
+ color: #17365d;
+ font-weight: 800;
+}
+
+.event-line {
+ display: grid;
+ grid-template-columns: 80px minmax(110px, .45fr) minmax(0, 1fr);
+ gap: 12px;
+ align-items: center;
+ padding: 10px 14px;
+ border-bottom: 1px solid #eef2f7;
+ color: #334155;
+}
+
+.event-line time {
+ color: #64748b;
+ font-size: 12px;
+}
+
+.event-line strong {
+ color: #17365d;
+ font-size: 13px;
+}
+
+.event-line span {
+ min-width: 0;
+ color: #475569;
+ font-size: 13px;
+}
+
+.event-running {
+ background: #eff6ff;
+}
+
+.event-completed {
+ background: #f7fef9;
+}
+
+.event-failed {
+ background: #fef2f2;
+}
+
+.run-summary {
+ display: grid;
+ grid-template-columns: repeat(auto-fit, minmax(130px, 1fr));
+ gap: 0;
+ align-content: start;
+ overflow: hidden;
+}
+
+.run-summary div {
+ min-height: 78px;
+ padding: 14px;
+ border-right: 1px solid #eef2f7;
+ border-bottom: 1px solid #eef2f7;
+}
+
+.run-summary strong,
+.run-summary span {
+ display: block;
+}
+
+.run-summary strong {
+ color: #17365d;
+ font-size: 20px;
+ word-break: break-word;
+}
+
+.run-summary span {
+ margin-top: 7px;
+ color: #64748b;
+ font-size: 12px;
+}
+
.gis-grid {
grid-template-columns: minmax(0, 1.35fr) minmax(min(100%, 420px), .65fr);
}
@@ -162,7 +493,8 @@ body {
font-weight: 600;
}
-.alarm-marker {
+.alarm-marker,
+.run-marker {
position: absolute;
z-index: 3;
transform: translate(-50%, -50%);
@@ -179,7 +511,8 @@ body {
cursor: default;
}
-.alarm-marker::after {
+.alarm-marker::after,
+.run-marker::after {
content: "";
position: absolute;
inset: -7px;
@@ -188,6 +521,27 @@ body {
opacity: .25;
}
+.run-marker {
+ z-index: 5;
+}
+
+.uav-marker {
+ position: absolute;
+ z-index: 6;
+ width: 58px;
+ height: 30px;
+ transform: translate(-50%, -50%);
+ border: 2px solid #fff;
+ border-radius: 999px;
+ background: #2563eb;
+ color: #fff;
+ display: grid;
+ place-items: center;
+ font-size: 12px;
+ font-weight: 800;
+ box-shadow: 0 12px 30px rgba(37, 99, 235, .3);
+}
+
.severity-critical {
background: #dc2626;
}
@@ -204,6 +558,10 @@ body {
background: #16a34a;
}
+.severity-resource {
+ background: #0f766e;
+}
+
.marker-card {
position: absolute;
left: 34px;
@@ -221,7 +579,8 @@ body {
transition: opacity .16s ease, transform .16s ease;
}
-.alarm-marker:hover .marker-card {
+.alarm-marker:hover .marker-card,
+.run-marker:hover .marker-card {
opacity: 1;
transform: translate(4px, -50%);
}
@@ -376,6 +735,9 @@ body {
gap: 14px;
}
+ .run-header,
+ .run-board,
+ .run-bottom,
.gis-grid,
.gis-workbench {
grid-template-columns: 1fr;
@@ -396,4 +758,9 @@ body {
align-items: flex-start;
flex-direction: column;
}
+
+ .event-line,
+ .visual-row {
+ grid-template-columns: 1fr;
+ }
}
diff --git a/platform/backend/src/main/java/com/ai/trackwalker/demo/DemoRunController.java b/platform/backend/src/main/java/com/ai/trackwalker/demo/DemoRunController.java
new file mode 100644
index 0000000..81ec7e7
--- /dev/null
+++ b/platform/backend/src/main/java/com/ai/trackwalker/demo/DemoRunController.java
@@ -0,0 +1,42 @@
+package com.ai.trackwalker.demo;
+
+import com.ai.trackwalker.api.ApiResponse;
+import org.springframework.web.bind.annotation.*;
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+import java.util.Map;
+
+@RestController
+@RequestMapping("/api/v1/demo/runs")
+public class DemoRunController {
+ private final DemoRunService service;
+
+ public DemoRunController(DemoRunService service) {
+ this.service = service;
+ }
+
+ @PostMapping
+ public ApiResponse> startRun() {
+ return ApiResponse.ok(service.startRun());
+ }
+
+ @GetMapping("/{runId}")
+ public ApiResponse> snapshot(@PathVariable String runId) {
+ return ApiResponse.ok(service.snapshot(runId));
+ }
+
+ @GetMapping("/{runId}/events")
+ public SseEmitter events(@PathVariable String runId) {
+ return service.subscribe(runId);
+ }
+
+ @GetMapping("/{runId}/events/history")
+ public ApiResponse> history(@PathVariable String runId) {
+ return ApiResponse.ok(Map.of("events", service.history(runId)));
+ }
+
+ @PostMapping("/{runId}/cancel")
+ public ApiResponse> cancel(@PathVariable String runId) {
+ return ApiResponse.ok(service.cancel(runId));
+ }
+}
diff --git a/platform/backend/src/main/java/com/ai/trackwalker/demo/DemoRunEventPublisher.java b/platform/backend/src/main/java/com/ai/trackwalker/demo/DemoRunEventPublisher.java
new file mode 100644
index 0000000..32dff20
--- /dev/null
+++ b/platform/backend/src/main/java/com/ai/trackwalker/demo/DemoRunEventPublisher.java
@@ -0,0 +1,179 @@
+package com.ai.trackwalker.demo;
+
+import com.ai.trackwalker.common.Ids;
+import com.ai.trackwalker.common.Jsonb;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.stereotype.Component;
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+import java.io.IOException;
+import java.sql.Timestamp;
+import java.time.Instant;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CopyOnWriteArrayList;
+
+@Component
+public class DemoRunEventPublisher {
+ private static final long SSE_TIMEOUT_MS = 30 * 60 * 1000L;
+
+ private final JdbcTemplate jdbc;
+ private final Map> emitters = new ConcurrentHashMap<>();
+
+ public DemoRunEventPublisher(JdbcTemplate jdbc) {
+ this.jdbc = jdbc;
+ }
+
+ public SseEmitter subscribe(String runId) {
+ SseEmitter emitter = new SseEmitter(SSE_TIMEOUT_MS);
+ emitters.computeIfAbsent(runId, ignored -> new CopyOnWriteArrayList<>()).add(emitter);
+ emitter.onCompletion(() -> removeEmitter(runId, emitter));
+ emitter.onTimeout(() -> removeEmitter(runId, emitter));
+ emitter.onError(error -> removeEmitter(runId, emitter));
+
+ for (Map event : history(runId)) {
+ send(emitter, event);
+ }
+ return emitter;
+ }
+
+ public Map emit(
+ String runId,
+ DemoStep step,
+ String status,
+ int progress,
+ String title,
+ String message,
+ String entityType,
+ String entityId,
+ Map metrics,
+ Map visualPayload
+ ) {
+ String eventId = Ids.next("dre");
+ Instant now = Instant.now();
+ Map safeMetrics = metrics == null ? Map.of() : metrics;
+ Map safeVisualPayload = visualPayload == null ? Map.of() : visualPayload;
+
+ jdbc.update(
+ "insert into demo_run_events(id, run_id, step_key, step_name, status, progress, title, message, entity_type, entity_id, metrics, visual_payload, created_at) values (?,?,?,?,?,?,?,?,?,?,?::jsonb,?::jsonb,?)",
+ eventId,
+ runId,
+ step.key(),
+ step.label(),
+ status,
+ progress,
+ title,
+ message,
+ entityType,
+ entityId,
+ Jsonb.write(safeMetrics),
+ Jsonb.write(safeVisualPayload),
+ Timestamp.from(now)
+ );
+
+ String runStatus = runStatus(step, status);
+ if ("failed".equals(runStatus)) {
+ jdbc.update("update demo_runs set status=?, current_step=?, progress=?, completed_at=?, error_message=? where id=?",
+ runStatus, step.key(), progress, Timestamp.from(now), message, runId);
+ } else if ("completed".equals(runStatus) || "cancelled".equals(runStatus)) {
+ jdbc.update("update demo_runs set status=?, current_step=?, progress=?, completed_at=? where id=?",
+ runStatus, step.key(), progress, Timestamp.from(now), runId);
+ } else {
+ jdbc.update("update demo_runs set status=?, current_step=?, progress=? where id=?", runStatus, step.key(), progress, runId);
+ }
+
+ Map event = Map.ofEntries(
+ Map.entry("event_id", eventId),
+ Map.entry("run_id", runId),
+ Map.entry("step_key", step.key()),
+ Map.entry("step_name", step.label()),
+ Map.entry("status", status),
+ Map.entry("progress", progress),
+ Map.entry("title", title),
+ Map.entry("message", message == null ? "" : message),
+ Map.entry("entity_type", entityType == null ? "" : entityType),
+ Map.entry("entity_id", entityId == null ? "" : entityId),
+ Map.entry("metrics", safeMetrics),
+ Map.entry("visual_payload", safeVisualPayload),
+ Map.entry("created_at", now.toString())
+ );
+ broadcast(runId, event);
+ return event;
+ }
+
+ public List