Files
planet/docs/plans/custom-source-live-mock-plan.md
2026-05-07 18:06:06 +08:00

385 lines
13 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Custom Source Live Mock 计划
**状态**:实施中
**创建日期**2026-05-01
**任务名**`Custom Source Live Mock`
**核心目标**:把自定义源升级为同时支持 REST 与 WebSocket 的可映射采集入口,并提供本地 AIS mock WebSocket 服务,用于验证 Earth 船只实时新增与 upsert 链路。
## 背景
真实 AIS 接口变化频率不可控,无法稳定验证 Earth 页面“不刷新也能看到新船只”的实时链路。当前系统已经有自定义源基础设施:
- `datasource_configs` 保存 endpoint、auth、headers、config。
- `datasource_mapping_templates` 保存目标 schema 的确定性映射模板。
- `run-mapped` 支持保存后的自定义 REST 源通过 active mapping 写入目标数据。
但现有能力主要面向 REST sample 和批量 mapping缺少以下能力
- 自定义源不能明确选择 `REST``WebSocket` 采集模式。
- WebSocket 长连接、订阅消息、重连、消息路径提取还没有通用 runtime。
- `vessel_ais` 自定义数据写入后需要进入 AIS raw observation 和 `vessels` WS channel才能真实验证 Earth 实时 upsert。
- 删除自定义源时没有清晰的数据清理选项。
- 设置中心里“采集调度 / 凭证 / 自定义源”入口混杂,用户很难判断该在哪里配置。
## 已确认决策
| 项目 | 决策 |
|-----|------|
| 计划名称 | `Custom Source Live Mock` |
| 自定义源传输类型 | 支持 `REST``WebSocket` |
| 采集写入方式 | 先映射到目标 schema再由 destination handler 写入 |
| AIS mock 目标 | 优先打通 `vessel_ais`,验证 Earth 船只实时新增和同 MMSI upsert |
| mock 服务 runtime | 使用 `bun` 启动本地 mock WS 服务 |
| 凭证配置 | 支持 headers、bearer、api key、basic并保留 query/header API key 位置配置 |
| 删除策略 | 删除自定义源时允许选择是否删除该源写入的数据 |
| 合并语义 | 自定义源必须选择“合并到哪个内置数据”,作为内置源的补充数据进入同一聚合链路 |
| UI 方向 | 自定义源创建和维护放在“配置中心 > 采集器设置”的采集器下拉框内联入口;数据源页保留总览与运行控制 |
## 范围
### 本阶段要做
- 自定义源可选择 `REST``WebSocket`
- 自定义源支持请求头、凭证、query params、body、WS subscribe message。
- WebSocket 自定义源支持长连接、重连、消息解析、mapping、写入。
- `vessel_ais` 自定义源写入 AIS raw observations并广播 `vessels` channel。
- 提供 mock AIS WS 服务,持续发送新增 MMSI 和位置变更。
- 删除自定义源时提供“是否删除该源数据”的选项。
- 梳理设置中心信息架构,明确后续 UI 重构方向。
### 暂不做
- 不新增任意动态数据库表。
- 不允许用户提交可执行脚本作为 mapping。
- 不让 LLM 进入正式采集链路。
- 不把 mock 数据直接写 legacy `vessel_position`,优先写 AIS raw observations保持可追踪和可删除。
- 不在本阶段完成完整 `Earth Live Sync`,但要为后续 summary invalidation 留出 hook。
## 现状入口
| 能力 | 当前位置 |
|-----|----------|
| 自定义源配置模型 | `backend/app/models/datasource_config.py` |
| 自定义源 mapping 模型 | `backend/app/models/datasource_mapping.py` |
| 自定义源 API | `backend/app/api/v1/datasource_config.py` |
| 目标 schema registry | `backend/app/core/target_schema_registry.py` |
| mapping engine | `backend/app/services/datasource_mapping.py` |
| 数据源总览 UI | `frontend/src/pages/DataSources/DataSources.tsx` |
| 采集器设置 UI | `frontend/src/pages/Settings/Settings.tsx` |
## 目标架构
```mermaid
flowchart LR
A[Custom Source Config] --> B{source_type}
B -->|rest| C[Mapped REST Runner]
B -->|websocket| D[Mapped WS Runner]
C --> E[Mapping Engine]
D --> E
E --> F[Target Schema Validator]
F --> G{Destination Handler}
G -->|vessel_ais| H[AIS Raw Observations]
H --> I[AIS Aggregation]
H --> J[vessels WS Channel]
J --> K[Earth Vessel Upsert]
```
## 数据配置设计
短期可以继续复用 `DataSourceConfig`,避免大迁移。语义约定如下:
| 字段 | 用途 |
|-----|------|
| `name` | 自定义源唯一名称,例如 `mock_ais_ws` |
| `source_type` | `rest``websocket` |
| `endpoint` | `http(s)://...``ws(s)://...` |
| `auth_type` | `none``bearer``api_key``basic` |
| `auth_config` | token、api_key、key name、basic username/password 等 |
| `headers` | 静态请求头 |
| `config` | method、params、body、timeout、retry、WS 订阅消息、重连策略、消息路径等 |
建议 `config` 结构:
```json
{
"transport": "websocket",
"delivery_mode": "realtime_stream",
"merge_target_source": "barentswatch_vessels",
"target_schema": "vessel_ais",
"method": "GET",
"params": {},
"body": null,
"timeout": 30,
"retry": 3,
"ws_subscribe_message": {"type": "subscribe", "channel": "vessels"},
"ws_message_path": "$.data",
"ws_items_path": "$.vessels[*]",
"ws_reconnect": true,
"reconnect_delay_seconds": 3,
"debug_max_messages": null,
"delete_policy": "config_only"
}
```
## 后端实施计划
### Phase 1 — 自定义源类型与连接测试
- 允许 `source_type``rest``websocket`
- REST 连接测试保留现有 HTTP 请求逻辑。
- WebSocket 连接测试新增:
- 校验 endpoint 必须是 `ws://``wss://`
- 注入 headers 和 auth。
- 连接后可选发送 `ws_subscribe_message`
- 读取一条消息或超时返回诊断。
### Phase 2 — Mapped REST Runner 补齐
现有 `run-mapped` 继续作为 REST 一次性采集入口,补齐:
- `GET/POST` method。
- query params。
- JSON body。
- headers 和 auth 注入。
- sample limit 与响应大小限制。
- `vessel_ais` destination handler。
### Phase 3 — Mapped WebSocket Runner
新增通用 WebSocket runner读取 `DataSourceConfig + active mapping`
- 建立长连接。
- 发送可选订阅消息。
- 循环接收消息。
- JSON parse。
-`ws_message_path/ws_items_path` 提取 item 或 list。
- 使用 mapping engine 转换。
- 使用 target schema validator 校验。
- 调用 destination handler 写入。
- 更新采集任务状态:
- `connecting`
- `streaming`
- `reconnecting`
- `stopped`
- 维护运行指标:
- `messages_seen`
- `records_written`
- `unique_entities`
- `last_message_at`
- `last_error`
- 后台长连接不读取 `config.debug_max_messages`;该字段只用于显式的一次性调试运行,避免正式 WS 流被测试上限截断。
### Phase 4 — Destination Handler
为 target schema 建立明确写入处理器。
`vessel_ais` handler
- 写入 `AISRawObservation`
- `source = datasource.name`
- `delivery_mode` 来自 config默认 WS 为 `realtime_stream`、REST 为 `polling`
- `transport` 来自 `source_type`
- 生成幂等 observation hash。
- 更新 AIS source health。
- 广播 `vessels` channelpayload 使用当前 Earth 已支持的 upsert 格式。
`generic_records` handler
- 写入通用 collected data 或后续 generic store。
- 不直接进入 Earth。
### Phase 5 — 删除与数据清理
删除自定义源时新增清理策略:
| 选项 | 行为 |
|-----|------|
| 只删除配置 | 删除 `datasource_configs`,保留 mapping 和历史数据需要另行处理 |
| 删除配置和 mapping | 删除配置及对应 `datasource_mapping_templates` |
| 删除配置、mapping 和该源数据 | 同时删除该源写入的数据 |
数据删除范围:
- `collected_data.source == datasource.name`
- `ais_raw_observations.source == datasource.name`
- `ais_source_health.source == datasource.name`
不建议直接删除 legacy `vessel_position`,因为当前 legacy 表不带 source无法安全归因。自定义 AIS 源应优先只写 raw observations。
删除数据后应触发:
- `vessels` channel 的 reload/invalidation 事件,提示 Earth 重新拉船只聚合。
- 后续接入 `Earth Live Sync` 后,触发 `earth_summary` invalidation。
### Phase 6 — Mock AIS WebSocket 服务
新增脚本:
`scripts/mock-ais-ws-server.ts`
运行方式建议:
```bash
bun run mock:ais-ws
```
服务行为:
- 监听 `ws://localhost:8787/ais`
- 接受任意客户端连接。
- 可记录收到的 subscribe message。
- 每 1-2 秒发送一条 AIS-like JSON。
- 每隔 N 条生成新 MMSI验证船只数量增长。
- 已存在 MMSI 随时间改变 `lat/lon/cog/heading`,验证同 MMSI upsert。
- 支持固定 seed保证测试可复现。
示例 payload
```json
{
"type": "vessel",
"data": {
"mmsi": "999000001",
"name": "MOCK VESSEL 001",
"lat": 31.23,
"lon": 121.47,
"sog": 12.4,
"cog": 86,
"heading": 90,
"received_at": "2026-05-01T00:00:00Z"
}
}
```
## 前端实施计划
### 信息架构调整
自定义源不作为割裂的新入口,而是作为内置采集器的补充源,直接纳入“配置中心 > 采集器设置”的采集器选择器:
- 采集器下拉框同时展示内置采集器和自定义补充源。
- 下拉框右侧提供加号按钮,用于添加自定义源。
- 新建自定义源时必须选择“合并到内置数据”,例如合并到 `barentswatch_vessels`
- 选择自定义源后右侧基础配置区域沿用正常采集器配置形态支持连接测试、保存、endpoint、headers、auth、高级 JSON。
- 自定义源比内置源多一个“删除自定义源”按钮。
- 删除时弹出确认框,可勾选“同时删除该自定义源生成的所有数据”。
数据源页保留:
- 内置源总览。
- 内置源最近状态。
- 内置源手动触发。
- 不展示自定义源管理入口;自定义源创建、维护、删除统一在采集器设置中完成。
### 自定义源表单
新增或重构自定义源表单:
- 源名称。
- 类型:`REST` / `WebSocket`
- 合并到内置数据:必选,用于声明该源补充哪个内置数据域。
- endpoint。
- method/body/params仅 REST 显示。
- subscribe message/message path/items path仅 WS 显示。
- auth type。
- headers。
- target schema。
- sample/test 按钮。
- mapping assistant/preview。
- 保存并运行。
### 删除确认
删除自定义源时弹出确认:
- 默认只删除配置。
- 可勾选删除 mapping。
- 可勾选删除该源写入的数据。
- 显示将删除的数据范围和不可恢复提示。
## 验证方案
### Mock WS 验证路径
1. 启动 mock 服务:
```bash
bun run mock:ais-ws
```
2. 新建自定义源:
| 字段 | 值 |
|-----|----|
| name | `mock_ais_ws` |
| source_type | `websocket` |
| endpoint | `ws://localhost:8787/ais` |
| merge_target_source | `barentswatch_vessels` |
| target_schema | `vessel_ais` |
| ws_message_path | `$.data` |
3. 保存 active mapping
```json
{
"source": {
"items_path": "$"
},
"fields": {
"mmsi": {"path": "$.mmsi", "type": "integer"},
"name": {"path": "$.name", "type": "string"},
"lat": {"path": "$.lat", "type": "float"},
"lon": {"path": "$.lon", "type": "float"},
"sog": {"path": "$.sog", "type": "float", "default": null},
"cog": {"path": "$.cog", "type": "float", "default": null},
"heading": {"path": "$.heading", "type": "integer", "default": null},
"received_at": {"path": "$.received_at", "type": "datetime", "default": null}
}
}
```
4. 启动自定义源。
5. 打开 Earth 船只图层,不刷新页面观察:
- `vessels` WS channel 收到 `source = mock_ais_ws`
- HUD 船只数在新 MMSI 到达时增加。
- 地球出现 `MOCK VESSEL`
- 同 MMSI 后续消息更新位置和航向,不重复叠加。
### 自动化测试
后端测试:
- WebSocket 自定义源连接测试。
- WS message path 和 items path 提取。
- mapping 到 `vessel_ais`
- 写入 AIS raw observation。
- 广播 `vessels` channel。
- 删除自定义源时按策略删除 mapping 和源数据。
前端测试:
- REST/WS 表单条件显示。
- 删除确认选项。
- mock 源配置保存 payload。
- mapping preview 展示错误和成功记录。
## 风险与约束
- WebSocket 自定义源是长连接,不能沿用一次性 REST 进度条。
- 如果 mock 源写 legacy vessel 表,删除会变得不安全,因此先只写 raw observations。
- 自定义 WS 可能消息量很大,必须有 backpressure、日志限流和任务取消能力。
- 任意外部 WS 不能信任 payload必须经过 mapping 和 schema validation。
- headers/auth 不能进入 LLM mapping prompt。
## 交付顺序
1. Mock AIS WS 服务。
2. 后端自定义 WS runner。
3. `vessel_ais` destination handler 和 `vessels` broadcast。
4. 删除自定义源及数据清理。
5. 设置中心采集器下拉框内联自定义源 UI。
6. 配置中心信息架构重整。
7.`Earth Live Sync` 对接 summary invalidation。