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

18 KiB
Raw Permalink Blame History

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_authhandler 签名与 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 顶部增加:

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

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: Commitgit 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 底部(或新文件)增加:

#[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 routerroute_layer(require_auth) 之内)加:
.route("/ws/logs", routing::get(ws::ws_logs_handler))

放在 /ws 路由附近。

  • Step 3: 验证cargo build + cargo clippy --all-targets --all-features -- -D warnings
  • Step 4: Commitgit 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

#[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 增加:

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
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 加:
.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: Commitgit 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: Commitgit 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: Commitgit 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: Commitgit 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.tomlwebui/package.json minor bump1.5.0 → 1.6.0)。
  • Step 3: Commitgit 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 checknpm run buildcargo buildcargo test --libcargo clippy -- -D warnings 全绿。

风险与开放项

  • chrono vs timetracing-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。