18 KiB
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: Rust(tracing-subscriber Layer trait、tokio broadcast、Axum WebSocket)、Svelte 5(runes)、bits-ui。
关键设计决策(务必遵守):
- 广播层全局:
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。 - 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"。 /ws/logs独立于/ws:不与聊天 WS 复用连接。路由注册在 protected router 内(require_auth),handler 签名与ws_handler类似但更简单(无 session 注册)。查询参数level(可选,最低级别过滤)和search(可选,关键字)。- 记忆 PUT 语义:
PUT /api/memories/{key}body{content: string, importance?: f64}。若 key 已存在则更新 content/importance/updated_at;若不存在则 404(不创建新条目——创建由 agent 内部完成)。需先查询 key 是否存在。 - 记忆 DELETE 语义:
DELETE /api/memories/{key}删除后返回{"deleted": true}。key 不存在也返回 200(幂等)。 - GET /api/logs 保留:文件尾读取不变,用于进入页面时拉历史与重连对齐。
- 前端日志页:进入时先
GET /api/logs拉历史尾 → 建立/ws/logs接管实时;断线重连重新拉尾对齐。level 过滤(全部/INF/WRN/ERR)、关键字搜索、暂停滚动、下载。 - 前端记忆页:行内编辑(textarea + importance 滑块 + 保存/取消)、删除确认分级(Knowledge 普通 / Timeline 强警告)、搜索 + 分类筛选 + 分页加载。
- 前端任务页:定时任务增加 cron 表达式显示、下次运行倒计时、最近运行状态点(绿/琥珀/红);后台任务增加脉冲动画。
参考: 规格 docs/superpowers/specs/2026-07-23-webui-refactor-design.md §6.4/§6.5/§6.6/§7.4/§7.5;P1 计划 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 crate(tracing-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: 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 底部(或新文件)增加:
#[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)之内)加:
.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
#[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 完整逻辑:
get_memory_by_key(&key)→ None → 404- 存在 → 构造更新后的 MemoryEntry(保留 id/category/session_id/created_at,更新 content/importance/updated_at)
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: 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/logsWebSocket 接管实时。 - 断线重连:重新拉尾对齐。
- 下载:将当前 lines 导出为 .log 文件(Blob + URL.createObjectURL)。
数据流:
-
onMount:先api("/api/logs?lines=200")填充历史 → 建立 WSnew 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)+ 排序(最近更新)。
-
记忆卡片列表:key(mono 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)。
- 可展开运行记录(已有,保留)。
-
后台任务增强:
- 运行中状态用脉冲动画(
.pulseclass)。 - 显示耗时(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 buildcargo buildcargo test --libcargo clippy --all-targets --all-features -- -D warningsExpected: 全部通过。
- Step 2: 版本号 —
Cargo.toml与webui/package.jsonminor bump(1.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/importance,key 不存在返回 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-timefeature 已引入timecrate。优先用time::OffsetDateTime::now_local()避免新增 chrono 依赖。若now_local()在多线程环境报错(timecrate 的已知限制),回退到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。