docs: add P2 logs and data implementation plan

This commit is contained in:
xiaoxixi 2026-07-26 21:47:32 +08:00
parent 038854afa1
commit dff155a93c

View File

@ -0,0 +1,413 @@
# P2 日志与数据 Implementation Plan
> **For agentic workers:** REQUIRED: Use superpowers:subagent-driven-development (if subagents available) or superpowers:executing-plans to implement this plan. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** 为 PicoBot WebUI 增加实时日志流tracing 广播层 + `/ws/logs`、记忆写入端点PUT/DELETE、以及日志页/记忆页/任务页的前端重构。
**Architecture:** 新增 tracing 广播层(自定义 `Layer` impl格式化后发送到 `tokio::sync::broadcast`,容量 1024慢客户端丢旧不反压。全局 `OnceLock<broadcast::Sender<LogEvent>>``init_logging()` 时初始化。`/ws/logs` 独立 WebSocket 端点复用现有 `require_auth` 中间件,连接后订阅广播、按 level/search 过滤推送。记忆写入复用已有 `Storage::upsert_memory`/`delete_memory`。前端三页重构为 Svelte 5 runes 组件。
**Tech Stack:** Rusttracing-subscriber Layer trait、tokio broadcast、Axum WebSocket、Svelte 5runes、bits-ui。
**关键设计决策(务必遵守):**
1. **广播层全局**`src/logging.rs` 增加 `static LOG_TX: OnceLock<broadcast::Sender<LogEvent>>` + `pub fn log_sender() -> Option<broadcast::Sender<LogEvent>>``init_logging()` 初始化 broadcast channel 并挂载自定义 Layer。无订阅者时 send 为 no-op。
2. **LogEvent 结构**`#[derive(Clone)] pub struct LogEvent { pub ts: String, pub level: String, pub target: String, pub message: String }`。ts 为 RFC 3339。level 为 "TRACE"/"DEBUG"/"INFO"/"WARN"/"ERROR"。
3. **`/ws/logs` 独立于 `/ws`**:不与聊天 WS 复用连接。路由注册在 protected router 内(`require_auth`handler 签名与 `ws_handler` 类似但更简单(无 session 注册)。查询参数 `level`(可选,最低级别过滤)和 `search`(可选,关键字)。
4. **记忆 PUT 语义**`PUT /api/memories/{key}` body `{content: string, importance?: f64}`。若 key 已存在则更新 content/importance/updated_at若不存在则 404不创建新条目——创建由 agent 内部完成)。需先查询 key 是否存在。
5. **记忆 DELETE 语义**`DELETE /api/memories/{key}` 删除后返回 `{"deleted": true}`。key 不存在也返回 200幂等
6. **GET /api/logs 保留**:文件尾读取不变,用于进入页面时拉历史与重连对齐。
7. **前端日志页**:进入时先 `GET /api/logs` 拉历史尾 → 建立 `/ws/logs` 接管实时断线重连重新拉尾对齐。level 过滤(全部/INF/WRN/ERR、关键字搜索、暂停滚动、下载。
8. **前端记忆页**行内编辑textarea + importance 滑块 + 保存/取消、删除确认分级Knowledge 普通 / Timeline 强警告)、搜索 + 分类筛选 + 分页加载。
9. **前端任务页**:定时任务增加 cron 表达式显示、下次运行倒计时、最近运行状态点(绿/琥珀/红);后台任务增加脉冲动画。
**参考:** 规格 `docs/superpowers/specs/2026-07-23-webui-refactor-design.md` §6.4/§6.5/§6.6/§7.4/§7.5P1 计划 `docs/superpowers/plans/2026-07-24-p1-observability.md`(前端组件/页面模式)。
**验证约定:** Rust 改动 → 定向测试 + `cargo test --lib` + `cargo clippy --all-targets --all-features -- -D warnings` + `cargo build`;前端 → `cd webui && npm run check && npm run build`。前端无单测框架,勿虚构。
---
## Chunk 1: 后端 — tracing 广播 + /ws/logs
### Task 1.1: LogEvent + 广播 Layer + 全局访问器
**Files:**
- Modify: `src/logging.rs`
- [ ] **Step 1: 定义 LogEvent 与全局 Sender**
`src/logging.rs` 顶部增加:
```rust
use tokio::sync::broadcast;
use std::sync::OnceLock;
use tracing::field::Visit;
use tracing_subscriber::layer::Context;
use tracing_subscriber::Layer;
const LOG_BROADCAST_CAP: usize = 1024;
#[derive(Clone, Debug)]
pub struct LogEvent {
pub ts: String,
pub level: String,
pub target: String,
pub message: String,
}
static LOG_TX: OnceLock<broadcast::Sender<LogEvent>> = OnceLock::new();
pub fn log_sender() -> Option<broadcast::Sender<LogEvent>> {
LOG_TX.get().cloned()
}
```
- [ ] **Step 2: 实现自定义 Layer**
```rust
struct BroadcastLayer;
struct MessageVisitor(String);
impl Visit for MessageVisitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = format!("{:?}", value);
}
}
fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
if field.name() == "message" {
self.0 = value.to_string();
}
}
}
impl<S: tracing::Subscriber> Layer<S> for BroadcastLayer {
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
let Some(tx) = LOG_TX.get() else { return };
if tx.receiver_count() == 0 { return; }
let mut visitor = MessageVisitor(String::new());
event.record(&mut visitor);
let metadata = event.metadata();
let _ = tx.send(LogEvent {
ts: chrono::Local::now().to_rfc3339(),
level: metadata.level().to_string(),
target: metadata.target().to_string(),
message: visitor.0,
});
}
}
```
注意:需确认 `chrono` 是否已在依赖中。若无需添加 `chrono = "0.4"` 到 Cargo.toml。或者用 `time` cratetracing-subscriber 的 `local-time` feature 已引入 `time`)。优先用 `time::OffsetDateTime::now_local()` 格式化为 RFC 3339避免新增依赖。
- [ ] **Step 3: 修改 init_logging() 挂载广播层**
`init_logging()` 中,`tracing_subscriber::registry()` 链增加 `.with(BroadcastLayer)`,并在函数开头初始化 broadcast
```rust
pub fn init_logging() {
let (tx, _rx) = broadcast::channel(LOG_BROADCAST_CAP);
let _ = LOG_TX.set(tx);
// ... 现有代码 ...
tracing_subscriber::registry()
.with(env_filter)
.with(console_layer)
.with(file_layer)
.with(BroadcastLayer)
.init();
// ...
}
```
- [ ] **Step 4: 验证**`cargo build` + `cargo clippy --all-targets --all-features -- -D warnings`
- [ ] **Step 5: Commit**`git add src/logging.rs Cargo.toml && git commit -m "feat(logging): tracing broadcast layer for real-time log streaming"`
### Task 1.2: /ws/logs WebSocket 端点
**Files:**
- Modify: `src/gateway/ws.rs`(或新建 `src/gateway/ws_logs.rs`
- Modify: `src/gateway/mod.rs`(路由注册)
- [ ] **Step 1: 实现 ws_logs handler**
`src/gateway/ws.rs` 底部(或新文件)增加:
```rust
#[derive(Debug, Default, Deserialize)]
pub struct WsLogsQuery {
level: Option<String>,
search: Option<String>,
}
pub async fn ws_logs_handler(
ws: WebSocketUpgrade,
Query(query): Query<WsLogsQuery>,
Extension(_identity): Extension<super::auth::AuthIdentity>,
) -> Response {
ws.on_upgrade(|socket| async move {
handle_logs_socket(socket, query).await;
})
}
async fn handle_logs_socket(ws: WebSocket, query: WsLogsQuery) {
let Some(tx) = crate::logging::log_sender() else { return };
let mut rx = tx.subscribe();
let (mut ws_sender, mut ws_receiver) = ws.split();
let min_level = query.level.as_deref().map(parse_min_level).unwrap_or(0);
let search = query.search.filter(|s| !s.is_empty()).map(|s| s.to_ascii_lowercase());
loop {
tokio::select! {
result = rx.recv() => {
match result {
Ok(event) => {
if level_rank(&event.level) < min_level { continue; }
if let Some(needle) = &search {
if !event.message.to_ascii_lowercase().contains(needle)
&& !event.target.to_ascii_lowercase().contains(needle) {
continue;
}
}
let json = serde_json::json!({
"ts": event.ts,
"level": event.level,
"target": event.target,
"message": event.message,
});
if ws_sender.send(WsMessage::Text(json.to_string().into())).await.is_err() {
break;
}
}
Err(broadcast::error::RecvError::Lagged(_)) => continue,
Err(broadcast::error::RecvError::Closed) => break,
}
}
msg = ws_receiver.next() => {
match msg {
Some(Ok(WsMessage::Close(_))) | None => break,
_ => {}
}
}
}
}
}
fn parse_min_level(level: &str) -> u8 {
match level.to_ascii_uppercase().as_str() {
"DEBUG" => 1,
"INFO" => 2,
"WARN" => 3,
"ERROR" => 4,
_ => 0,
}
}
fn level_rank(level: &str) -> u8 {
match level.as_str() {
"TRACE" => 0,
"DEBUG" => 1,
"INFO" => 2,
"WARN" => 3,
"ERROR" => 4,
_ => 2,
}
}
```
- [ ] **Step 2: 注册路由** — 在 `src/gateway/mod.rs` 的 protected router`route_layer(require_auth)` 之内)加:
```rust
.route("/ws/logs", routing::get(ws::ws_logs_handler))
```
放在 `/ws` 路由附近。
- [ ] **Step 3: 验证**`cargo build` + `cargo clippy --all-targets --all-features -- -D warnings`
- [ ] **Step 4: Commit**`git add src/gateway && git commit -m "feat(gateway): /ws/logs real-time log streaming endpoint"`
---
## Chunk 2: 后端 — 记忆写入端点
### Task 2.1: PUT /api/memories/{key} + DELETE /api/memories/{key}
**Files:**
- Modify: `src/gateway/http.rs`
- Modify: `src/gateway/mod.rs`(路由)
- [ ] **Step 1: PUT handler**
```rust
#[derive(Deserialize)]
pub struct PutMemoryBody {
content: String,
importance: Option<f64>,
}
pub async fn put_memory(
State(state): State<Arc<GatewayState>>,
Path(key): Path<String>,
Json(body): Json<PutMemoryBody>,
) -> Result<Json<Value>, ApiError> {
let existing = state.storage
.list_memories(None, None, 1)
.await
.map_err(ApiError::internal)?;
// 需要按 key 查询——检查是否有 get_memory_by_key 方法
// 若无,用 search 或直接 SQL 查询
// 简化:直接用 upsert但需先确认 key 存在
// 实际实现:查询 SELECT * FROM memories WHERE key = ?
// 若不存在返回 404
// 若存在,更新 content/importance/updated_at
...
}
```
**注意**:当前 Storage 无 `get_memory_by_key` 方法。需在 `src/storage/memory.rs` 增加:
```rust
pub async fn get_memory_by_key(&self, key: &str) -> Result<Option<MemoryEntry>, StorageError> {
let row = sqlx::query(
"SELECT id, key, content, category, importance, session_id, created_at, updated_at FROM memories WHERE key = ?"
)
.bind(key)
.fetch_optional(self.pool())
.await?;
match row {
Some(row) => Ok(Some(parse_memory_row(&row)?)),
None => Ok(None),
}
}
```
PUT handler 完整逻辑:
1. `get_memory_by_key(&key)` → None → 404
2. 存在 → 构造更新后的 MemoryEntry保留 id/category/session_id/created_at更新 content/importance/updated_at
3. `upsert_memory(&entry)` → 200 `{"updated": true, "key": key}`
- [ ] **Step 2: DELETE handler**
```rust
pub async fn delete_memory(
State(state): State<Arc<GatewayState>>,
Path(key): Path<String>,
) -> Result<Json<Value>, ApiError> {
state.storage.delete_memory(&key).await.map_err(ApiError::internal)?;
Ok(Json(json!({ "deleted": true })))
}
```
- [ ] **Step 3: 注册路由** — protected router 加:
```rust
.route("/api/memories/{key}", routing::put(http::put_memory).delete(http::delete_memory))
```
注意:现有 `GET /api/memories`(无 path param不受影响。
- [ ] **Step 4: 验证**`cargo build` + `cargo test --lib` + `cargo clippy --all-targets --all-features -- -D warnings`
- [ ] **Step 5: Commit**`git add src/storage/memory.rs src/gateway && git commit -m "feat(gateway): PUT/DELETE /api/memories/{key} write endpoints"`
---
## Chunk 3: 前端 — 日志页重构
### Task 3.1: LogsPage 实时流式重构
**Files:**
- Modify: `webui/src/pages/LogsPage.svelte`
布局(规格 §6.4
- 工具栏level 过滤 chips全部/INF/WRN/ERR、关键字搜索、暂停滚动按钮、下载按钮。
- 连接状态条:「实时推送中 · N 行/分」或「已断开 — 重连中」。
- 日志行:时间戳 + level 着色INF=signal, WRN=accent, ERR=danger, DBG=muted+ target + 消息。
- 自动跟随尾部(除非暂停)。
- 进入页面:`GET /api/logs?lines=200` 拉历史尾 → 建立 `/ws/logs` WebSocket 接管实时。
- 断线重连:重新拉尾对齐。
- 下载:将当前 lines 导出为 .log 文件Blob + URL.createObjectURL
数据流:
- `onMount`:先 `api("/api/logs?lines=200")` 填充历史 → 建立 WS `new WebSocket(\`${wsBase}/ws/logs?level=...\`)`。
- WS `onmessage`:解析 JSON `{ts, level, target, message}`,追加到 `lines` 数组(上限 2000 行,超出 shift
- WS `onclose`:设 disconnected 状态3s 后重连(重新拉尾 + 重建 WS
- level/search 过滤:前端 `$derived` 过滤已存储行WS 查询参数在连接时固定;切换 filter 需重建 WS 或纯前端过滤——**选择纯前端过滤**WS 不带 filter 参数,简化重连逻辑)。
- 暂停:`paused` 状态,暂停时不自动滚动到底部,新行仍追加。
- [ ] **Step 1: 重写 LogsPage.svelte**(上述布局与数据流)。
- [ ] **Step 2: 验证**`cd webui && npm run check && npm run build`
- [ ] **Step 3: Commit**`git add webui/src/pages/LogsPage.svelte && git commit -m "feat(webui): real-time streaming logs page"`
---
## Chunk 4: 前端 — 记忆页重构
### Task 4.1: MemoryPage 可编辑/可删除重构
**Files:**
- Modify: `webui/src/pages/MemoryPage.svelte`
布局(规格 §6.5
- 统计条:总量 / Knowledge / Timeline。
- 搜索框 + 分类筛选(全部/Knowledge/Timeline+ 排序(最近更新)。
- 记忆卡片列表keymono accent、content、category badge、importance、updated_at、session_id。
- 行内编辑:点击「编辑」→ content 变 textarea + importance range input + 保存/取消按钮。保存调 `PUT /api/memories/{key}`
- 删除:点击「删除」→ 确认对话框。Knowledge 普通确认Timeline 强警告(红色面板 + 说明文字 + 按钮「我了解,确认删除」)。调 `DELETE /api/memories/{key}`
- 分页limit=100底部「加载更多」按钮offset 通过增大 limit 实现——现有 API 不支持 offset用 limit 递增近似)。
- [ ] **Step 1: 重写 MemoryPage.svelte**(上述布局)。
- [ ] **Step 2: 验证**`cd webui && npm run check && npm run build`
- [ ] **Step 3: Commit**`git add webui/src/pages/MemoryPage.svelte && git commit -m "feat(webui): editable memory page with delete confirmation"`
---
## Chunk 5: 前端 — 任务页重构
### Task 5.1: TasksPage 增强
**Files:**
- Modify: `webui/src/pages/TasksPage.svelte`
布局(规格 §6.6
- 定时任务卡片增强:
- 显示 cron 表达式mono
- 下次运行倒计时「3 分钟后」格式,每 30s 刷新)。
- 最近运行状态点(最多 10 个圆点:绿=completed / 琥珀=timeout / 红=error
- 可展开运行记录(已有,保留)。
- 后台任务增强:
- 运行中状态用脉冲动画(`.pulse` class
- 显示耗时created_at 到现在的差值,或 duration 字段)。
- [ ] **Step 1: 增强 TasksPage.svelte**(上述改动,保留现有结构)。
- [ ] **Step 2: 验证**`cd webui && npm run check && npm run build`
- [ ] **Step 3: Commit**`git add webui/src/pages/TasksPage.svelte && git commit -m "feat(webui): enhanced tasks page with countdown and status dots"`
---
## Chunk 6: 收尾
### Task 6.1: P2 收尾验证 + 版本号
- [ ] **Step 1: 全量验证**
- `cd webui && npm run check && npm run build`
- `cargo build`
- `cargo test --lib`
- `cargo clippy --all-targets --all-features -- -D warnings`
Expected: 全部通过。
- [ ] **Step 2: 版本号**`Cargo.toml``webui/package.json` minor bump1.5.0 → 1.6.0)。
- [ ] **Step 3: Commit**`git add -A && git commit -m "chore(release): P2 logs and data"`(仅暂存版本号文件 + 可能的文档;确认无构建产物)。
---
## P2 完成标志
- `/ws/logs` 实时推送日志帧ts/level/target/message设备鉴权保护慢客户端丢旧不阻塞。
- `GET /api/logs` 保留,用于历史拉取与重连对齐。
- `PUT /api/memories/{key}` 更新 content/importancekey 不存在返回 404。
- `DELETE /api/memories/{key}` 幂等删除。
- 日志页:实时流 + level 过滤 + 搜索 + 暂停 + 下载 + 断线重连。
- 记忆页:行内编辑 + 分级删除确认 + 搜索 + 分类筛选。
- 任务页cron 显示 + 倒计时 + 运行状态点 + 脉冲动画。
- `npm run check``npm run build``cargo build``cargo test --lib``cargo clippy -- -D warnings` 全绿。
## 风险与开放项
- **chrono vs time**tracing-subscriber 的 `local-time` feature 已引入 `time` crate。优先用 `time::OffsetDateTime::now_local()` 避免新增 chrono 依赖。若 `now_local()` 在多线程环境报错(`time` crate 的已知限制),回退到 `std::time::SystemTime` + 手动格式化或用 `chrono`
- **广播层性能**`receiver_count() == 0` 短路确保无订阅者时零开销。有订阅者时每条日志一次 format + send对高频 DEBUG 日志可能有微量开销,但 EnvFilter 默认 INFO 已过滤大部分。
- **WS 重连对齐**:前端断线重连时重新拉 `GET /api/logs` 尾部可能丢失断线期间的少量日志broadcast 不持久化)。可接受——日志页为运维辅助,非审计。
- **记忆分页**:现有 API 无 offset/cursor用 limit 递增近似。大量记忆(>1000时性能可接受SQLite 单表 LIMIT 查询)。
- **记忆 PUT 不创建**:设计决策——新记忆只由 agent 内部创建WebUI 仅编辑已有条目。若需创建能力,后续 P3 可加 POST。