PicoBot/docs/superpowers/plans/2026-07-24-p1-observability.md

35 KiB
Raw Permalink Blame History

P1 观测 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 增加运行状况观测能力:进程级 Metrics 采集token/费用/工具调用/turn/延迟)、GET /api/status 聚合快照、GET /api/toolsGET /api/skills 只读端点,以及前端概览页(运行仪表盘)与工具&Skills 页。

Architecture: 新增进程级 MetricsOnceLock<Arc<Metrics>> 全局访问器,结构体可直接单测;选择全局而非穿线,避免改动 5 个 AgentLoop 构造点且天然跨配置重载存活——precedented by mcp::MCP_SERVER_STATUS。AgentLoop 在工具执行/模型调用/turn 完成处记录到全局 Metrics。/api/status handler 聚合 Metrics + 各服务只读内省(需为 MessageBus/TaskSupervisor/SessionManager/OutboundDispatcher 补轻量内省方法 + ws 连接计数 + 进程 uptime 静态)。前端概览页每 2s 轮询 /api/status,工具页拉取 /api/tools+/api/skills

Tech Stack: RustAxum handler、tokio、原子计数、Svelte 5runes、手写 SVG sparkline。

关键设计决策(务必遵守):

  1. Metrics 进程级全局src/observability/metrics.rs 定义 Metrics 结构体 + pub fn global_metrics() -> Arc<Metrics>OnceLock 首次调用初始化。AgentLoop 与 /api/status 都通过 global_metrics() 访问,不穿线、不给 GatewayState 加字段。
  2. cost 可选:在 LLMProviderConfig 增加可选 price_input_per_million: Option<f64> / price_output_per_million: Option<f64>serde default None。Metrics 记录 tokencost 仅在 AgentLoop 知道单价时累加(从 provider config 取),缺省为 0。不引入硬编码价格表。
  3. uptimesrc/gateway/mod.rs 增加进程级 static STARTED: OnceLock<std::time::Instant> + fn process_uptime_secs() -> u64,跨重载存活。
  4. active_lanesOutboundDispatcher 增加 active_lanes: Arc<AtomicUsize> 字段构造时注入lane 生成时 +1、lane 任务结束时 -1暴露 active_lane_count()。Gateway 侧保留一个 clone 供 /api/status 读取。
  5. failed_7d 近似/api/status 的调度器失败数用「last_status 为失败态error/timeout/delivery_error的任务数」近似避免 2s 轮询时逐任务查运行记录。在响应字段命名为 failed_jobs(语义=最近一次运行失败的任务数UI 标注清楚。
  6. provider status 派生Metrics 维护 per-provider 最近调用错误率;status = 最近窗口错误率 > 阈值(如最近 10 次中 ≥3 次失败)→ "degraded",否则 "ok"

参考: 规格 docs/superpowers/specs/2026-07-23-webui-refactor-design.md §6.2/§6.3/§7.2/§7.3/§7.6P0 计划 docs/superpowers/plans/2026-07-23-p0-webui-foundation.md(前端组件/页面模式)。

验证约定: Rust 改动 → 定向测试 + cargo test --lib + cargo clippy --all-targets --all-features -- -D warnings + cargo build;前端 → cd webui && npm run check && npm run build。前端无单测框架,勿虚构。

Implementation Progress (2026-07-26)

Branch: feat/webui-p1

Completed and reviewed:

  • Task 1.1 Metrics collector — 11474c0
  • Task 1.2 optional provider pricing — 115e77f
  • Task 1.3 AgentLoop instrumentation — 5d0cf5b
  • Task 2.1 runtime introspection — aa989cb
  • Task 2.2 protected GET /api/statusb14edc4, with review fixes in 87cf500
  • Task 2.3 protected GET /api/tools119df57
  • Task 2.4 protected GET /api/skills + SkillsLoader accessor — 5e2771c
  • Task 3.1 Sparkline/CapacityMeter/MetricTile components — 6e719e7
  • Task 3.2 OverviewPage runtime dashboard — 0b02fd5
  • Task 3.3 ToolsPage tools & skills browser — 266fb8c
  • Task 3.4 App.svelte wiring — 4621210
  • Task 4.1 P1 release — 7676502

Final verification baseline: cargo test --lib 350 passed; cargo clippy --all-targets --all-features -- -D warnings clean; cargo build passed; npm run check 0 errors/warnings; npm run build passed.

Status: P1 COMPLETE. Version bumped to 1.5.0.


Chunk 1: Metrics 后端

Task 1.1: Metrics 模块(结构体 + 记录/快照 + 全局访问器 + 单测)

Files:

  • Create: src/observability/metrics.rs

  • Modify: src/observability/mod.rspub mod metrics;,若 observability 模块存在;否则按代码库实际模块组织放置——先读 src/observability/ 确认)

  • Step 1: 写失败测试(在 metrics.rs 底部 #[cfg(test)]

覆盖record_turn 累加 tokens/turns/延迟窗口record_tool_call 累加总量与 per-toolrecord_provider 累加 per-provider token/cost/延迟/错误snapshot 返回正确聚合p95 计算provider status 派生ok vs degraded。示例

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn records_turns_tokens_and_p95() {
        let m = Metrics::new();
        for i in 1..=100 {
            m.record_turn(Some(&Usage { prompt_tokens: 10, completion_tokens: 20, total_tokens: 30, ..Default::default() }), i);
        }
        let s = m.snapshot();
        assert_eq!(s.turns, 100);
        assert_eq!(s.tokens_in, 1000);
        assert_eq!(s.tokens_out, 2000);
        assert!(s.turn_latency_p95_ms >= 95); // p95 of 1..=100
    }

    #[test]
    fn records_per_tool_counts() {
        let m = Metrics::new();
        m.record_tool_call("bash", true);
        m.record_tool_call("bash", true);
        m.record_tool_call("read_file", false);
        let s = m.snapshot();
        assert_eq!(s.tool_calls, 3);
        assert_eq!(s.per_tool.get("bash"), Some(&2));
        assert_eq!(s.per_tool.get("read_file"), Some(&1));
    }

    #[test]
    fn derives_provider_status() {
        let m = Metrics::new();
        for _ in 0..7 { m.record_provider("openai", "gpt-4o", None, 100, false); }
        for _ in 0..3 { m.record_provider("openai", "gpt-4o", None, 100, true); }
        let s = m.snapshot();
        let p = s.providers.iter().find(|p| p.name == "openai").unwrap();
        assert_eq!(p.status, "degraded"); // 3/10 recent errors
    }
}
  • Step 2: 运行测试确认失败cargo test --lib metrics → FAIL模块不存在
  • Step 3: 实现 Metrics
