全链路节点可视化

This commit is contained in:
2026-07-21 17:00:43 +08:00
parent 221e77d9da
commit bb541bb406
14 changed files with 2528 additions and 25 deletions
+11 -4
View File
@@ -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 返回。
+29 -2
View File
@@ -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 源。需要以下基础镜像可用:
+826
View File
@@ -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 指标,展示阶段耗时和服务延迟。
- 接入模型评估指标,展示样本回流对模型版本的影响。
+3
View File
@@ -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/ {
+458 -10
View File
@@ -5,8 +5,12 @@
<h1>铁路无人机智能巡检平台</h1>
<p>任务接入AI识别GIS规则告警工单样本闭环</p>
</div>
<el-button type="primary" :loading="running" @click="handleRunDemo">运行全链路示例</el-button>
<div class="header-actions">
<el-button v-if="running && activeRunId" plain @click="handleCancelRun">取消运行</el-button>
<el-button type="primary" :loading="running" @click="handleRunDemo">运行全链路示例</el-button>
</div>
</el-header>
<el-main>
<section class="stats">
<el-card v-for="item in statCards" :key="item.label" shadow="never">
@@ -15,6 +19,152 @@
</el-card>
</section>
<section v-if="runConsoleVisible" class="run-console">
<el-card shadow="never">
<template #header>
<div class="card-title">
<span>全链路运行看板</span>
<small>{{ activeRunId || "等待运行" }} · {{ currentRunTitle }}</small>
</div>
</template>
<div class="run-header">
<div>
<strong>{{ runStatusLabel }}</strong>
<span>{{ currentRunMessage }}</span>
</div>
<el-progress :percentage="runProgress" :stroke-width="12" :status="runProgressStatus" />
</div>
<div class="run-timeline">
<div
v-for="stage in demoStages"
:key="stage.key"
class="run-stage"
:class="`stage-${stageStatus(stage.key)}`"
>
<i>{{ stageIndex(stage.key) }}</i>
<span>{{ stage.name }}</span>
</div>
</div>
<div class="run-board">
<div class="run-map">
<svg class="track-layer" viewBox="0 0 1000 420" role="img" aria-label="全链路运行空间态势">
<path class="corridor corridor-outer" d="M80 305 C245 225 375 245 515 180 C665 112 785 126 925 68" />
<path class="corridor corridor-inner" d="M110 328 C275 248 395 268 535 203 C685 135 805 149 945 91" />
<path class="rail-shadow" d="M78 278 C250 198 380 218 520 153 C670 85 795 96 934 42" />
<path class="rail-line" d="M78 278 C250 198 380 218 520 153 C670 85 795 96 934 42" />
<path class="wire-line" d="M88 235 C258 165 395 186 536 122 C682 57 808 68 938 18" />
<polygon class="flood-zone" points="590,176 690,128 794,143 820,232 700,268 616,238" />
<polygon class="protected-zone" points="205,250 390,214 434,284 253,326" />
<line class="grid-line" x1="92" y1="340" x2="930" y2="74" />
<text class="mileage-label" x="76" y="366">K123+000</text>
<text class="mileage-label" x="850" y="108">K130+000</text>
</svg>
<div v-if="routeProgress > 0" class="uav-marker" :style="{ left: `${uavPosition.x}%`, top: `${uavPosition.y}%` }">
<span>{{ routeProgress }}%</span>
</div>
<div
v-for="marker in runMarkers"
:key="marker.id"
class="run-marker"
:class="`severity-${marker.severity}`"
:style="{ left: `${marker.x}%`, top: `${marker.y}%` }"
>
<span>{{ marker.index }}</span>
<div class="marker-card">
<strong>{{ marker.scene }}</strong>
<em>{{ marker.mileage }} / {{ marker.distance || "线路邻近" }}</em>
<small>{{ severityLabel(marker.severity) }} {{ statusLabel(marker.status) }}</small>
</div>
</div>
<div class="map-legend">
<span><i class="legend-rail"></i>线路中心线</span>
<span><i class="legend-protected"></i>防护区</span>
<span><i class="legend-flood"></i>重点防洪点</span>
<span><i class="legend-alarm"></i>阶段点位</span>
</div>
</div>
<aside class="stage-detail">
<div class="stage-title">
<strong>{{ currentStageName }}</strong>
<el-tag :type="stageTag(currentEvent?.status)" effect="plain">{{ eventStatusLabel(currentEvent?.status) }}</el-tag>
</div>
<div class="metric-grid">
<div v-for="metric in currentMetrics" :key="metric.label">
<strong>{{ metric.value }}</strong>
<span>{{ metric.label }}</span>
</div>
</div>
<div v-if="currentVisual.models?.length" class="visual-block">
<h3>模型推理矩阵</h3>
<div v-for="model in currentVisual.models" :key="model.model_group" class="visual-row">
<span>{{ model.model_group }}</span>
<em>{{ model.status }} · {{ model.result_count }} · {{ model.avg_confidence }}</em>
</div>
</div>
<div v-if="currentVisual.checks?.length" class="visual-block">
<h3>预处理检查</h3>
<div v-for="check in currentVisual.checks" :key="check.name" class="visual-row">
<span>{{ check.name }}</span>
<em>{{ check.status }} · {{ check.quality }}</em>
</div>
</div>
<div v-if="currentVisual.rule_hits?.length" class="visual-block">
<h3>规则命中</h3>
<div v-for="hit in currentVisual.rule_hits" :key="hit.alarm_id" class="visual-row">
<span>{{ hit.scene }}</span>
<em>{{ formatList(hit.rule_hits) }}</em>
</div>
</div>
<div v-if="currentVisual.workorders?.length" class="visual-block">
<h3>工单流转</h3>
<div v-for="wo in currentVisual.workorders" :key="wo.workorder_id" class="visual-row">
<span>{{ wo.workorder_id }}</span>
<em>{{ wo.scene }} · {{ statusLabel(wo.status) }}</em>
</div>
</div>
<div v-if="currentVisual.summary" class="visual-block">
<h3>完成汇总</h3>
<div v-for="item in objectEntries(currentVisual.summary)" :key="item.label" class="visual-row">
<span>{{ item.label }}</span>
<em>{{ item.value }}</em>
</div>
</div>
</aside>
</div>
<div class="run-bottom">
<div class="event-stream">
<div class="stream-title">事件流</div>
<div v-for="event in recentRunEvents" :key="event.event_id" class="event-line" :class="`event-${event.status}`">
<time>{{ formatTime(event.created_at) }}</time>
<strong>{{ event.step_name }}</strong>
<span>{{ event.title }}</span>
</div>
</div>
<div class="run-summary">
<div v-for="item in runSummaryItems" :key="item.label">
<strong>{{ item.value }}</strong>
<span>{{ item.label }}</span>
</div>
</div>
</div>
</el-card>
</section>
<section class="grid">
<el-card shadow="never">
<template #header>巡检任务</template>
@@ -155,9 +305,20 @@
</template>
<script setup lang="ts">
import { computed, onMounted, ref } from "vue";
import { computed, onBeforeUnmount, onMounted, ref } from "vue";
import { ElMessage } from "element-plus";
import { alarms, closeWorkorder, models, overview, runDemo, tasks, workorders } from "./services/api";
import {
alarms,
cancelDemoRun,
closeWorkorder,
demoEventsUrl,
demoRunSnapshot,
models,
overview,
startDemoRun,
tasks,
workorders
} from "./services/api";
type Row = Record<string, any>;
@@ -179,6 +340,12 @@ const alarmRows = ref<Row[]>([]);
const workorderRows = ref<Row[]>([]);
const modelRows = ref<Row[]>([]);
const running = ref(false);
const runConsoleVisible = ref(false);
const activeRunId = ref("");
const demoEvents = ref<Row[]>([]);
const demoRun = ref<Row | null>(null);
let demoEventSource: EventSource | null = null;
const labels: Record<string, string> = {
tasks: "巡检任务",
@@ -192,9 +359,103 @@ const labels: Record<string, string> = {
exemptions: "豁免区"
};
const demoStages = [
{ key: "run.initializing", name: "运行初始化" },
{ key: "task.created", name: "巡检任务创建" },
{ key: "uav.dispatching", name: "航线任务下发" },
{ key: "resource.ingesting", name: "多源数据接入" },
{ key: "preprocess.running", name: "数据预处理" },
{ key: "ai.inferencing", name: "AI推理分析" },
{ key: "rule.evaluating", name: "GIS规则判定" },
{ key: "alarm.generating", name: "告警生成" },
{ key: "workorder.dispatching", name: "工单派发" },
{ key: "workorder.closing", name: "工单闭环" },
{ key: "sample.feedback", name: "样本回流" },
{ key: "run.completed", name: "完成汇总" }
];
const statCards = computed(() => Object.entries(labels).map(([key, label]) => ({ label, value: stats.value[key] ?? 0 })));
const recentWorkorders = computed(() => workorderRows.value.slice(0, 8));
const currentEvent = computed(() => demoEvents.value[demoEvents.value.length - 1] ?? null);
const currentVisual = computed<Row>(() => safeObject(currentEvent.value?.visual_payload));
const recentRunEvents = computed(() => [...demoEvents.value].reverse().slice(0, 14));
const runProgress = computed(() => Number(currentEvent.value?.progress ?? demoRun.value?.progress ?? 0));
const currentRunTitle = computed(() => currentEvent.value?.title ?? "等待启动");
const currentRunMessage = computed(() => currentEvent.value?.message ?? "点击运行后将实时展示任务、数据、AI、规则、告警和工单闭环过程");
const currentStageName = computed(() => currentEvent.value?.step_name ?? "全链路阶段详情");
const runStatusLabel = computed(() => {
const status = currentEvent.value?.status ?? demoRun.value?.status ?? "pending";
if (currentEvent.value?.step_key === "run.completed") {
return "运行完成";
}
return eventStatusLabel(status);
});
const runProgressStatus = computed<"" | "success" | "exception" | "warning">(() => {
const status = currentEvent.value?.status;
if (status === "failed") {
return "exception";
}
if (status === "cancelled") {
return "warning";
}
if (currentEvent.value?.step_key === "run.completed") {
return "success";
}
return "";
});
const currentMetrics = computed(() => objectEntries(safeObject(currentEvent.value?.metrics)).slice(0, 8));
const runSummaryItems = computed(() => {
const summary = safeObject(currentVisual.value.summary ?? demoRun.value?.summary);
const entries = objectEntries(summary);
if (entries.length > 0) {
return entries.slice(0, 8);
}
return [
{ label: "阶段事件", value: demoEvents.value.length },
{ label: "总进度", value: `${runProgress.value}%` },
{ label: "当前阶段", value: currentStageName.value }
];
});
const routeProgress = computed(() => {
const payload = latestPayloadWith("route_progress");
return Number(payload?.route_progress ?? 0);
});
const uavPosition = computed(() => {
const progress = routeProgress.value / 100;
return {
x: clamp(8 + progress * 84, 5, 94),
y: clamp(70 - progress * 54, 8, 88)
};
});
const runMarkers = computed<GisMarker[]>(() => {
const markerPayload = latestPayloadWith("markers");
if (Array.isArray(markerPayload?.markers) && markerPayload.markers.length > 0) {
return markerPayload.markers.slice(0, 18).map((marker: Row, index: number) => normalizeMarker(marker, index));
}
const resourcePayload = latestPayloadWith("resources");
if (Array.isArray(resourcePayload?.resources) && resourcePayload.resources.length > 0) {
return resourcePayload.resources.map((resource: Row, index: number) => ({
id: String(resource.resource_id || `${resource.resource_type}-${index}`),
index: index + 1,
scene: resourceLabel(resource.resource_type),
mileage: "采集点",
distance: resource.status,
severity: "resource",
status: resource.status,
x: Number(resource.x ?? 20 + index * 20),
y: Number(resource.y ?? 60 - index * 8)
}));
}
return gisMarkers.value;
});
const activeLineLabel = computed(() => {
const firstTask = taskRows.value[0];
@@ -262,22 +523,131 @@ async function refresh() {
}
async function handleRunDemo() {
closeDemoEventSource();
running.value = true;
runConsoleVisible.value = true;
demoEvents.value = [];
demoRun.value = null;
try {
await runDemo();
await refresh();
ElMessage.success("全链路示例已完成");
} finally {
const run = await startDemoRun();
activeRunId.value = run.run_id;
demoRun.value = run;
connectDemoEvents(run.run_id);
ElMessage.success("全链路示例已启动");
} catch (error) {
running.value = false;
ElMessage.error("全链路示例启动失败");
}
}
async function handleCancelRun() {
if (!activeRunId.value) {
return;
}
await cancelDemoRun(activeRunId.value);
ElMessage.warning("已提交取消请求");
}
async function handleCloseWorkOrder(row: Row) {
await closeWorkorder(String(row.alarm_id));
await refresh();
ElMessage.success("工单已闭环");
}
function connectDemoEvents(runId: string) {
demoEventSource = new EventSource(demoEventsUrl(runId));
demoEventSource.addEventListener("demo-run-event", async (event: Event) => {
const message = event as MessageEvent;
const data = safeObject(message.data);
upsertEvent(data);
const terminal = ["run.completed", "run.failed", "run.cancelled"].includes(String(data.step_key));
if (terminal) {
running.value = false;
closeDemoEventSource();
demoRun.value = await demoRunSnapshot(runId);
await refresh();
if (data.step_key === "run.completed") {
ElMessage.success("全链路运行完成");
}
if (data.step_key === "run.failed") {
ElMessage.error("全链路运行失败");
}
}
});
demoEventSource.onerror = async () => {
if (!activeRunId.value || !running.value) {
return;
}
try {
const snapshot = await demoRunSnapshot(activeRunId.value);
demoRun.value = snapshot;
const status = String(snapshot.status);
if (["completed", "failed", "cancelled"].includes(status)) {
running.value = false;
closeDemoEventSource();
await refresh();
}
} catch {
// EventSource will retry automatically.
}
};
}
function closeDemoEventSource() {
if (demoEventSource) {
demoEventSource.close();
demoEventSource = null;
}
}
function upsertEvent(event: Row) {
if (!event.event_id) {
return;
}
const index = demoEvents.value.findIndex((item) => item.event_id === event.event_id);
if (index >= 0) {
demoEvents.value[index] = event;
} else {
demoEvents.value.push(event);
}
}
function latestPayloadWith(field: string) {
for (let i = demoEvents.value.length - 1; i >= 0; i--) {
const payload = safeObject(demoEvents.value[i].visual_payload);
if (payload[field] !== undefined) {
return payload;
}
}
return {};
}
function stageStatus(key: string) {
const events = demoEvents.value.filter((event) => event.step_key === key);
if (events.length === 0) {
return "pending";
}
return String(events[events.length - 1].status);
}
function stageIndex(key: string) {
return demoStages.findIndex((stage) => stage.key === key) + 1;
}
function normalizeMarker(marker: Row, index: number): GisMarker {
return {
id: String(marker.id ?? marker.alarm_id ?? index),
index: index + 1,
scene: String(marker.scene ?? "-"),
mileage: String(marker.mileage ?? "K123+456"),
distance: String(marker.distance ?? marker.distance_to_track_m ?? ""),
severity: String(marker.severity ?? "medium"),
status: String(marker.status ?? "pending"),
x: Number(marker.x ?? 20 + index * 5),
y: Number(marker.y ?? 60 - index * 4)
};
}
function safeObject(value: unknown): Row {
if (!value) {
return {};
@@ -292,6 +662,42 @@ function safeObject(value: unknown): Row {
}
}
function objectEntries(value: unknown) {
const objectValue = safeObject(value);
return Object.entries(objectValue).map(([label, rawValue]) => ({
label,
value: formatValue(rawValue)
}));
}
function formatValue(value: unknown) {
if (Array.isArray(value)) {
return `${value.length}`;
}
if (value && typeof value === "object") {
return Object.entries(value as Row).map(([key, item]) => `${key}:${item}`).join(" ");
}
return String(value ?? "-");
}
function formatList(value: unknown) {
if (Array.isArray(value)) {
return value.join(" / ");
}
return String(value ?? "-");
}
function formatTime(value: unknown) {
if (!value) {
return "--:--:--";
}
const date = new Date(String(value));
if (Number.isNaN(date.getTime())) {
return String(value).slice(11, 19);
}
return date.toLocaleTimeString("zh-CN", { hour12: false });
}
function clamp(value: number, min: number, max: number) {
return Math.min(Math.max(value, min), max);
}
@@ -301,7 +707,8 @@ function severityLabel(value: unknown) {
critical: "严重",
high: "高",
medium: "中",
low: "低"
low: "低",
resource: "资源"
};
return labelsBySeverity[String(value)] ?? String(value ?? "-");
}
@@ -311,7 +718,8 @@ function severityTag(value: unknown) {
critical: "danger",
high: "warning",
medium: "info",
low: "success"
low: "success",
resource: "success"
};
return tags[String(value)] ?? "info";
}
@@ -322,21 +730,61 @@ function statusLabel(value: unknown) {
created: "已派发",
processing: "处置中",
closed: "已闭环",
suppressed: "已抑制"
suppressed: "已抑制",
completed: "已完成",
running: "执行中",
ingesting: "接入中",
failed: "失败",
cancelled: "已取消"
};
return labelsByStatus[String(value)] ?? String(value ?? "-");
}
function eventStatusLabel(value: unknown) {
const labelsByStatus: Record<string, string> = {
pending: "未开始",
running: "执行中",
completed: "已完成",
failed: "失败",
cancelled: "已取消"
};
return labelsByStatus[String(value)] ?? String(value ?? "未开始");
}
function statusTag(value: unknown) {
const tags: Record<string, "danger" | "warning" | "info" | "success"> = {
created: "warning",
processing: "info",
pending: "warning",
running: "info",
completed: "success",
closed: "success",
failed: "danger",
cancelled: "warning",
suppressed: "info"
};
return tags[String(value)] ?? "info";
}
function stageTag(value: unknown) {
const tags: Record<string, "danger" | "warning" | "info" | "success"> = {
running: "info",
completed: "success",
failed: "danger",
cancelled: "warning"
};
return tags[String(value)] ?? "info";
}
function resourceLabel(value: unknown) {
const labelsByType: Record<string, string> = {
image: "可见光资源",
thermal: "红外资源",
tif: "防洪TIF资源"
};
return labelsByType[String(value)] ?? String(value ?? "资源");
}
onMounted(refresh);
onBeforeUnmount(closeDemoEventSource);
</script>
+19
View File
@@ -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`;
}
+370 -3
View File
@@ -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;
}
}
@@ -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));
}
}
@@ -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<String, CopyOnWriteArrayList<SseEmitter>> 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<String, Object> event : history(runId)) {
send(emitter, event);
}
return emitter;
}
public Map<String, Object> emit(
String runId,
DemoStep step,
String status,
int progress,
String title,
String message,
String entityType,
String entityId,
Map<String, Object> metrics,
Map<String, Object> visualPayload
) {
String eventId = Ids.next("dre");
Instant now = Instant.now();
Map<String, Object> safeMetrics = metrics == null ? Map.of() : metrics;
Map<String, Object> 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<String, Object> 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<Map<String, Object>> history(String runId) {
return jdbc.query("""
select id as event_id,
run_id,
step_key,
step_name,
status,
progress,
title,
message,
entity_type,
entity_id,
metrics::text as metrics,
visual_payload::text as visual_payload,
created_at
from demo_run_events
where run_id=?
order by created_at
""",
(rs, rowNum) -> Map.ofEntries(
Map.entry("event_id", rs.getString("event_id")),
Map.entry("run_id", rs.getString("run_id")),
Map.entry("step_key", rs.getString("step_key")),
Map.entry("step_name", rs.getString("step_name")),
Map.entry("status", rs.getString("status")),
Map.entry("progress", rs.getInt("progress")),
Map.entry("title", rs.getString("title")),
Map.entry("message", rs.getString("message") == null ? "" : rs.getString("message")),
Map.entry("entity_type", rs.getString("entity_type") == null ? "" : rs.getString("entity_type")),
Map.entry("entity_id", rs.getString("entity_id") == null ? "" : rs.getString("entity_id")),
Map.entry("metrics", Jsonb.map(rs.getString("metrics"))),
Map.entry("visual_payload", Jsonb.map(rs.getString("visual_payload"))),
Map.entry("created_at", rs.getTimestamp("created_at").toInstant().toString())
),
runId
);
}
private void broadcast(String runId, Map<String, Object> event) {
for (SseEmitter emitter : emitters.getOrDefault(runId, new CopyOnWriteArrayList<>())) {
send(emitter, event);
}
}
private void send(SseEmitter emitter, Map<String, Object> event) {
try {
emitter.send(SseEmitter.event()
.id(String.valueOf(event.get("event_id")))
.name("demo-run-event")
.data(event));
} catch (IOException | IllegalStateException ignored) {
emitters.values().forEach(list -> list.remove(emitter));
}
}
private void removeEmitter(String runId, SseEmitter emitter) {
CopyOnWriteArrayList<SseEmitter> runEmitters = emitters.get(runId);
if (runEmitters != null) {
runEmitters.remove(emitter);
}
}
private static String runStatus(DemoStep step, String eventStatus) {
if (step == DemoStep.RUN_COMPLETED) {
return "completed";
}
if (step == DemoStep.RUN_FAILED || "failed".equals(eventStatus)) {
return "failed";
}
if (step == DemoStep.RUN_CANCELLED || "cancelled".equals(eventStatus)) {
return "cancelled";
}
return "running";
}
}
@@ -0,0 +1,459 @@
package com.ai.trackwalker.demo;
import com.ai.trackwalker.api.dto.Requests;
import com.ai.trackwalker.common.Ids;
import com.ai.trackwalker.common.Jsonb;
import com.ai.trackwalker.service.PlatformService;
import jakarta.annotation.PreDestroy;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.sql.Timestamp;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@Service
public class DemoRunService {
private final PlatformService platformService;
private final JdbcTemplate jdbc;
private final DemoRunEventPublisher eventPublisher;
private final ExecutorService executor = Executors.newCachedThreadPool();
public DemoRunService(PlatformService platformService, JdbcTemplate jdbc, DemoRunEventPublisher eventPublisher) {
this.platformService = platformService;
this.jdbc = jdbc;
this.eventPublisher = eventPublisher;
}
public Map<String, Object> startRun() {
String runId = Ids.next("run");
Instant now = Instant.now();
jdbc.update("insert into demo_runs(id, status, current_step, progress, summary, started_at) values (?,?,?,?,?::jsonb,?)",
runId, "running", DemoStep.INITIALIZING.key(), 0, "{}", Timestamp.from(now));
eventPublisher.emit(runId, DemoStep.INITIALIZING, "running", 2, "全链路运行实例已创建",
"正在初始化线路、场景与阶段时间轴", "demo_run", runId,
Map.of("stage_count", DemoStep.timeline().size()),
Map.of("line_id", "line-demo", "mileage_start", "K123+000", "mileage_end", "K130+000", "scenes", demoScenes()));
executor.submit(() -> executeRun(runId));
return Map.of(
"run_id", runId,
"status", "running",
"event_stream_url", "/api/v1/demo/runs/" + runId + "/events",
"snapshot_url", "/api/v1/demo/runs/" + runId
);
}
public Map<String, Object> snapshot(String runId) {
Map<String, Object> run = jdbc.queryForMap("select *, summary::text as summary_text from demo_runs where id=?", runId);
List<Map<String, Object>> events = eventPublisher.history(runId);
Map<String, Map<String, Object>> latestByStep = events.stream()
.collect(java.util.stream.Collectors.toMap(
event -> String.valueOf(event.get("step_key")),
event -> event,
(left, right) -> right
));
List<Map<String, Object>> steps = DemoStep.timeline().stream()
.map(step -> {
Map<String, Object> latest = latestByStep.get(String.valueOf(step.get("key")));
return Map.<String, Object>of(
"key", step.get("key"),
"name", step.get("name"),
"progress", latest == null ? 0 : latest.get("progress"),
"status", latest == null ? "pending" : latest.get("status")
);
})
.toList();
return Map.ofEntries(
Map.entry("run_id", run.get("id")),
Map.entry("status", run.get("status")),
Map.entry("current_step", run.get("current_step") == null ? "" : run.get("current_step")),
Map.entry("progress", run.get("progress")),
Map.entry("task_id", run.get("task_id") == null ? "" : run.get("task_id")),
Map.entry("summary", Jsonb.map(String.valueOf(run.get("summary_text")))),
Map.entry("started_at", ((Timestamp) run.get("started_at")).toInstant().toString()),
Map.entry("completed_at", run.get("completed_at") == null ? "" : ((Timestamp) run.get("completed_at")).toInstant().toString()),
Map.entry("error_message", run.get("error_message") == null ? "" : run.get("error_message")),
Map.entry("steps", steps),
Map.entry("events", events)
);
}
public List<Map<String, Object>> history(String runId) {
return eventPublisher.history(runId);
}
public SseEmitter subscribe(String runId) {
return eventPublisher.subscribe(runId);
}
public Map<String, Object> cancel(String runId) {
jdbc.update("update demo_runs set cancel_requested=true where id=? and status='running'", runId);
return Map.of("run_id", runId, "cancel_requested", true);
}
@PreDestroy
public void shutdown() {
executor.shutdownNow();
}
private void executeRun(String runId) {
try {
pause();
ensureActive(runId);
eventPublisher.emit(runId, DemoStep.INITIALIZING, "completed", DemoStep.INITIALIZING.progress(), "运行初始化完成",
"线路、场景、阶段时间轴和事件通道已准备就绪", "demo_run", runId,
Map.of("stage_count", DemoStep.timeline().size()),
Map.of("line_id", "line-demo", "mileage_start", "K123+000", "mileage_end", "K130+000", "scenes", demoScenes()));
eventPublisher.emit(runId, DemoStep.TASK_CREATED, "running", 6, "正在创建巡检任务",
"写入巡检任务、线路里程和场景清单", "inspection_task", "",
Map.of("scene_count", demoScenes().size()),
Map.of("line_id", "line-demo", "mileage_start", "K123+000", "mileage_end", "K130+000"));
Map<String, Object> task = platformService.createTask(new Requests.CreateTaskRequest(
"DEMO-" + runId.substring(Math.max(0, runId.length() - 6)),
"line-demo",
"K123+000",
"K130+000",
"route-demo",
demoScenes(),
"high",
null
));
String taskId = String.valueOf(task.get("task_id"));
jdbc.update("update demo_runs set task_id=? where id=?", taskId, runId);
eventPublisher.emit(runId, DemoStep.TASK_CREATED, "completed", DemoStep.TASK_CREATED.progress(), "巡检任务创建完成",
"任务已进入数据接入与分析流程", "inspection_task", taskId,
Map.of("task_count", 1),
Map.of("task_id", taskId, "line_id", "line-demo", "mileage_start", "K123+000", "mileage_end", "K130+000", "priority", "high"));
ensureActive(runId);
emitUavDispatch(runId, taskId);
ensureActive(runId);
List<Requests.ResourceItem> resources = demoResources();
eventPublisher.emit(runId, DemoStep.RESOURCE_INGESTING, "running", 22, "正在接入多源巡检数据",
"接入可见光、红外与防洪 TIF 资源", "inspection_task", taskId,
Map.of("resource_count", resources.size()),
Map.of("resources", resourceVisuals(resources, List.of())));
List<String> resourceIds = platformService.ingestResources(new Requests.ResourceCompleteRequest(taskId, resources));
eventPublisher.emit(runId, DemoStep.RESOURCE_INGESTING, "completed", DemoStep.RESOURCE_INGESTING.progress(), "多源数据接入完成",
"资源已写入数据链路并完成空间位置登记", "inspection_task", taskId,
Map.of("resource_count", resourceIds.size()),
Map.of("resources", resourceVisuals(resources, resourceIds)));
ensureActive(runId);
emitPreprocess(runId, taskId, resourceIds);
ensureActive(runId);
eventPublisher.emit(runId, DemoStep.AI_INFERENCING, "running", 45, "正在创建分析任务并调用模型",
"视觉、红外、点云/TIF 模型开始推理", "inspection_task", taskId,
Map.of("resource_count", resourceIds.size(), "model_count", 3),
Map.of("models", modelVisuals("running", 0), "resources", resourceVisuals(resources, resourceIds)));
Map<String, Object> job = platformService.createAnalysisJob(new Requests.CreateAnalysisJobRequest(taskId, resourceIds, demoScenes(), "offline", "high"));
String jobId = String.valueOf(job.get("analysis_job_id"));
Map<String, Object> analysis = platformService.runAnalysis(jobId);
int aiResultCount = count("ai_results where task_id=?", taskId);
eventPublisher.emit(runId, DemoStep.AI_INFERENCING, "completed", DemoStep.AI_INFERENCING.progress(), "AI推理分析完成",
"模型已生成识别、温度和形变分析结果", "analysis_job", jobId,
Map.of("ai_results", aiResultCount, "avg_confidence", avgConfidence(taskId)),
Map.of("models", modelVisuals("completed", aiResultCount), "summary", analysis.get("summary")));
ensureActive(runId);
List<Map<String, Object>> alarms = alarmsByTask(taskId);
eventPublisher.emit(runId, DemoStep.RULE_EVALUATING, "completed", DemoStep.RULE_EVALUATING.progress(), "GIS规则判定完成",
"已完成线路范围、防护区、接触网 ROI、温度和防洪阈值判定", "inspection_task", taskId,
Map.of("rule_hits", ruleHitCount(alarms), "suppressed", count("alarms where task_id=? and suppressed=true", taskId)),
Map.of("rule_hits", ruleVisuals(alarms), "markers", alarmMarkers(alarms)));
eventPublisher.emit(runId, DemoStep.ALARM_GENERATING, "completed", DemoStep.ALARM_GENERATING.progress(), "告警生成完成",
"有效告警已完成分级、证据关联和空间定位", "inspection_task", taskId,
Map.of("alarms", alarms.size(), "severity", severityCounts(taskId)),
Map.of("alarms", alarms.stream().limit(8).toList(), "markers", alarmMarkers(alarms)));
ensureActive(runId);
List<Map<String, Object>> workorders = workordersByTask(taskId);
eventPublisher.emit(runId, DemoStep.WORKORDER_DISPATCHING, "completed", DemoStep.WORKORDER_DISPATCHING.progress(), "工单派发完成",
"告警已自动生成处置工单并派发专业复核人员", "inspection_task", taskId,
Map.of("workorders", workorders.size(), "pending", pendingCount(workorders)),
Map.of("workorders", workorders.stream().limit(8).toList(), "flow", workorderFlow(workorders)));
ensureActive(runId);
int closed = closeSampleWorkorders(workorders);
workorders = workordersByTask(taskId);
eventPublisher.emit(runId, DemoStep.WORKORDER_CLOSING, "completed", DemoStep.WORKORDER_CLOSING.progress(), "工单闭环状态已更新",
"部分工单已模拟完成现场处置和复核关闭", "inspection_task", taskId,
Map.of("closed", closedCount(workorders), "pending", pendingCount(workorders)),
Map.of("workorders", workorders.stream().limit(8).toList(), "closed_this_run", closed, "flow", workorderFlow(workorders)));
ensureActive(runId);
Map<String, Object> samples = Map.of(
"candidate_samples", aiResultCount,
"closed_verified_samples", closedCount(workorders),
"labeling_queue", Math.max(alarms.size() - closedCount(workorders), 0)
);
eventPublisher.emit(runId, DemoStep.SAMPLE_FEEDBACK, "completed", DemoStep.SAMPLE_FEEDBACK.progress(), "样本回流已生成",
"告警证据、复核结果与模型版本已形成样本回流队列", "inspection_task", taskId,
samples,
Map.of("samples", samples, "model_groups", List.of("vision-detector", "thermal-analyzer", "pointcloud-analyzer")));
Map<String, Object> summary = summaryForTask(taskId);
jdbc.update("update demo_runs set summary=?::jsonb where id=?", Jsonb.write(summary), runId);
eventPublisher.emit(runId, DemoStep.RUN_COMPLETED, "completed", DemoStep.RUN_COMPLETED.progress(), "全链路运行完成",
"巡检任务、数据接入、AI识别、规则告警、工单闭环和样本回流已完成", "inspection_task", taskId,
summary,
Map.of("summary", summary, "markers", alarmMarkers(alarms), "workorders", workorders.stream().limit(8).toList()));
} catch (CancelledRunException ignored) {
// The cancellation event has already been emitted.
} catch (Exception e) {
eventPublisher.emit(runId, DemoStep.RUN_FAILED, "failed", 100, "全链路运行失败",
e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage(), "demo_run", runId,
Map.of(), Map.of("error", e.getClass().getName()));
}
}
private void emitUavDispatch(String runId, String taskId) {
eventPublisher.emit(runId, DemoStep.UAV_DISPATCHING, "running", 12, "航线任务已下发",
"无人机平台已接收线路巡检航线", "inspection_task", taskId,
Map.of("route_progress", 20),
Map.of("route_progress", 20, "uav_status", "已下发", "line_id", "line-demo"));
pause();
eventPublisher.emit(runId, DemoStep.UAV_DISPATCHING, "running", 16, "无人机巡检执行中",
"正在沿 K123+000 至 K130+000 执行采集", "inspection_task", taskId,
Map.of("route_progress", 68),
Map.of("route_progress", 68, "uav_status", "飞行中", "position", Map.of("mileage", "K126+800")));
pause();
eventPublisher.emit(runId, DemoStep.UAV_DISPATCHING, "completed", DemoStep.UAV_DISPATCHING.progress(), "航线任务完成",
"无人机任务状态回调已确认完成", "inspection_task", taskId,
Map.of("route_progress", 100),
Map.of("route_progress", 100, "uav_status", "已完成", "position", Map.of("mileage", "K130+000")));
}
private void emitPreprocess(String runId, String taskId, List<String> resourceIds) {
List<Map<String, Object>> checks = new ArrayList<>();
checks.add(Map.of("name", "图像/视频解码", "status", "completed", "quality", "清晰度达标"));
eventPublisher.emit(runId, DemoStep.PREPROCESS_RUNNING, "running", 32, "图像资源预处理完成",
"可见光与红外资源已完成解码和质量检查", "inspection_task", taskId,
Map.of("finished_checks", checks.size(), "resource_count", resourceIds.size()),
Map.of("checks", checks));
pause();
checks.add(Map.of("name", "坐标与时间同步", "status", "completed", "quality", "WGS84 坐标有效"));
eventPublisher.emit(runId, DemoStep.PREPROCESS_RUNNING, "running", 35, "空间坐标归一完成",
"采集点、航线和线路里程已建立关联", "inspection_task", taskId,
Map.of("finished_checks", checks.size(), "resource_count", resourceIds.size()),
Map.of("checks", checks));
pause();
checks.add(Map.of("name", "TIF/点云指标解析", "status", "completed", "quality", "体积变化指标有效"));
eventPublisher.emit(runId, DemoStep.PREPROCESS_RUNNING, "completed", DemoStep.PREPROCESS_RUNNING.progress(), "数据预处理完成",
"标准化资源已进入 AI 分析队列", "inspection_task", taskId,
Map.of("finished_checks", checks.size(), "resource_count", resourceIds.size()),
Map.of("checks", checks));
}
private int closeSampleWorkorders(List<Map<String, Object>> workorders) {
int limit = Math.min(2, workorders.size());
for (int i = 0; i < limit; i++) {
Map<String, Object> row = workorders.get(i);
platformService.workOrderCallback(new Requests.WorkOrderCallbackRequest(
String.valueOf(row.get("alarm_id")),
String.valueOf(row.get("workorder_id")),
"closed",
"现场处置完成",
"全链路可视化示例自动闭环",
"巡检平台",
null
));
}
return limit;
}
private void ensureActive(String runId) {
Boolean cancelled = jdbc.queryForObject("select cancel_requested from demo_runs where id=?", Boolean.class, runId);
if (Boolean.TRUE.equals(cancelled)) {
eventPublisher.emit(runId, DemoStep.RUN_CANCELLED, "cancelled", 100, "全链路运行已取消",
"用户取消了当前全链路示例运行", "demo_run", runId, Map.of(), Map.of());
throw new CancelledRunException();
}
}
private void pause() {
try {
Thread.sleep(550);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
private List<String> demoScenes() {
return List.of("违法开挖", "塔吊", "人员入侵", "防洪点精密监测", "接触网异物(高危)", "杆上设备发热");
}
private List<Requests.ResourceItem> demoResources() {
return List.of(
new Requests.ResourceItem(null, "image", "s3://rail-inspection/demo/visible-001.jpg", null, Map.of("longitude", 116.1, "latitude", 39.1, "distance_to_track_m", 72.5)),
new Requests.ResourceItem(null, "thermal", "s3://rail-inspection/demo/thermal-001.tiff", null, Map.of("temperature_c", 86.3, "longitude", 116.101, "latitude", 39.101)),
new Requests.ResourceItem(null, "tif", "s3://rail-inspection/demo/flood-001.tif", null, Map.of("volume_change_m3", 1.8, "longitude", 116.102, "latitude", 39.102))
);
}
private List<Map<String, Object>> resourceVisuals(List<Requests.ResourceItem> resources, List<String> ids) {
List<Map<String, Object>> visuals = new ArrayList<>();
for (int i = 0; i < resources.size(); i++) {
Requests.ResourceItem item = resources.get(i);
visuals.add(Map.of(
"resource_id", ids.size() > i ? ids.get(i) : "",
"resource_type", item.resourceType(),
"storage_url", item.storageUrl(),
"status", ids.size() > i ? "completed" : "ingesting",
"x", 22 + i * 24,
"y", 62 - i * 12
));
}
return visuals;
}
private List<Map<String, Object>> modelVisuals(String status, int resultCount) {
return List.of(
Map.of("model_group", "vision-detector", "status", status, "result_count", resultCount == 0 ? 0 : 12, "avg_confidence", 0.90),
Map.of("model_group", "thermal-analyzer", "status", status, "result_count", resultCount == 0 ? 0 : 3, "avg_confidence", 0.90),
Map.of("model_group", "pointcloud-analyzer", "status", status, "result_count", resultCount == 0 ? 0 : 3, "avg_confidence", 0.93)
);
}
private List<Map<String, Object>> alarmsByTask(String taskId) {
return jdbc.queryForList("""
select id as alarm_id,
task_id,
scene,
category,
severity,
confidence,
location::text as location,
rule_hits::text as rule_hits,
status,
created_at,
updated_at
from alarms
where task_id=?
order by created_at desc
""", taskId);
}
private List<Map<String, Object>> workordersByTask(String taskId) {
return jdbc.queryForList("""
select wo.id as workorder_id,
wo.alarm_id,
wo.status,
wo.assignee,
wo.close_result,
wo.comment,
wo.created_at,
wo.updated_at,
a.task_id,
a.scene,
a.category,
a.severity,
a.location::text as location
from work_orders wo
join alarms a on a.id=wo.alarm_id
where a.task_id=?
order by wo.updated_at desc
""", taskId);
}
private List<Map<String, Object>> alarmMarkers(List<Map<String, Object>> alarms) {
List<Map<String, Object>> markers = new ArrayList<>();
int total = Math.max(alarms.size() - 1, 1);
for (int i = 0; i < Math.min(alarms.size(), 18); i++) {
Map<String, Object> alarm = alarms.get(i);
double progress = alarms.size() == 1 ? 0.5 : (double) i / total;
Map<String, Object> location = Jsonb.map(String.valueOf(alarm.get("location")));
markers.add(Map.of(
"id", alarm.get("alarm_id"),
"scene", alarm.get("scene"),
"severity", alarm.get("severity"),
"status", alarm.get("status"),
"mileage", location.getOrDefault("mileage", "K123+456"),
"distance", location.getOrDefault("distance_to_track_m", ""),
"x", Math.max(5, Math.min(94, 8 + progress * 84 + (i % 3 - 1) * 2.4)),
"y", Math.max(8, Math.min(88, 70 - progress * 54 + (i % 2 == 0 ? -9 : 10)))
));
}
return markers;
}
private List<Map<String, Object>> ruleVisuals(List<Map<String, Object>> alarms) {
return alarms.stream()
.limit(8)
.map(alarm -> Map.<String, Object>of(
"alarm_id", alarm.get("alarm_id"),
"scene", alarm.get("scene"),
"severity", alarm.get("severity"),
"rule_hits", Jsonb.stringList(String.valueOf(alarm.get("rule_hits")))
))
.toList();
}
private Map<String, Object> workorderFlow(List<Map<String, Object>> workorders) {
return Map.of(
"alarm_generated", workorders.size(),
"workorder_dispatched", workorders.size(),
"processing", pendingCount(workorders),
"closed", closedCount(workorders)
);
}
private int count(String sqlSuffix, String taskId) {
return jdbc.queryForObject("select count(*) from " + sqlSuffix, Integer.class, taskId);
}
private double avgConfidence(String taskId) {
Number value = jdbc.queryForObject("select coalesce(avg(confidence),0) from ai_results where task_id=?", Number.class, taskId);
return value == null ? 0 : Math.round(value.doubleValue() * 100.0) / 100.0;
}
private int ruleHitCount(List<Map<String, Object>> alarms) {
return alarms.stream().mapToInt(alarm -> Jsonb.stringList(String.valueOf(alarm.get("rule_hits"))).size()).sum();
}
private Map<String, Object> severityCounts(String taskId) {
return jdbc.queryForList("select severity, count(*) as count from alarms where task_id=? group by severity", taskId)
.stream()
.collect(java.util.stream.Collectors.toMap(
row -> String.valueOf(row.get("severity")),
row -> row.get("count")
));
}
private int pendingCount(List<Map<String, Object>> workorders) {
return (int) workorders.stream().filter(row -> !"closed".equals(String.valueOf(row.get("status")))).count();
}
private int closedCount(List<Map<String, Object>> workorders) {
return (int) workorders.stream().filter(row -> "closed".equals(String.valueOf(row.get("status")))).count();
}
private Map<String, Object> summaryForTask(String taskId) {
List<Map<String, Object>> workorders = workordersByTask(taskId);
return Map.of(
"tasks", 1,
"resources", count("inspection_resources where task_id=?", taskId),
"analysis_jobs", count("analysis_jobs where task_id=?", taskId),
"ai_results", count("ai_results where task_id=?", taskId),
"alarms", count("alarms where task_id=? and suppressed=false", taskId),
"suppressed", count("alarms where task_id=? and suppressed=true", taskId),
"workorders", workorders.size(),
"closed_workorders", closedCount(workorders),
"pending_workorders", pendingCount(workorders),
"sample_candidates", count("ai_results where task_id=?", taskId)
);
}
private static final class CancelledRunException extends RuntimeException {
}
}
@@ -0,0 +1,51 @@
package com.ai.trackwalker.demo;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
public enum DemoStep {
INITIALIZING("run.initializing", "运行初始化", 3),
TASK_CREATED("task.created", "巡检任务创建", 10),
UAV_DISPATCHING("uav.dispatching", "航线任务下发", 18),
RESOURCE_INGESTING("resource.ingesting", "多源数据接入", 28),
PREPROCESS_RUNNING("preprocess.running", "数据预处理", 38),
AI_INFERENCING("ai.inferencing", "AI推理分析", 55),
RULE_EVALUATING("rule.evaluating", "GIS规则判定", 66),
ALARM_GENERATING("alarm.generating", "告警生成", 76),
WORKORDER_DISPATCHING("workorder.dispatching", "工单派发", 84),
WORKORDER_CLOSING("workorder.closing", "工单闭环", 91),
SAMPLE_FEEDBACK("sample.feedback", "样本回流", 96),
RUN_COMPLETED("run.completed", "完成汇总", 100),
RUN_CANCELLED("run.cancelled", "运行取消", 100),
RUN_FAILED("run.failed", "运行失败", 100);
private final String key;
private final String label;
private final int progress;
DemoStep(String key, String label, int progress) {
this.key = key;
this.label = label;
this.progress = progress;
}
public String key() {
return key;
}
public String label() {
return label;
}
public int progress() {
return progress;
}
public static List<Map<String, Object>> timeline() {
return Arrays.stream(values())
.filter(step -> step != RUN_CANCELLED && step != RUN_FAILED)
.map(step -> Map.<String, Object>of("key", step.key, "name", step.label, "progress", step.progress))
.toList();
}
}
@@ -56,6 +56,14 @@ public class PlatformService {
@Transactional
public Map<String, Object> completeResources(Requests.ResourceCompleteRequest request) {
List<String> resourceIds = ingestResources(request);
Map<String, Object> job = createAnalysisJob(new Requests.CreateAnalysisJobRequest(request.taskId(), resourceIds, null, "offline", "normal"));
runAnalysis(String.valueOf(job.get("analysis_job_id")));
return Map.of("accepted", resourceIds.size(), "analysis_job_id", job.get("analysis_job_id"));
}
@Transactional
public List<String> ingestResources(Requests.ResourceCompleteRequest request) {
List<String> resourceIds = request.resources().stream().map(resource -> resource.resourceId() == null ? Ids.next("res") : resource.resourceId()).toList();
Instant now = Instant.now();
for (int i = 0; i < request.resources().size(); i++) {
@@ -72,9 +80,7 @@ public class PlatformService {
);
}
jdbc.update("update inspection_tasks set status='data_uploading', updated_at=? where id=?", Timestamp.from(now), request.taskId());
Map<String, Object> job = createAnalysisJob(new Requests.CreateAnalysisJobRequest(request.taskId(), resourceIds, null, "offline", "normal"));
runAnalysis(String.valueOf(job.get("analysis_job_id")));
return Map.of("accepted", resourceIds.size(), "analysis_job_id", job.get("analysis_job_id"));
return resourceIds;
}
@Transactional
@@ -0,0 +1,31 @@
create table if not exists 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,
cancel_requested boolean not null default false,
started_at timestamptz not null,
completed_at timestamptz,
error_message text
);
create table if not exists 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
);
create index if not exists idx_demo_runs_status on demo_runs(status);
create index if not exists idx_demo_run_events_run on demo_run_events(run_id, created_at);
+41 -3
View File
@@ -6,14 +6,52 @@ $frontend = if ($env:RAIL_FRONTEND_URL) { $env:RAIL_FRONTEND_URL } else { "http:
Write-Host "Checking backend health via $base..."
Invoke-RestMethod "$base/actuator/health" | ConvertTo-Json -Depth 10
Write-Host "Running full-chain demo..."
Invoke-RestMethod "$base/api/v1/demo/run-full-chain" -Method POST -ContentType "application/json" -Body "{}" | ConvertTo-Json -Depth 20
Write-Host "Starting visual full-chain demo run..."
$run = (Invoke-RestMethod "$base/api/v1/demo/runs" -Method POST -ContentType "application/json" -Body "{}").data
$run | ConvertTo-Json -Depth 10
Write-Host "Waiting for visual run completion..."
$snapshot = $null
for ($i = 0; $i -lt 90; $i++) {
Start-Sleep -Seconds 1
$snapshot = (Invoke-RestMethod "$base/api/v1/demo/runs/$($run.run_id)").data
if ($snapshot.status -in @("completed", "failed", "cancelled")) {
break
}
}
if ($null -eq $snapshot) {
throw "No demo run snapshot received"
}
$summary = [ordered]@{
run_id = $snapshot.run_id
status = $snapshot.status
progress = $snapshot.progress
event_count = $snapshot.events.Count
step_count = $snapshot.steps.Count
completed_step_count = @($snapshot.steps | Where-Object { $_.status -eq "completed" }).Count
summary = $snapshot.summary
}
$summary | ConvertTo-Json -Depth 20
if ($snapshot.status -ne "completed") {
throw "Visual full-chain demo did not complete. Status: $($snapshot.status)"
}
if (@($snapshot.steps | Where-Object { $_.status -eq "completed" }).Count -ne 12) {
throw "Visual full-chain demo did not complete all 12 stages"
}
Write-Host "Fetching overview..."
Invoke-RestMethod "$base/api/v1/stats/overview" | ConvertTo-Json -Depth 20
Write-Host "Fetching workorders..."
Invoke-RestMethod "$base/api/v1/workorders" | ConvertTo-Json -Depth 20
$workorders = (Invoke-RestMethod "$base/api/v1/workorders").data.workorders
[ordered]@{
workorder_count = $workorders.Count
latest_status = if ($workorders.Count -gt 0) { $workorders[0].status } else { "" }
} | ConvertTo-Json -Depth 10
Write-Host "Checking frontend..."
(Invoke-WebRequest -UseBasicParsing $frontend).StatusCode