Compare commits

...

8 Commits

16 changed files with 1438 additions and 9 deletions

View File

@ -0,0 +1,599 @@
# 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/tools``GET /api/skills` 只读端点,以及前端概览页(运行仪表盘)与工具&Skills 页。
**Architecture:** 新增进程级 `Metrics``OnceLock<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. **uptime**`src/gateway/mod.rs` 增加进程级 `static STARTED: OnceLock<std::time::Instant>` + `fn process_uptime_secs() -> u64`,跨重载存活。
4. **active_lanes**`OutboundDispatcher` 增加 `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-24)
**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/status``b14edc4`, with review fixes in `87cf500`
**Current verification baseline:** `cargo test --lib` 349 passed; `cargo clippy --all-targets --all-features -- -D warnings` and `cargo build` passed after Task 2.2.
**Resume at:** Task 2.3, protected `GET /api/tools`. No Task 2.3 code has been started. Continue with subagent-driven development and run separate spec-compliance and code-quality reviews before marking it complete.
**Restore on another computer:**
```bash
git fetch origin
git switch --track origin/feat/webui-p1
```
---
## Chunk 1: Metrics 后端
### Task 1.1: Metrics 模块(结构体 + 记录/快照 + 全局访问器 + 单测)
**Files:**
- Create: `src/observability/metrics.rs`
- Modify: `src/observability/mod.rs``pub 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。示例
```rust
#[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**
```rust
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: clippy**`cargo clippy --all-targets --all-features -- -D warnings`
- [ ] **Step 6: Commit**`git add src/observability && git commit -m "feat(observability): process-global Metrics collector"`
### Task 1.2: 配置可选价格字段
**Files:**
- Modify: `src/config/mod.rs``LLMProviderConfig` 结构体)
- [ ] **Step 1: 加可选字段** — 在 `LLMProviderConfig` 增加serde 可选,缺省 None不影响现有配置解析
```rust
#[serde(default)]
pub price_input_per_million: Option<f64>,
#[serde(default)]
pub price_output_per_million: Option<f64>,
```
- [ ] **Step 2: 加 cost 辅助方法**(在 LLMProviderConfig impl
```rust
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: Commit**`git 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-provider`stream_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_usage`turn 延迟用进入 `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: Commit**`git 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.rs`MessageBus 队列深度)
- Modify: `src/bus/dispatcher.rs`active lane 计数)
- Modify: `src/task_supervisor.rs`(运行任务数)
- Modify: `src/session/session.rs`(会话数 + 活动 turn 数)
- Modify: `src/gateway/ws.rs`ws 连接计数)
- Modify: `src/gateway/mod.rs`uptime 静态 + ws_connections 字段 + dispatcher lane 计数注入)
- [ ] **Step 1: MessageBus 队列深度** — 在 `src/bus/mod.rs` 加:
```rust
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` 加:
```rust
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 的判定):
```rust
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
}
```
`inner``Arc<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.rs``handle_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: Commit**`git 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.rs`handler
- Modify: `src/gateway/mod.rs`(注册路由)
- [ ] **Step 1: handler** — 在 `src/gateway/http.rs` 加(聚合所有来源;用 `serde_json::json!`
```rust
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<_>>(),
})))
}
```
`ReloadPhase``Serialize`snake_case`McpServerStatus` 字段见 `src/mcp/mod.rs:27-34``state.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 -- gateway``curl` 需鉴权,目检 JSON 结构;或写一个轻量集成断言。无鉴权 token 时跳过 curl靠编译 + 结构审查。)
- [ ] **Step 4: Commit**`git 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)`
```rust
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: Commit**`git 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: 暴露 SkillsLoader**`SessionManager` 私有字段 `skills_loader: Arc<SkillsLoader>``session.rs:1423`)无公开访问器。加:
```rust
pub fn skills_loader(&self) -> Arc<crate::skills::SkillsLoader> {
self.skills_loader.clone()
}
```
- [ ] **Step 2: source 派生辅助**`Skill.path` 是技能目录source 由路径前缀派生(`~/.agents/skills``agent``~/.picobot/skills``picobot``{workspace}/skills``workspace`)。由于三个目录字段是 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`
```rust
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: Commit**`git 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.svelte``CapacityMeter.svelte``MetricTile.svelte`
- [ ] **Step 1: Sparkline.svelte** — props `values: number[]``color: string`(默认 `var(--signal)`)。手写 SVG归一化 values 到 viewBox画折线/柱。无第三方库。
- [ ] **Step 2: CapacityMeter.svelte** — props `depth: number``cap: number``segments?: number`(默认 8。按 depth/cap 比例点亮分段块;接近满(>90%)用 `var(--accent)`/`var(--danger)`,否则 `var(--signal)`
- [ ] **Step 3: MetricTile.svelte** — props `label``value``sub`(可选小字)、`sparkValues`/`sparkColor`(可选)。大等宽数字(`var(--font-mono)`+ label-caps + 可选 sparkline。用 `.panel` 容器。
- [ ] **Step 4: 验证**`cd webui && npm run check && npm run build`
- [ ] **Step 5: Commit**`git 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
- 主状态条:`RUNNING`phase=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格式化为倒计时或时间
- 数据:`onMount``setInterval` 每 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: Commit**`git 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 标签:服务器名 + 连接状态徽标 + 工具数。
- 数据:`onMount``api("/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: Commit**`git 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 `OverviewPage``ToolsPage`;把 `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: Commit**`git 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.toml``webui/package.json` minor bump1.4.0 → 1.5.0)。检查 README 是否有需更新的观测能力描述(如有则最小化更新)。
- [ ] **Step 3: Commit**`git 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 check``npm run build``cargo build``cargo test --lib``cargo 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阈值可调。

View File

@ -491,16 +491,41 @@ impl AgentLoop {
iteration: u32,
turn: Option<&AgentTurnContext>,
) -> Result<ChatCompletionResponse, AgentError> {
let mut provider_stream = self.provider.stream(request).await.map_err(|error| {
let metrics = crate::observability::metrics::global_metrics();
let provider_name = self.provider.name().to_string();
let provider_model = self.provider.model_id().to_string();
let start = Instant::now();
let mut provider_stream = match self.provider.stream(request).await {
Ok(stream) => stream,
Err(error) => {
tracing::error!(error = %error, "LLM request failed");
AgentError::LlmError(error.to_string())
})?;
metrics.record_provider(
&provider_name,
&provider_model,
None,
start.elapsed().as_millis() as u64,
true,
);
return Err(AgentError::LlmError(error.to_string()));
}
};
let mut accumulator = ProviderResponseAccumulator::default();
while let Some(chunk) = provider_stream.next().await {
let chunk = chunk.map_err(|error| {
let chunk = match chunk {
Ok(chunk) => chunk,
Err(error) => {
tracing::error!(error = %error, "LLM stream failed");
AgentError::LlmError(error.to_string())
})?;
metrics.record_provider(
&provider_name,
&provider_model,
None,
start.elapsed().as_millis() as u64,
true,
);
return Err(AgentError::LlmError(error.to_string()));
}
};
if let Some(turn) = turn {
let event = match &chunk {
ProviderChunk::Reasoning(delta) => Some(TurnEvent::ReasoningDelta {
@ -521,7 +546,11 @@ impl AgentLoop {
}
accumulator.push(chunk);
}
Ok(accumulator.finish())
let response = accumulator.finish();
let latency_ms = start.elapsed().as_millis() as u64;
metrics.record_provider(&provider_name, &provider_model, None, latency_ms, false);
metrics.record_provider_tokens(&provider_name, &response.usage);
Ok(response)
}
fn annotate_message(
@ -615,6 +644,8 @@ impl AgentLoop {
mut messages: Vec<ChatMessage>,
turn: Option<AgentTurnContext>,
) -> Result<AgentProcessResult, AgentError> {
let turn_start = Instant::now();
#[cfg(debug_assertions)]
tracing::debug!(
history_len = messages.len(),
@ -703,6 +734,10 @@ impl AgentLoop {
assistant_message.provider_state = response.provider_state;
Self::annotate_message(&mut assistant_message, turn.as_ref(), iteration, true);
emitted_messages.push(assistant_message.clone());
crate::observability::metrics::global_metrics().record_turn(
Some(&accumulated_usage),
turn_start.elapsed().as_millis() as u64,
);
return Ok(AgentProcessResult {
final_response: assistant_message,
emitted_messages,
@ -844,6 +879,10 @@ impl AgentLoop {
true,
);
emitted_messages.push(assistant_message.clone());
crate::observability::metrics::global_metrics().record_turn(
Some(&accumulated_usage),
turn_start.elapsed().as_millis() as u64,
);
Ok(AgentProcessResult {
final_response: assistant_message,
emitted_messages,
@ -876,6 +915,11 @@ impl AgentLoop {
let mut final_message = ChatMessage::assistant(fallback);
Self::annotate_message(&mut final_message, turn.as_ref(), summary_iteration, true);
emitted_messages.push(final_message.clone());
let turn_usage = (accumulated_usage.total_tokens > 0).then_some(&accumulated_usage);
crate::observability::metrics::global_metrics().record_turn(
turn_usage,
turn_start.elapsed().as_millis() as u64,
);
Ok(AgentProcessResult {
final_response: final_message,
emitted_messages,
@ -1016,6 +1060,9 @@ impl AgentLoop {
});
}
crate::observability::metrics::global_metrics()
.record_tool_call(&tool_name, result.success);
// Apply duration
Ok(ToolExecutionOutcome { duration, ..result })
}

View File

@ -832,6 +832,8 @@ mod tests {
token_limit: 4096,
workspace_dir: std::env::temp_dir(),
input_types: vec!["text".into()],
price_input_per_million: None,
price_output_per_million: None,
},
Arc::new(ToolRegistry::new()),
None,

View File

@ -1,5 +1,6 @@
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use tokio::sync::mpsc;
@ -22,6 +23,7 @@ pub struct OutboundDispatcher {
channel_manager: ChannelManager,
task_supervisor: TaskSupervisor,
write_locks: ConversationWriteLocks,
active_lanes: Arc<AtomicUsize>,
}
impl OutboundDispatcher {
@ -30,12 +32,14 @@ impl OutboundDispatcher {
channel_manager: ChannelManager,
task_supervisor: TaskSupervisor,
write_locks: ConversationWriteLocks,
active_lanes: Arc<AtomicUsize>,
) -> Self {
Self {
bus,
channel_manager,
task_supervisor,
write_locks,
active_lanes,
}
}
@ -125,9 +129,14 @@ impl OutboundDispatcher {
chat_id: String,
) -> bool {
let target_lock = self.write_locks.for_target(&channel_name, &chat_id);
let active_lanes = self.active_lanes.clone();
self.task_supervisor.spawn(
format!("outbound-lane:{channel_name}:{chat_id}"),
async move {
active_lanes.fetch_add(1, Ordering::Relaxed);
let _guard = LaneGuard {
counter: active_lanes,
};
loop {
let msg = match tokio::time::timeout(LANE_IDLE_TIMEOUT, receiver.recv()).await {
Ok(Some(msg)) => msg,
@ -180,6 +189,18 @@ impl OutboundDispatcher {
}
}
/// Decrements the active-lane counter exactly once when a lane task ends,
/// whether it exits normally, is cancelled, or is aborted.
struct LaneGuard {
counter: Arc<AtomicUsize>,
}
impl Drop for LaneGuard {
fn drop(&mut self) {
self.counter.fetch_sub(1, Ordering::Relaxed);
}
}
#[cfg(test)]
mod tests {
use super::*;
@ -275,6 +296,7 @@ mod tests {
manager,
supervisor.clone(),
ConversationWriteLocks::default(),
Arc::new(AtomicUsize::new(0)),
);
let task = tokio::spawn(async move { dispatcher.run().await });
bus.publish_outbound(outbound("slow", "slow-1"))
@ -318,6 +340,7 @@ mod tests {
manager,
supervisor.clone(),
ConversationWriteLocks::default(),
Arc::new(AtomicUsize::new(0)),
);
let task = tokio::spawn(async move { dispatcher.run().await });
@ -348,6 +371,7 @@ mod tests {
manager,
supervisor.clone(),
ConversationWriteLocks::default(),
Arc::new(AtomicUsize::new(0)),
);
let task = tokio::spawn(async move { dispatcher.run().await });
@ -360,6 +384,50 @@ mod tests {
supervisor.shutdown(Duration::from_secs(1)).await;
}
#[tokio::test]
async fn active_lane_count_tracks_task_lifetime() {
let bus = MessageBus::new(8);
let manager = ChannelManager::with_bus(
Arc::new(crate::channels::CliChatChannel::new()),
bus.clone(),
);
manager
.register_channel(
"recording",
Arc::new(RecordingChannel {
sent: Mutex::new(Vec::new()),
notify: Notify::new(),
}),
)
.await;
let supervisor = TaskSupervisor::new();
let active_lanes = Arc::new(AtomicUsize::new(0));
let dispatcher = OutboundDispatcher::new(
bus.clone(),
manager,
supervisor.clone(),
ConversationWriteLocks::default(),
active_lanes.clone(),
);
let dispatcher_task = tokio::spawn(async move { dispatcher.run().await });
bus.publish_outbound(outbound("counted", "message"))
.await
.unwrap();
tokio::time::timeout(Duration::from_secs(1), async {
while active_lanes.load(Ordering::Relaxed) != 1 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
supervisor.shutdown(Duration::from_secs(1)).await;
assert_eq!(active_lanes.load(Ordering::Relaxed), 0);
dispatcher_task.abort();
let _ = dispatcher_task.await;
}
#[tokio::test]
async fn permanent_send_failure_is_not_retried() {
let channel = PermanentFailureChannel {

View File

@ -106,6 +106,29 @@ impl MessageBus {
pub async fn consume_control(&self) -> Option<ControlMessage> {
self.control_rx.lock().await.recv().await
}
/// Snapshot of the current depth and capacity of each bus queue.
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,
}
}
}
/// Read-only snapshot of MessageBus queue utilization.
#[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,
}
// ============================================================================
@ -130,3 +153,38 @@ impl std::fmt::Display for BusError {
}
impl std::error::Error for BusError {}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
#[tokio::test]
async fn queue_depths_report_retained_sender_usage() {
let bus = MessageBus::new(3);
let empty = bus.queue_depths();
assert_eq!(empty.inbound_depth, 0);
assert_eq!(empty.inbound_cap, 3);
assert_eq!(empty.outbound_depth, 0);
assert_eq!(empty.outbound_cap, 3);
assert_eq!(empty.control_depth, 0);
assert_eq!(empty.control_cap, 3);
bus.publish_outbound(OutboundMessage {
channel: "test".to_string(),
chat_id: "chat".to_string(),
content: "queued".to_string(),
reply_to: None,
media: vec![],
metadata: HashMap::new(),
delivery: None,
})
.await
.unwrap();
assert_eq!(bus.queue_depths().outbound_depth, 1);
bus.consume_outbound().await.unwrap();
assert_eq!(bus.queue_depths().outbound_depth, 0);
}
}

View File

@ -508,6 +508,19 @@ pub struct LLMProviderConfig {
pub token_limit: usize,
pub workspace_dir: PathBuf,
pub input_types: Vec<String>,
pub price_input_per_million: Option<f64>,
pub price_output_per_million: Option<f64>,
}
impl LLMProviderConfig {
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,
}
}
}
pub fn get_default_config_path() -> PathBuf {
@ -653,6 +666,8 @@ impl Config {
token_limit: agent.token_limit,
workspace_dir: expand_path(&self.workspace_dir),
input_types: model.input_type.clone(),
price_input_per_million: None,
price_output_per_million: None,
})
}
}

View File

@ -169,6 +169,14 @@ impl ApiError {
message: error.to_string(),
}
}
fn internal_with_message(error: impl std::fmt::Display, message: impl Into<String>) -> Self {
tracing::error!(error = %error, "WebUI API request failed");
Self {
status: StatusCode::INTERNAL_SERVER_ERROR,
message: message.into(),
}
}
}
pub async fn upload_file(
@ -689,6 +697,129 @@ pub struct LimitQuery {
limit: Option<usize>,
}
fn scheduler_snapshot(jobs: &[crate::storage::ScheduledJob]) -> Value {
let enabled = jobs.iter().filter(|job| job.enabled).count();
let failed_jobs = jobs
.iter()
.filter(|job| {
matches!(
job.last_status.as_deref(),
Some("error" | "timeout" | "delivery_error")
)
})
.count();
let next_run_at = jobs
.iter()
.filter(|job| job.enabled)
.map(|job| job.next_run_at)
.min();
json!({
"jobs": jobs.len(),
"enabled": enabled,
"failed_jobs": failed_jobs,
"next_run_at": next_run_at,
})
}
fn scheduler_status_error(error: impl std::fmt::Display) -> ApiError {
ApiError::internal_with_message(error, "failed to load scheduler status")
}
fn channel_snapshot(mut channels: Vec<(String, bool)>) -> Vec<Value> {
channels.sort_by(|left, right| left.0.cmp(&right.0));
channels
.into_iter()
.map(|(name, running)| {
json!({
"name": name,
"status": if running { "connected" } else { "stopped" },
})
})
.collect()
}
fn sorted_providers(
mut providers: Vec<crate::observability::metrics::ProviderSnapshot>,
) -> Vec<crate::observability::metrics::ProviderSnapshot> {
providers.sort_by(|left, right| left.name.cmp(&right.name));
providers
}
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 channel_states = Vec::new();
for name in state.channel_manager.list_channel_names().await {
let running = state
.channel_manager
.get_channel(&name)
.await
.is_some_and(|channel| channel.is_running());
channel_states.push((name, running));
}
let channels = channel_snapshot(channel_states);
let jobs = state
.storage
.list_scheduled_jobs()
.await
.map_err(scheduler_status_error)?;
let scheduler = scheduler_snapshot(&jobs);
let mcp = crate::mcp::get_mcp_status()
.into_iter()
.map(|status| {
json!({
"name": status.name,
"connected": status.connected,
"tools": status.tools.len(),
})
})
.collect::<Vec<_>>();
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": sorted_providers(metrics.providers),
"channels": channels,
"scheduler": scheduler,
"mcp": mcp,
})))
}
pub async fn get_tasks(
State(state): State<Arc<GatewayState>>,
Query(query): Query<LimitQuery>,
@ -779,6 +910,117 @@ pub async fn get_memories(
mod tests {
use super::*;
fn scheduled_job(
id: &str,
enabled: bool,
next_run_at: i64,
last_status: Option<&str>,
) -> crate::storage::ScheduledJob {
crate::storage::ScheduledJob {
id: id.to_string(),
name: id.to_string(),
schedule: crate::scheduler::Schedule::Every { every_ms: 60_000 },
prompt: String::new(),
channel: "cli_chat".to_string(),
chat_id: "test".to_string(),
model: None,
job_kind: crate::storage::JobKind::Task,
delivery_policy: crate::storage::DeliveryPolicy::Never,
enabled,
delete_after_run: false,
next_run_at,
last_run_at: None,
last_status: last_status.map(str::to_string),
last_error: None,
created_at: 0,
updated_at: 0,
}
}
#[test]
fn scheduler_snapshot_classifies_failures_and_next_enabled_run() {
let jobs = vec![
scheduled_job("healthy", true, 300, Some("ok")),
scheduled_job("error", true, 200, Some("error")),
scheduled_job("timeout", false, 100, Some("timeout")),
scheduled_job("delivery", true, 400, Some("delivery_error")),
scheduled_job("other", false, 50, Some("cancelled")),
];
assert_eq!(
scheduler_snapshot(&jobs),
json!({
"jobs": 5,
"enabled": 3,
"failed_jobs": 3,
"next_run_at": 200,
})
);
}
#[test]
fn scheduler_snapshot_has_null_next_run_without_enabled_jobs() {
let jobs = vec![scheduled_job("disabled", false, 100, None)];
assert_eq!(
scheduler_snapshot(&jobs),
json!({
"jobs": 1,
"enabled": 0,
"failed_jobs": 0,
"next_run_at": null,
})
);
}
#[test]
fn channel_snapshot_is_sorted_and_maps_running_state() {
let channels = channel_snapshot(vec![
("zeta".to_string(), false),
("alpha".to_string(), true),
]);
assert_eq!(
channels,
vec![
json!({ "name": "alpha", "status": "connected" }),
json!({ "name": "zeta", "status": "stopped" }),
]
);
}
#[test]
fn provider_snapshot_is_sorted_by_name() {
let provider = |name: &str| crate::observability::metrics::ProviderSnapshot {
name: name.to_string(),
model: String::new(),
status: "ok".to_string(),
latency_ms: 0,
latencies: Vec::new(),
tokens_in: 0,
tokens_out: 0,
cost: 0.0,
};
let providers = sorted_providers(vec![provider("zeta"), provider("alpha")]);
assert_eq!(
providers
.iter()
.map(|provider| provider.name.as_str())
.collect::<Vec<_>>(),
vec!["alpha", "zeta"]
);
}
#[test]
fn scheduler_status_error_hides_internal_details() {
let error = scheduler_status_error("database contained private payload");
assert_eq!(error.status, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(error.message, "failed to load scheduler status");
assert!(!error.message.contains("private payload"));
}
#[test]
fn secrets_are_redacted_and_restored() {
let current = json!({"api_key":"real", "nested":{"access_token":"token"}, "safe":"yes"});

View File

@ -8,6 +8,7 @@ pub mod ws;
use axum::{Router, middleware, routing};
use std::net::SocketAddr;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
use tokio::net::TcpListener;
use crate::bus::{MessageBus, OutboundDispatcher};
@ -21,6 +22,18 @@ use crate::scheduler::Scheduler;
use crate::session::{SessionManager, SessionManagerServices};
use crate::task_supervisor::TaskSupervisor;
/// Process boot clock. A process-level static so uptime survives config reload,
/// which swaps GatewayState generations without restarting the process.
static STARTED: std::sync::OnceLock<std::time::Instant> = std::sync::OnceLock::new();
/// Seconds elapsed since the gateway process started.
pub fn process_uptime_secs() -> u64 {
STARTED
.get_or_init(std::time::Instant::now)
.elapsed()
.as_secs()
}
pub struct GatewayState {
pub config: Config,
pub config_path: std::path::PathBuf,
@ -33,6 +46,10 @@ pub struct GatewayState {
pub connection_shutdown: tokio_util::sync::CancellationToken,
pub auth: auth::AuthManager,
pub uploads: uploads::UploadRegistry,
/// Live WebSocket connection count.
pub ws_connections: Arc<AtomicUsize>,
/// Active outbound dispatcher lane count (shared with the dispatcher).
pub outbound_lanes: Arc<AtomicUsize>,
pub(crate) reload: reload::ReloadHandle,
pub(crate) admission: reload::RuntimeAdmission,
}
@ -231,6 +248,8 @@ impl GatewayState {
connection_shutdown,
auth,
uploads,
ws_connections: Arc::new(AtomicUsize::new(0)),
outbound_lanes: Arc::new(AtomicUsize::new(0)),
reload,
admission,
})
@ -312,6 +331,7 @@ impl GatewayState {
self.channel_manager.clone(),
self.task_supervisor.clone(),
self.delivery_coordinator.write_locks(),
self.outbound_lanes.clone(),
);
self.task_supervisor
@ -341,6 +361,7 @@ pub async fn run(
host: Option<String>,
port: Option<u16>,
) -> Result<(), Box<dyn std::error::Error>> {
STARTED.get_or_init(std::time::Instant::now);
let config_path = crate::config::resolve_default_config_path();
let startup_process_env = Config::startup_process_env();
let startup_cwd = std::env::current_dir()?;
@ -572,6 +593,7 @@ fn build_router(state: Arc<GatewayState>) -> Router {
)
.route("/api/logs", routing::get(http::get_logs))
.route("/api/tasks", routing::get(http::get_tasks))
.route("/api/status", routing::get(http::get_status))
.route("/api/jobs", routing::get(http::get_jobs))
.route("/api/jobs/{id}/runs", routing::get(http::get_job_runs))
.route("/api/memories", routing::get(http::get_memories))

View File

@ -7,6 +7,7 @@ use axum::response::Response;
use futures_util::{SinkExt, StreamExt};
use serde::Deserialize;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::mpsc;
use tokio::time::{Duration, timeout};
@ -42,6 +43,8 @@ async fn handle_socket(
client_id: Option<String>,
identity: super::auth::AuthIdentity,
) {
let _connection_guard = ConnectionGuard::new(state.ws_connections.clone());
// Create channel for sending outbound messages to this client
let (sender, mut receiver) = mpsc::channel::<WsOutbound>(100);
@ -131,6 +134,23 @@ async fn handle_socket(
tracing::info!(session_id = %session_id, "CLI session ended");
}
struct ConnectionGuard {
counter: Arc<AtomicUsize>,
}
impl ConnectionGuard {
fn new(counter: Arc<AtomicUsize>) -> Self {
counter.fetch_add(1, Ordering::Relaxed);
Self { counter }
}
}
impl Drop for ConnectionGuard {
fn drop(&mut self) {
self.counter.fetch_sub(1, Ordering::Relaxed);
}
}
#[cfg(test)]
mod tests {
use super::*;
@ -145,4 +165,14 @@ mod tests {
assert!(valid_client_id(Some("x".repeat(65))).is_none());
assert!(valid_client_id(Some(String::new())).is_none());
}
#[test]
fn connection_guard_tracks_its_scope() {
let connections = Arc::new(AtomicUsize::new(0));
{
let _guard = ConnectionGuard::new(connections.clone());
assert_eq!(connections.load(Ordering::Relaxed), 1);
}
assert_eq!(connections.load(Ordering::Relaxed), 0);
}
}

View File

@ -0,0 +1,295 @@
use std::collections::{HashMap, VecDeque};
use std::sync::atomic::{AtomicU64, Ordering::Relaxed};
use std::sync::{Arc, Mutex, OnceLock};
use crate::providers::Usage;
const WINDOW: usize = 100;
const DEGRADE_THRESHOLD: usize = 3;
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>,
}
pub struct Metrics {
tokens_in: AtomicU64,
tokens_out: AtomicU64,
turns: AtomicU64,
tool_calls: AtomicU64,
per_tool: Mutex<HashMap<String, u64>>,
turn_latencies: Mutex<VecDeque<u64>>,
providers: Mutex<HashMap<String, ProviderStat>>,
}
#[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,
pub latency_ms: u64,
pub latencies: Vec<u64>,
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: AtomicU64::new(0),
tokens_out: AtomicU64::new(0),
turns: AtomicU64::new(0),
tool_calls: AtomicU64::new(0),
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) {
self.turns.fetch_add(1, Relaxed);
if let Some(u) = usage {
self.tokens_in.fetch_add(u64::from(u.prompt_tokens), Relaxed);
self.tokens_out
.fetch_add(u64::from(u.completion_tokens), 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) {
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 += u64::from(usage.prompt_tokens);
stat.tokens_out += u64::from(usage.completion_tokens);
}
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 {
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".to_string()
} else {
"ok".to_string()
},
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();
pub fn global_metrics() -> Arc<Metrics> {
GLOBAL.get_or_init(|| Arc::new(Metrics::new())).clone()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn records_turns_tokens_and_p95() {
let m = Metrics::new();
for i in 1..=100u64 {
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);
}
#[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 records_provider_tokens_and_cost() {
let m = Metrics::new();
m.record_provider("openai", "gpt-4o", Some(0.05), 120, false);
m.record_provider_tokens(
"openai",
&Usage {
prompt_tokens: 100,
completion_tokens: 50,
total_tokens: 150,
..Default::default()
},
);
let s = m.snapshot();
assert_eq!(s.cost, 0.05);
let p = s.providers.iter().find(|p| p.name == "openai").unwrap();
assert_eq!(p.tokens_in, 100);
assert_eq!(p.tokens_out, 50);
assert_eq!(p.latency_ms, 120);
assert_eq!(p.status, "ok");
}
#[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");
}
#[test]
fn tool_call_count_accessor() {
let m = Metrics::new();
assert_eq!(m.tool_call_count("bash"), 0);
m.record_tool_call("bash", true);
assert_eq!(m.tool_call_count("bash"), 1);
}
#[test]
fn global_metrics_returns_same_instance() {
let a = global_metrics();
let b = global_metrics();
assert!(Arc::ptr_eq(&a, &b));
}
}

View File

@ -3,6 +3,8 @@
//! This module provides an Observer pattern for emitting and collecting
//! telemetry events during agent execution.
pub mod metrics;
use std::time::Duration;
use crate::bus::MediaRef;

View File

@ -225,6 +225,8 @@ mod tests {
token_limit: 8_192,
workspace_dir: PathBuf::from("."),
input_types: vec!["text".to_string(), "image".to_string()],
price_input_per_million: None,
price_output_per_million: None,
};
let session = Arc::new(Mutex::new(
Session::new(

View File

@ -432,6 +432,7 @@ mod cancelled_partial_tests {
channels,
supervisor.clone(),
ConversationWriteLocks::default(),
Arc::new(std::sync::atomic::AtomicUsize::new(0)),
);
let dispatcher_task = tokio::spawn(async move { dispatcher.run().await });
@ -2157,6 +2158,27 @@ impl SessionManager {
}
}
/// Number of live sessions currently tracked by the manager.
pub async fn session_count(&self) -> usize {
self.inner.lock().await.sessions.len()
}
/// Number of sessions with an actively executing Turn.
pub async fn active_turn_count(&self) -> usize {
let sessions: Vec<_> = {
let inner = self.inner.lock().await;
inner.sessions.values().cloned().collect()
};
let mut count = 0;
for session in sessions {
let session = session.lock().await;
if session.current_cancel.is_some() {
count += 1;
}
}
count
}
pub async fn create_session(
&self,
channel: &str,

View File

@ -114,6 +114,13 @@ impl TaskSupervisor {
true
}
/// Number of currently running managed tasks, pruning finished handles.
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()
}
pub fn cancel(&self) {
let mut state = self.inner.state.lock().unwrap_or_else(|e| e.into_inner());
state.stopping = true;
@ -199,4 +206,18 @@ mod tests {
assert!(cleaned_up.load(std::sync::atomic::Ordering::SeqCst));
}
#[tokio::test]
async fn running_count_prunes_finished_tasks() {
let supervisor = TaskSupervisor::new();
assert!(supervisor.spawn("pending", std::future::pending()));
assert_eq!(supervisor.running_count(), 1);
assert!(supervisor.spawn("finished", async {}));
tokio::task::yield_now().await;
assert_eq!(supervisor.running_count(), 1);
supervisor.shutdown(Duration::from_secs(1)).await;
assert_eq!(supervisor.running_count(), 0);
}
}

View File

@ -27,6 +27,8 @@ fn load_config() -> Option<LLMProviderConfig> {
token_limit: 128_000,
workspace_dir: std::path::PathBuf::from("/tmp/test-workspace"),
input_types: vec!["text".to_string()],
price_input_per_million: None,
price_output_per_million: None,
})
}

View File

@ -27,6 +27,8 @@ fn load_openai_config() -> Option<LLMProviderConfig> {
token_limit: 128_000,
workspace_dir: std::path::PathBuf::from("/tmp/test-workspace"),
input_types: vec!["text".to_string()],
price_input_per_million: None,
price_output_per_million: None,
})
}