use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Instant;

use crate::providers::Usage;

const WINDOW: usize = 100;          // 延迟/错误滚动窗口大小
const DEGRADE_THRESHOLD: usize = 3; // 最近 10 次中失败 ≥3 → degraded
const DEGRADE_WINDOW: usize = 10;

#[derive(Default)]
struct ProviderStat {
    model: String,
    tokens_in: u64,
    tokens_out: u64,
    cost: f64,
    calls: u64,
    last_latency_ms: u64,
    latencies: VecDeque<u64>,
    recent_results: VecDeque<bool>, // true = error
}

pub struct Metrics {
    tokens_in: std::sync::atomic::AtomicU64,
    tokens_out: std::sync::atomic::AtomicU64,
    turns: std::sync::atomic::AtomicU64,
    tool_calls: std::sync::atomic::AtomicU64,
    per_tool: Mutex<HashMap<String, u64>>,
    turn_latencies: Mutex<VecDeque<u64>>,
    providers: Mutex<HashMap<String, ProviderStat>>, // key = provider name
}

#[derive(serde::Serialize)]
pub struct MetricsSnapshot {
    pub tokens_in: u64,
    pub tokens_out: u64,
    pub cost: f64,
    pub turns: u64,
    pub tool_calls: u64,
    pub turn_latency_p95_ms: u64,
    pub per_tool: HashMap<String, u64>,
    pub providers: Vec<ProviderSnapshot>,
}

#[derive(serde::Serialize)]
pub struct ProviderSnapshot {
    pub name: String,
    pub model: String,
    pub status: String, // "ok" | "degraded"
    pub latency_ms: u64,
    pub latencies: Vec<u64>, // sparkline
    pub tokens_in: u64,
    pub tokens_out: u64,
    pub cost: f64,
}

impl Default for Metrics {
    fn default() -> Self { Self::new() }
}

impl Metrics {
    pub fn new() -> Self {
        Self {
            tokens_in: Default::default(),
            tokens_out: Default::default(),
            turns: Default::default(),
            tool_calls: Default::default(),
            per_tool: Mutex::new(HashMap::new()),
            turn_latencies: Mutex::new(VecDeque::new()),
            providers: Mutex::new(HashMap::new()),
        }
    }

    pub fn record_turn(&self, usage: Option<&Usage>, latency_ms: u64) {
        use std::sync::atomic::Ordering::Relaxed;
        self.turns.fetch_add(1, Relaxed);
        if let Some(u) = usage {
            self.tokens_in.fetch_add(u.prompt_tokens as u64, Relaxed);
            self.tokens_out.fetch_add(u.completion_tokens as u64, Relaxed);
        }
        let mut q = self.turn_latencies.lock().unwrap_or_else(|e| e.into_inner());
        q.push_back(latency_ms);
        while q.len() > WINDOW { q.pop_front(); }
    }

