PicoBot/docs/superpowers/plans/2026-07-26-p2-logs-data.md

414 lines
18 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.

# 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。