    pub fn record_tool_call(&self, name: &str, _success: bool) {
        use std::sync::atomic::Ordering::Relaxed;
        self.tool_calls.fetch_add(1, Relaxed);
        *self.per_tool.lock().unwrap_or_else(|e| e.into_inner()).entry(name.to_string()).or_insert(0) += 1;
    }

    pub fn record_provider(&self, name: &str, model: &str, cost: Option<f64>, latency_ms: u64, is_error: bool) {
        let mut map = self.providers.lock().unwrap_or_else(|e| e.into_inner());
        let stat = map.entry(name.to_string()).or_insert_with(|| ProviderStat { model: model.to_string(), ..Default::default() });
        stat.model = model.to_string();
        stat.calls += 1;
        stat.last_latency_ms = latency_ms;
        stat.latencies.push_back(latency_ms);
        while stat.latencies.len() > WINDOW { stat.latencies.pop_front(); }
        stat.recent_results.push_back(is_error);
        while stat.recent_results.len() > DEGRADE_WINDOW { stat.recent_results.pop_front(); }
        if let Some(c) = cost { stat.cost += c; }
    }

    pub fn record_provider_tokens(&self, name: &str, usage: &Usage) {
        let mut map = self.providers.lock().unwrap_or_else(|e| e.into_inner());
        let stat = map.entry(name.to_string()).or_default();
        stat.tokens_in += usage.prompt_tokens as u64;
        stat.tokens_out += usage.completion_tokens as u64;
    }

    pub fn tool_call_count(&self, name: &str) -> u64 {
        *self.per_tool.lock().unwrap_or_else(|e| e.into_inner()).get(name).unwrap_or(&0)
    }

    pub fn snapshot(&self) -> MetricsSnapshot {
        use std::sync::atomic::Ordering::Relaxed;
        let latencies = self.turn_latencies.lock().unwrap_or_else(|e| e.into_inner());
        let p95 = percentile_95(&latencies);
        let providers = self.providers.lock().unwrap_or_else(|e| e.into_inner());
        let mut cost = 0.0;
        let provider_snaps = providers.iter().map(|(name, s)| {
            cost += s.cost;
            let errors = s.recent_results.iter().filter(|e| **e).count();
            ProviderSnapshot {
                name: name.clone(),
                model: s.model.clone(),
                status: if s.recent_results.len() >= DEGRADE_WINDOW && errors >= DEGRADE_THRESHOLD { "degraded".into() } else { "ok".into() },
                latency_ms: s.last_latency_ms,
                latencies: s.latencies.iter().copied().collect(),
                tokens_in: s.tokens_in,
                tokens_out: s.tokens_out,
                cost: s.cost,
            }
        }).collect();
        MetricsSnapshot {
            tokens_in: self.tokens_in.load(Relaxed),
            tokens_out: self.tokens_out.load(Relaxed),
            cost,
            turns: self.turns.load(Relaxed),
            tool_calls: self.tool_calls.load(Relaxed),
            turn_latency_p95_ms: p95,
            per_tool: self.per_tool.lock().unwrap_or_else(|e| e.into_inner()).clone(),
            providers: provider_snaps,
        }
    }
}

fn percentile_95(values: &VecDeque<u64>) -> u64 {
    if values.is_empty() { return 0; }
    let mut sorted: Vec<u64> = values.iter().copied().collect();
    sorted.sort_unstable();
    let idx = ((sorted.len() as f64 * 0.95).ceil() as usize).saturating_sub(1).min(sorted.len() - 1);
    sorted[idx]
}

static GLOBAL: OnceLock<Arc<Metrics>> = OnceLock::new();

/// Process-global metrics, shared across config reloads (same process).
pub fn global_metrics() -> Arc<Metrics> {
    GLOBAL.get_or_init(|| Arc::new(Metrics::new())).clone()
}

Usage 字段:prompt_tokens/completion_tokens/total_tokens 等,见 src/providers/traits.rs:118-129#[derive(Default)]record_provider 的 degraded 判定按「最近 DEGRADE_WINDOW 次中错误 ≥ DEGRADE_THRESHOLD」测试需相应构造样本量。若测试断言与实现阈值不符以实现为准调整测试样本。

  • Step 4: 运行测试确认通过cargo test --lib metrics → PASS
  • Step 5: clippycargo clippy --all-targets --all-features -- -D warnings
  • Step 6: Commitgit add src/observability && git commit -m "feat(observability): process-global Metrics collector"

Task 1.2: 配置可选价格字段

Files:

  • Modify: src/config/mod.rsLLMProviderConfig 结构体)

  • Step 1: 加可选字段 — 在 LLMProviderConfig 增加serde 可选,缺省 None不影响现有配置解析

#[serde(default)]
pub price_input_per_million: Option<f64>,
#[serde(default)]
pub price_output_per_million: Option<f64>,
  • Step 2: 加 cost 辅助方法(在 LLMProviderConfig impl
pub fn cost_of(&self, prompt_tokens: u32, completion_tokens: u32) -> Option<f64> {
    match (self.price_input_per_million, self.price_output_per_million) {
        (Some(pi), Some(po)) => Some(prompt_tokens as f64 / 1e6 * pi + completion_tokens as f64 / 1e6 * po),
        _ => None,
    }
}
  • Step 3: 验证cargo test --lib config + cargo clippy -- -D warnings(确认现有配置测试仍过,新字段不破坏反序列化)。
  • Step 4: Commitgit add src/config/mod.rs && git commit -m "feat(config): optional provider pricing for cost metrics"

Task 1.3: AgentLoop 埋点(记录到全局 Metrics

Files:

  • Modify: src/agent/agent_loop.rs

埋点位置(来自探查):

  • 工具执行:execute_one_tool(约 970-1021 行),tool_name 已知(约 977 行),成功与否在 result.success(约 1004/1015 行)。在已有 observer ObserverEvent::ToolCall 附近加 crate::observability::metrics::global_metrics().record_tool_call(&tool_name, success);

  • 模型调用延迟 + per-providerstream_completion(约 488-525 行)是唯一模型调用入口。用 Instant(已 import包裹 self.provider.stream(request).await 累加循环,得到 latency_ms;调用 global_metrics().record_provider(self.provider.name(), self.provider.model_id(), cost, latency_ms, is_error)record_provider_tokens(name, &response.usage)。cost 从 provider config 取——但 AgentLoop 当前不持有 provider config 单价;简化cost 传 None除非能拿到 config。若 AgentLoop 无法拿到单价cost 由 record_provider 传 None保持 0:若需 cost可在构造 AgentLoop 时把 (price_in, price_out) 一并传入为降低侵入P1 先传 Nonecost 留待有 config 接入时补(在计划风险项注明)。

  • turn 完成:在 process_inner 的三个返回点(约 706、847、879 行)已知 accumulated_usageturn 延迟用进入 process_inner 时的 Instant 到返回的 elapsed。加 global_metrics().record_turn(accumulated_usage_opt, latency_ms)

  • Step 1: 工具埋点 — 在 execute_one_tool 的成功/失败分支(已有 observer ToolCall 事件处)加 record_tool_call(&tool_name, success)

  • Step 2: 模型调用埋点 — 在 stream_completion 用 Instant 测延迟,结束后 record_provider(name, model, None, latency_ms, is_error) + record_provider_tokens(name, &usage)。provider 名用 self.provider.name()model 用 self.provider.model_id()(见 src/providers/traits.rs:145-149)。注意:stream_completion 的错误通过 ? 提前返回(约 494-497、500-503 行)——必须在这些错误路径也 record_provider(..., is_error=true),否则 provider status 的 degraded 判定在生产中永远不触发(单测因直接调 record_provider 仍能过,会掩盖此 bug 可用一个在成功/失败都执行的收尾闭包或在小函数返回前统一记录。

  • Step 3: turn 埋点 — 在 process_inner 入口记 let turn_start = Instant::now();,三个返回点前 global_metrics().record_turn(usage_opt, turn_start.elapsed().as_millis() as u64)

  • Step 4: 验证cargo test --lib(现有 agent 测试仍过)+ cargo clippy -- -D warnings + cargo build

  • Step 5: Commitgit add src/agent/agent_loop.rs && git commit -m "feat(agent): record tool/turn/provider metrics"


Chunk 2: 内省接口 + 端点

Task 2.1: 服务内省方法

Files:

  • Modify: src/bus/mod.rsMessageBus 队列深度)

  • Modify: src/bus/dispatcher.rsactive lane 计数)

  • Modify: src/task_supervisor.rs(运行任务数)

  • Modify: src/session/session.rs(会话数 + 活动 turn 数)

  • Modify: src/gateway/ws.rsws 连接计数)

  • Modify: src/gateway/mod.rsuptime 静态 + ws_connections 字段 + dispatcher lane 计数注入)

  • Step 1: MessageBus 队列深度 — 在 src/bus/mod.rs 加:

pub fn queue_depths(&self) -> QueueDepths {
    QueueDepths {
        inbound_depth: (self.inbound_tx.max_capacity() - self.inbound_tx.capacity()) as u64,
        inbound_cap: self.inbound_tx.max_capacity() as u64,
        outbound_depth: (self.outbound_tx.max_capacity() - self.outbound_tx.capacity()) as u64,
        outbound_cap: self.outbound_tx.max_capacity() as u64,
        control_depth: (self.control_tx.max_capacity() - self.control_tx.capacity()) as u64,
        control_cap: self.control_tx.max_capacity() as u64,
    }
}

并定义 #[derive(serde::Serialize)] pub struct QueueDepths { pub inbound_depth: u64, pub inbound_cap: u64, pub outbound_depth: u64, pub outbound_cap: u64, pub control_depth: u64, pub control_cap: u64 }。(模式参考 session.rs:2139 已用 capacity()/max_capacity()。)

  • Step 2: OutboundDispatcher active lane 计数 — 给 OutboundDispatcher 加字段 active_lanes: Arc<std::sync::atomic::AtomicUsize>,构造时注入(在 gateway 创建 dispatcher 处 Arc::new(AtomicUsize::new(0)))。计数增减位置(避免 spawn 失败泄漏):把该 Arc clone 进 lane 任务,在已生成的 lane 任务体起始处 fetch_add(1, Relaxed),在 lane 循环的唯一退出口(约 131-146 行 break 之后、任务结束前)fetch_sub(1, Relaxed)——不要在调用 spawn_lane 之前 +1spawn 可能失败)。在 gateway 侧把同一 Arc clone 存入 GatewayState 新字段 pub outbound_lanes: Arc<std::sync::atomic::AtomicUsize>from_config 字面量初始化),供 /api/status 直接 state.outbound_lanes.load(Relaxed) 读取(无需额外方法)。

  • Step 3: TaskSupervisor 运行任务数 — 在 src/task_supervisor.rs 加:

pub fn running_count(&self) -> usize {
    let mut state = self.inner.state.lock().unwrap_or_else(|e| e.into_inner());
    state.tasks.retain(|task| !task.handle.is_finished());
    state.tasks.len()
}
  • Step 4: SessionManager 会话数 + 活动 turn 数 — 在 src/session/session.rs 加(镜像 wait_until_idle 约 2125-2158 的判定):
pub async fn session_count(&self) -> usize {
    self.inner.lock().await.sessions.len()
}
pub async fn active_turn_count(&self) -> usize {
    let sessions: Vec<_> = self.inner.lock().await.sessions.values().cloned().collect();
    let mut count = 0;
    for session in sessions {
        let s = session.lock().await;
        if s.current_cancel.is_some() { count += 1; }
    }
    count
}

innerArc<Mutex<SessionManagerInner>>sessions: HashMap<String, Arc<Mutex<Session>>>Session.current_cancel: Option<oneshot::Sender> 表示 turn 执行中。锁均为 tokio Mutex纯内存计数不跨 I/O 持锁,符合不变量。)

  • Step 5: ws 连接计数 + uptime 静态 — 在 src/gateway/mod.rs

    • static STARTED: OnceLock<std::time::Instant> = OnceLock::new();pub fn process_uptime_secs() -> u64 { STARTED.get_or_init(Instant::now).elapsed().as_secs() }(在 run() 起始调用一次 STARTED.get_or_init(Instant::now))。
    • GatewayState 加 pub ws_connections: Arc<std::sync::atomic::AtomicUsize>(在 from_config 结构体字面量初始化 Default::default())。
    • src/gateway/ws.rshandle_socket(约 39-132 行):进入时 state.ws_connections.fetch_add(1, Relaxed),在函数尾部清理处(约 120-131 行)fetch_sub(1, Relaxed)。为保证多 break 路径都减,用一个 RAII guard 或在唯一尾部清理段减(确认 handle_socket 是否有单一出口;若多出口,用 guard struct Drop 减)。
  • Step 6: 验证cargo build + cargo test --lib + cargo clippy -- -D warnings

  • Step 7: Commitgit add src/bus src/task_supervisor.rs src/session/session.rs src/gateway && git commit -m "feat(gateway): runtime introspection for status endpoint"

Task 2.2: GET /api/status 聚合 handler

Files:

  • Modify: src/gateway/http.rshandler

  • Modify: src/gateway/mod.rs(注册路由)

  • Step 1: handler — 在 src/gateway/http.rs 加(聚合所有来源;用 serde_json::json!

pub async fn get_status(State(state): State<Arc<GatewayState>>) -> Result<Json<Value>, ApiError> {
    let reload = state.reload.status();
    let metrics = crate::observability::metrics::global_metrics().snapshot();

    let depths = state.bus().queue_depths();
    let active_lanes = state.outbound_lanes.load(std::sync::atomic::Ordering::Relaxed);
    let ws_connections = state.ws_connections.load(std::sync::atomic::Ordering::Relaxed);
    let background_tasks = state.task_supervisor.running_count();
    let sessions_total = state.session_manager.session_count().await;
    let active_turns = state.session_manager.active_turn_count().await;

    // 渠道状态
    let mut channels = Vec::new();
    for name in state.channel_manager.list_channel_names().await {
        let running = state.channel_manager.get_channel(&name).await.map(|c| c.is_running()).unwrap_or(false);
        channels.push(json!({ "name": name, "status": if running { "connected" } else { "stopped" } }));
    }

    // 调度器(失败数用 last_status 近似)
    let jobs = state.storage.list_scheduled_jobs().await.map_err(ApiError::internal)?;
    let failed_jobs = jobs.iter().filter(|j| matches!(j.last_status.as_deref(), Some("error") | Some("timeout") | Some("delivery_error"))).count();
    let enabled = jobs.iter().filter(|j| j.enabled).count();
    let next_run = jobs.iter().filter(|j| j.enabled).map(|j| j.next_run_at).min();

    Ok(Json(json!({
        "generation": reload.generation,
        "version": env!("CARGO_PKG_VERSION"),
        "uptime_secs": crate::gateway::process_uptime_secs(),
        "phase": reload.phase,
        "ws_connections": ws_connections,
        "background_tasks": background_tasks,
        "sessions": { "total": sessions_total, "active_turns": active_turns },
        "metrics": {
            "tokens_in": metrics.tokens_in,
            "tokens_out": metrics.tokens_out,
            "cost": metrics.cost,
            "tool_calls": metrics.tool_calls,
            "turns": metrics.turns,
            "turn_latency_p95_ms": metrics.turn_latency_p95_ms,
        },
        "bus": {
            "inbound": { "depth": depths.inbound_depth, "cap": depths.inbound_cap },
            "outbound": { "depth": depths.outbound_depth, "cap": depths.outbound_cap },
            "control": { "depth": depths.control_depth, "cap": depths.control_cap },
            "active_lanes": active_lanes,
        },
        "providers": metrics.providers,
        "channels": channels,
        "scheduler": { "jobs": jobs.len(), "enabled": enabled, "failed_jobs": failed_jobs, "next_run_at": next_run },
        "mcp": crate::mcp::get_mcp_status().iter().map(|s| json!({ "name": s.name, "connected": s.connected, "tools": s.tools.len() })).collect::<Vec<_>>(),
    })))
}

ReloadPhaseSerializesnake_caseMcpServerStatus 字段见 src/mcp/mod.rs:27-34state.bus()/channel_manager/task_supervisor/session_manager/storage/reload 均为 GatewayState 字段;outbound_lane_count() 按 Task 2.1 实际暴露方式调用。)

  • Step 2: 注册路由 — 在 src/gateway/mod.rs 的 protected router约 574-575 行 /api/tasks//api/jobs 附近)加 .route("/api/status", routing::get(http::get_status))。确认在 route_layer(require_auth) 之内(设备鉴权)。
  • Step 3: 验证cargo build + cargo clippy -- -D warnings。(可选:cargo run -- gatewaycurl 需鉴权,目检 JSON 结构;或写一个轻量集成断言。无鉴权 token 时跳过 curl靠编译 + 结构审查。)
  • Step 4: Commitgit add src/gateway && git commit -m "feat(gateway): GET /api/status runtime snapshot"

Task 2.3: GET /api/tools handler

Files:

  • Modify: src/gateway/http.rs

  • Modify: src/gateway/mod.rs(路由)

  • Step 1: handler — 遍历 state.session_manager.tools().iter()(返回 Vec<(String, Arc<dyn Tool>)>,见 src/tools/registry.rs:69-76。source 由命名约定派生:名字含 __"mcp"MCP 工具名为 {server}__{tool},见 src/mcp/tool_wrapper.rs:26),否则 "builtin"。call_count 取 global_metrics().tool_call_count(&name)

pub async fn get_tools(State(state): State<Arc<GatewayState>>) -> Result<Json<Value>, ApiError> {
    let metrics = crate::observability::metrics::global_metrics();
    let tools: Vec<Value> = state.session_manager.tools().iter().iter().map(|(name, tool)| {
        let source = if name.contains("__") { "mcp" } else { "builtin" };
        json!({
            "name": name,
            "description": tool.description(),
            "parameters_schema": tool.parameters_schema(),
            "source": source,
            "read_only": tool.read_only(),
            "exclusive": tool.exclusive(),
            "concurrency_safe": tool.concurrency_safe(),
            "call_count": metrics.tool_call_count(name),
        })
    }).collect();
    Ok(Json(json!({ "tools": tools })))
}
  • Step 2: 注册路由 — protected router 加 .route("/api/tools", routing::get(http::get_tools))
  • Step 3: 验证cargo build + cargo clippy -- -D warnings
  • Step 4: Commitgit add src/gateway && git commit -m "feat(gateway): GET /api/tools with capability metadata"

Task 2.4: GET /api/skills handler+ SkillsLoader 访问器)

Files:

  • Modify: src/session/session.rs(暴露 SkillsLoader 访问器)

  • Modify: src/gateway/http.rs

  • Modify: src/gateway/mod.rs(路由)

  • Step 1: 暴露 SkillsLoaderSessionManager 私有字段 skills_loader: Arc<SkillsLoader>session.rs:1423)无公开访问器。加:

pub fn skills_loader(&self) -> Arc<crate::skills::SkillsLoader> {
    self.skills_loader.clone()
}
  • Step 2: source 派生辅助Skill.path 是技能目录source 由路径前缀派生(~/.agents/skillsagent~/.picobot/skillspicobot{workspace}/skillsworkspace)。由于三个目录字段是 SkillsLoader 私有,简化:在 handler 里用 Skill.path 的字符串包含关系粗略派生(包含 /.agents/skills→agent/.picobot/skills→picobot否则 workspace/other。或在 SkillsLoader 加 pub fn source_of(&self, path) -> &str 辅助(更干净,推荐)。选择在 SkillsLoader 加辅助方法。
  • Step 3: handler — 用 get_loaded_skills()src/skills/mod.rs:252,返回 Vec<Skill>不要list_skills())。默认不返回 content
pub async fn get_skills(State(state): State<Arc<GatewayState>>) -> Result<Json<Value>, ApiError> {
    let loader = state.session_manager.skills_loader();
    let skills: Vec<Value> = loader.get_loaded_skills().iter().map(|s| json!({
        "name": s.name,
        "description": s.description,
        "always": s.always,
        "source": loader.source_of(s.path.as_deref()),
    })).collect();
    Ok(Json(json!({ "skills": skills })))
}
  • Step 4: 注册路由 — protected router 加 .route("/api/skills", routing::get(http::get_skills))
  • Step 5: 验证cargo build + cargo test --lib + cargo clippy -- -D warnings
  • Step 6: Commitgit add src/session/session.rs src/skills/mod.rs src/gateway && git commit -m "feat(gateway): GET /api/skills"

Chunk 3: 前端

Task 3.1: 可视化组件Sparkline / CapacityMeter / MetricTile

Files:

  • Create: webui/src/lib/components/Sparkline.svelteCapacityMeter.svelteMetricTile.svelte

  • Step 1: Sparkline.svelte — props values: number[]color: string(默认 var(--signal))。手写 SVG归一化 values 到 viewBox画折线/柱。无第三方库。

  • Step 2: CapacityMeter.svelte — props depth: numbercap: numbersegments?: number(默认 8。按 depth/cap 比例点亮分段块;接近满(>90%)用 var(--accent)/var(--danger),否则 var(--signal)

  • Step 3: MetricTile.svelte — props labelvaluesub(可选小字)、sparkValues/sparkColor(可选)。大等宽数字(var(--font-mono)+ label-caps + 可选 sparkline。用 .panel 容器。

  • Step 4: 验证cd webui && npm run check && npm run build

  • Step 5: Commitgit add webui/src/lib/components && git commit -m "feat(webui): sparkline, capacity meter, metric tile components"

Task 3.2: 概览页(轮询 /api/status

Files:

  • Create: webui/src/pages/OverviewPage.svelte
  • Modify: webui/src/lib/api.js(如需 status 辅助;现有 api() 通用函数已够,可复用)

布局(规格 §6.2,参考 .superpowers/brainstorm/111044-1784795642/pages-runtime.html mockup

  • 主状态条:RUNNINGphase=active 时青绿脉冲、运行代、uptime格式化为 Xd Xh、版本、WS 连接、后台任务、上次重载相位。

  • 指标块行MetricTile会话数、今日 Tokenin+out+cost sub、工具调用、今日 Turns+p95 sub

  • Provider 表名称、模型、状态徽标ok=signal/degraded=accent、延迟 Sparkline、token、cost。

  • 消息总线inbound/outbound/control 三个 CapacityMeter + 活跃 lane 数 + 调度器状态 + MCP 连接数。

  • 渠道状态列表name + 状态徽标。

  • 调度器jobs/enabled/failed_jobs/next_run格式化为倒计时或时间

  • 数据:onMountsetInterval 每 2s api("/api/status");卸载清除。用 $state 持有 status$derived 派生展示值。错误时显示离线态。

  • Step 1: 实现 OverviewPage.svelte(轮询 + 上述布局,套用 Signal Deck tokens 与 P0 组件)。

  • Step 2: 验证npm run check && npm run build

  • Step 3: Commitgit add webui/src/pages/OverviewPage.svelte && git commit -m "feat(webui): overview runtime dashboard"

Task 3.3: 工具 & Skills 页

Files:

  • Create: webui/src/pages/ToolsPage.svelte

布局(规格 §6.3,参考 .superpowers/brainstorm/111044-1784795642/page-tools-v2.html mockup

  • 三标签:工具 / Skills / MCP。

  • 工具标签:搜索框 + 能力筛选 chips全部/只读/可并发/有副作用/独占)+ 图例工具卡片网格名称、来源徽标builtin/mcp、能力徽标◇只读=signal / ⇉可并发=info / △有副作用=accent / ■独占=danger、调用次数、描述、可展开 <details> 参数 schema<pre> mono

  • Skills 标签名称、描述、always 徽标、来源。

  • MCP 标签:服务器名 + 连接状态徽标 + 工具数。

  • 数据:onMountapi("/api/tools") + api("/api/skills")MCP 从 api("/api/status")mcp 字段取一次即可MCP 状态为连接时快照)。能力筛选与搜索为前端纯派生($derived 过滤)。

  • Step 1: 实现 ToolsPage.svelte(三标签 + 搜索 + 能力筛选 + 卡片)。

  • Step 2: 验证npm run check && npm run build

  • Step 3: Commitgit add webui/src/pages/ToolsPage.svelte && git commit -m "feat(webui): tools and skills browser"

Task 3.4: App.svelte 接线(替换占位)

Files:

  • Modify: webui/src/App.svelte

  • Step 1: 接线 — import OverviewPageToolsPage;把 overview/tools{:else} 占位(即将上线)替换为对应页面渲染({:else if current === "overview"}<OverviewPage /> {:else if current === "tools"}<ToolsPage />)。保留其余页面分支。

  • Step 2: 验证npm run check && npm run build

  • Step 3: 目检(可选,需 gateway— 概览页显示实时指标、工具页列出工具/Skills/MCP、亮暗主题正常。

  • Step 4: Commitgit add webui/src/App.svelte && git commit -m "feat(webui): wire overview and tools pages"


Chunk 4: 收尾

Task 4.1: P1 收尾验证 + 版本号

  • 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.4.0 → 1.5.0)。检查 README 是否有需更新的观测能力描述(如有则最小化更新)。
  • Step 3: Commitgit add -A && git commit -m "chore(release): P1 observability"(仅暂存版本号文件 + 可能的 README确认无构建产物

P1 完成标志

  • GET /api/status 返回完整运行快照generation/uptime/sessions/metrics/bus/providers/channels/scheduler/mcp/ws_connections/background_tasks设备鉴权保护。
  • GET /api/tools 返回工具 + 能力字段read_only/exclusive/concurrency_safe+ call_count + source。
  • GET /api/skills 返回 name/description/always/source不含 content
  • 概览页每 2s 轮询并渲染仪表盘;工具页三标签 + 搜索 + 能力筛选。
  • Metrics 进程级、跨重载存活、纯内存。
  • npm run checknpm run buildcargo buildcargo test --libcargo clippy -- -D warnings 全绿。

风险与开放项

  • cost 暂未接入单价AgentLoop 埋点 cost 传 None拿不到 provider config 单价),故 /api/status 的 cost 恒为 0除非后续把 (price_in, price_out) 传入 AgentLoop。可选价格字段已加Task 1.2),接线留待需要时补。若评审认为 P1 必须出 cost则在 Task 1.3 把单价随 AgentLoop 构造传入(需从 Session 持有的 provider config 取)。
  • failed_jobs 为近似:用 last_status 而非严格 7 天窗口(避免 2s 轮询逐任务查运行记录。UI 需标注语义。
  • active_lanes 计数:依赖 dispatcher lane 生成/结束正确 +1/-1多 break 路径需保证减计数(用 guard 或单一出口)。
  • Metrics 全局态选择进程级全局precedented by MCP非穿线若评审偏好显式依赖注入可改为经 SessionManagerServices 穿线成本5 个 AgentLoop 构造点 + 多处结构体字段)。
  • provider status 派生阈值:最近 10 次 ≥3 失败 → degraded阈值可调。