Compare commits

..
4 changed files with 339 additions and 830 deletions

View file

@ -1473,6 +1473,245 @@ window.marked = marked;
}
return d;
}
// Generic args-pretty-printer for unknown / extra-MCP tools. The
// built-in switch handles the common claude/hyperhive tools; this
// is the fallback so an `mcp__matrix__send_message` or similar
// doesn't dump raw JSON. Heuristics: single string-valued field →
// `Name field: "value"`; single dict-valued field → `Name field
// {…}`; otherwise compact JSON. Always trimmed to fit a row.
function fmtArgsGeneric(name, input) {
const keys = Object.keys(input || {});
if (keys.length === 0) return name + '()';
if (keys.length === 1) {
const k = keys[0];
const v = input[k];
if (typeof v === 'string') {
const oneline = v.replace(/\s+/g, ' ').trim();
return name + ' ' + k + ': ' + JSON.stringify(trim(oneline, 100));
}
if (typeof v === 'number' || typeof v === 'boolean') {
return name + ' ' + k + ': ' + JSON.stringify(v);
}
}
// Multi-field: render `k: v` pairs with strings/numbers inlined and
// anything else summarised by type so the row stays readable.
const pretty = keys.slice(0, 4).map((k) => {
const v = input[k];
if (v == null) return k + ': null';
if (typeof v === 'string') {
const oneline = v.replace(/\s+/g, ' ').trim();
return k + ': ' + JSON.stringify(trim(oneline, 40));
}
if (typeof v === 'number' || typeof v === 'boolean') return k + ': ' + v;
if (Array.isArray(v)) return k + `: [${v.length}]`;
return k + ': {…}';
});
const tail = keys.length > 4 ? ' …+' + (keys.length - 4) : '';
return name + ' ' + pretty.join(' · ') + tail;
}
// Per-tool glyph for the tool-use row prefix — gives each hive / MCP /
// built-in tool a distinctive icon instead of a generic wrench, so the
// scrollback is scannable at a glance. Exact tool names first, then
// MCP-server family fallbacks, then a generic wrench default.
const TOOL_ICONS = {
'mcp__hyperhive__send': '📤',
'mcp__hyperhive__recv': '📥',
'mcp__hyperhive__ask': '❓',
'mcp__hyperhive__answer': '✍️',
'mcp__hyperhive__remind': '⏰',
'mcp__hyperhive__set_status': '🏷️',
'mcp__hyperhive__get_loose_ends': '🪢',
'mcp__hyperhive__cancel_loose_end': '✂️',
'mcp__hyperhive__ack_until': '✅',
'mcp__hyperhive__get_agent_meta': '',
'mcp__hyperhive__request_next_turn': '⏩',
'mcp__hyperhive__restart': '↻',
'mcp__hyperhive__kill': '⏹️',
'mcp__hyperhive__start': '▶️',
'mcp__hyperhive__update': '🔄',
'mcp__hyperhive__list_containers': '📋',
'mcp__hyperhive__get_logs': '📜',
'mcp__hyperhive__get_host_journal': '📜',
'mcp__matrix__read_room': '📖',
'mcp__matrix__mark_read': '👁️',
'mcp__matrix__list_rooms': '📋',
'mcp__matrix__list_room_members': '📋',
'mcp__matrix__list_invites': '📋',
'mcp__bash__kill': '🛑',
Read: '📖', Write: '💾', Edit: '✏️', Glob: '🔍', Grep: '🔍',
};
function toolIcon(name) {
if (TOOL_ICONS[name]) return TOOL_ICONS[name];
if (typeof name === 'string') {
if (name.startsWith('mcp__matrix__')) return '💬';
if (name.startsWith('mcp__bash__')) return '🖥️';
if (name.includes('schedule')) return '⏱️';
// request_init_config / request_update_meta_inputs
if (name.startsWith('mcp__hyperhive__request_')) return '📦';
}
return '🔧';
}
// Shorten a matrix room id or alias for display. Room ids (!xxx:server)
// are trimmed to the first 8 local chars; aliases (#name:server) are
// returned as-is (they're already readable). Falls back to the raw
// value truncated.
function fmtRoom(r) {
if (!r) return '?';
if (r.startsWith('!')) return r.slice(0, r.indexOf(':') > 0 ? r.indexOf(':') : 9);
if (r.startsWith('#')) return r.split(':')[0] || r;
return trim(r, 20);
}
// Shorten @user:server → @user.
function fmtUser(u) {
if (!u) return '?';
const colon = u.indexOf(':');
return colon > 0 ? u.slice(0, colon) : u;
}
// Pretty-print a tool call: per-known-tool format, fallback to JSON
// for unknown tools.
function fmtToolUse(c) {
const name = c.name || '';
const input = c.input || {};
const short = name.startsWith('mcp__hyperhive__')
? name.slice('mcp__hyperhive__'.length) + '*'
: name.startsWith('mcp__bash__')
? name.slice('mcp__bash__'.length) + '*'
: name.startsWith('mcp__matrix__')
? name.slice('mcp__matrix__'.length) + '*'
: name;
switch (name) {
case 'Read': return short + ' ' + (input.file_path || '');
case 'Write': return short + ' ' + (input.file_path || '');
case 'Edit': return short + ' ' + (input.file_path || '');
case 'Glob': return short + ' ' + (input.pattern || '');
case 'Grep': return short + ' ' + (input.pattern || '');
case 'Bash': return short + (input.run_in_background ? ' [bg]' : '')
+ ' $ ' + (input.command || '');
case 'TodoWrite': return short + ' (' + ((input.todos || []).length) + ' items)';
case 'mcp__hyperhive__send': return short + ' → ' + (input.to || '?') + ': '
+ JSON.stringify(input.body || '').slice(0, 80);
case 'mcp__hyperhive__recv': {
// Surface the long-poll wait + batch size — a bare `recv()` row
// hides whether the agent is parking a turn (wait_seconds) or
// draining a burst (max).
const parts = [];
if (input.wait_seconds != null) parts.push('wait ' + input.wait_seconds + 's');
if (input.max != null) parts.push('max ' + input.max);
return short + (parts.length ? ' ' + parts.join(' · ') : '()');
}
case 'mcp__hyperhive__kill': return short + ' ' + (input.name || '');
case 'mcp__hyperhive__restart': return short + ' ' + (input.name || '');
case 'mcp__hyperhive__start': return short + ' ' + (input.name || '');
case 'mcp__hyperhive__update': return short + ' ' + (input.name || '');
case 'mcp__hyperhive__ack_until':
return short + ' ≤' + (input.up_to != null ? input.up_to : '?');
case 'mcp__hyperhive__get_logs':
return short + ' ' + (input.agent || '?')
+ (input.lines != null ? ' · ' + input.lines + 'L' : '');
case 'mcp__hyperhive__get_host_journal': {
const parts = [];
if (input.container) parts.push(input.container);
else if (input.unit) parts.push(input.unit);
if (input.grep) parts.push('/' + input.grep + '/');
if (input.lines != null) parts.push(input.lines + 'L');
return short + (parts.length ? ' ' + parts.join(' · ') : '()');
}
case 'mcp__hyperhive__remind': {
// Surface when the reminder fires + the first line of the message.
// `delay_seconds` → human-readable "+5m"; `at_unix_timestamp` →
// "at HH:MM"; message truncated to fit.
let when = '';
if (input.delay_seconds != null) {
const s = input.delay_seconds;
when = '+' + (s < 60 ? s + 's' : s < 3600 ? Math.round(s / 60) + 'm'
: (s / 3600).toFixed(1) + 'h');
} else if (input.at_unix_timestamp != null) {
when = 'at ' + new Date(input.at_unix_timestamp * 1000)
.toISOString().slice(11, 16) + 'Z';
}
const msg = String(input.message || input.file_path || '').replace(/\s+/g, ' ').trim();
return short + (when ? ' ' + when : '') + (msg ? ' "' + trim(msg, 60) + '"' : '');
}
case 'mcp__hyperhive__request_init_config':
return short + ' ' + (input.name || '?');
case 'mcp__hyperhive__request_update_meta_inputs': {
const ins = Array.isArray(input.inputs) && input.inputs.length
? '[' + input.inputs.slice(0, 4).join(', ')
+ (input.inputs.length > 4 ? ', …' : '') + ']'
: 'all';
return short + ' ' + ins;
}
case 'mcp__hyperhive__list_schedules':
return short + '()';
case 'mcp__hyperhive__cancel_schedule':
return short + ' #' + (input.id != null ? input.id : '?')
+ (Array.isArray(input.targets) && input.targets.length
? ' [' + input.targets.join(', ') + ']' : ' all');
case 'mcp__hyperhive__fire_schedule_now':
return short + ' #' + (input.id != null ? input.id : '?');
case 'mcp__hyperhive__edit_schedule': {
const parts = ['#' + (input.id != null ? input.id : '?')];
if (input.body != null) parts.push('body');
if (input.interval_seconds != null) parts.push('interval');
if (input.next_fire_at_unix != null) parts.push('next');
if (input.targets_add && input.targets_add.length) parts.push('+' + input.targets_add.length + ' tgt');
if (input.targets_remove && input.targets_remove.length) parts.push('-' + input.targets_remove.length + ' tgt');
return short + ' ' + parts.join(' · ');
}
case 'mcp__hyperhive__request_schedule_prompt': {
const tgts = Array.isArray(input.targets) ? input.targets : [];
const when = input.first_fire_at_unix != null
? new Date(input.first_fire_at_unix * 1000).toISOString().slice(11, 16) + 'Z'
: '?';
return short + ' → ' + (tgts.length ? tgts.join(', ') : '?') + ' at ' + when
+ (input.interval_seconds != null ? ' +' + input.interval_seconds + 's' : '');
}
case 'mcp__bash__run': {
// Rich renderer handles the full body; this summary covers any
// fallback path and the details summary line.
const firstLine = String(input.cmd || '').split('\n')[0];
return short + ' $ ' + trim(firstLine.trim(), 72);
}
case 'mcp__bash__status':
return short + ' id:' + (input.id || '?')
+ (input.wait_seconds != null ? ' · wait ' + input.wait_seconds + 's' : '');
case 'mcp__bash__kill':
return short + ' ' + (input.id || '?') + (input.force ? ' [force]' : '');
case 'mcp__hyperhive__set_status':
return short + ' "' + trim(String(input.text || ''), 60) + '"';
case 'mcp__hyperhive__get_loose_ends':
return short + (input.agent ? ' [' + input.agent + ']' : '()');
case 'mcp__hyperhive__get_agent_meta':
return short + (input.name ? ' ' + input.name : '()');
case 'mcp__hyperhive__cancel_loose_end':
return short + ' ' + (input.kind || '?') + ' #' + (input.id != null ? input.id : '?');
case 'mcp__matrix__read_room':
return short + ' ' + fmtRoom(input.room)
+ (input.limit != null ? ' [' + input.limit + ']' : '');
case 'mcp__matrix__mark_read':
return short + ' ' + fmtRoom(input.room);
case 'mcp__matrix__send_message':
return short + ' → ' + fmtRoom(input.room) + ': '
+ JSON.stringify(trim(String(input.body || ''), 50));
case 'mcp__matrix__send_dm':
return short + ' → ' + fmtUser(input.user_id) + ': '
+ JSON.stringify(trim(String(input.body || ''), 50));
case 'mcp__matrix__send_reply':
return short + ' → ' + fmtRoom(input.room) + ': '
+ JSON.stringify(trim(String(input.body || ''), 50));
case 'mcp__matrix__send_reaction':
return short + ' ' + fmtRoom(input.room) + ' ' + (input.key || '?');
case 'mcp__matrix__join_room':
return short + ' ' + fmtRoom(input.room);
case 'mcp__matrix__open_dm':
return short + ' ' + fmtUser(input.user_id);
case 'mcp__matrix__invite_user':
return short + ' ' + fmtUser(input.user_id) + ' → ' + fmtRoom(input.room);
case 'mcp__matrix__download_file':
return short + ' ' + fmtRoom(input.room);
default: return fmtArgsGeneric(short, input);
}
}
// Build a "rich" tool_use row for tools whose input has a body we
// want the operator to see in full. Returns null for any other tool
// so the caller falls back to the flat-row path.
@ -1482,7 +1721,6 @@ window.marked = marked;
function renderRichToolUse(c, api) {
const name = c.name || '';
const input = c.input || {};
const icon = c._icon || '🔧';
if (name === 'Write' || name === 'Edit') {
const path = input.file_path || '?';
let body;
@ -1508,7 +1746,7 @@ window.marked = marked;
// the CSS disclosure caret leads the text).
const summary = name + ' ' + path + ' · '
+ (minus ? '-' + minus + ' ' : '') + '+' + plus;
return api.detailsDiff('tool-use', summary, body, icon);
return api.detailsDiff('tool-use', summary, body, toolIcon(name));
}
// Message-bearing tools render default-open with a markdown body so
// the operator sees the content without an extra click. send / ask
@ -1519,7 +1757,7 @@ window.marked = marked;
const lines = body.split('\n').length;
return detailsOpenMd(api, 'tool-use',
'send → ' + to + (lines > 1 ? ` · ${lines}L` : ''),
body, icon);
body, toolIcon(name));
}
if (name === 'mcp__hyperhive__ask') {
const to = input.to || 'operator';
@ -1527,7 +1765,7 @@ window.marked = marked;
const lines = q.split('\n').length;
const d = detailsOpenMd(api, 'tool-use',
'ask → ' + to + (lines > 1 ? ` · ${lines}L` : ''),
q, icon);
q, toolIcon(name));
// When the ask targets the operator, mount an inline answer
// slot in the live terminal — see docs/web-ui.md::Per-agent
// page (Ask → operator inline-answer binding) for the slot
@ -1557,7 +1795,7 @@ window.marked = marked;
const lines = a.split('\n').length;
return detailsOpenMd(api, 'tool-use',
'answer #' + id + (lines > 1 ? ` · ${lines}L` : ''),
a, icon);
a, toolIcon(name));
}
// Bash task runner — show full command in an expandable pre block so
// multi-line scripts are readable. Summary uses the first line so the
@ -1566,7 +1804,7 @@ window.marked = marked;
const cmd = String(input.cmd || '');
const firstLine = cmd.split('\n')[0];
const summary = 'run* $ ' + trim(firstLine.trim(), 72);
return api.details('tool-use', summary, '$ ' + cmd, icon);
return api.details('tool-use', summary, '$ ' + cmd, toolIcon(name));
}
return null;
}
@ -1684,28 +1922,101 @@ window.marked = marked;
// then "✓ done") for what's really one event.
const updatePluginInstall = makeCoalescer('note');
function renderStream(v, api) {
// Backend pre-computes `_category` on all known event types.
// "drop" covers: type=result, type=rate_limit_event, and system/init.
if (v._category === 'drop') return;
// Drop claude's result line and rate-limit — noise. TurnEnd
// communicates pass/fail; rate-limit events are noisy status chatter.
if (v.type === 'rate_limit_event') return;
if (v.type === 'result') return;
// `system` events: `init` is silent startup noise; `api_retry`
// and `api_error` get human-readable notes; unknown subtypes get a
// muted line rather than a raw-JSON dump in the loud `sys` colour.
if (v.type === 'system') {
const cat = v._category;
const summary = v._summary;
// thinking_tok: collapse many ticks into one in-place counter row.
if (cat === 'thinking_tok') {
updateThinkingTokens(api, summary || 'thinking…');
if (v.subtype === 'init') return;
if (v.subtype === 'api_retry') {
const parts = ['⚠ api retry'];
if (v.attempt != null && v.max_retries != null)
parts.push(v.attempt + '/' + v.max_retries);
if (v.error) parts.push(String(v.error));
else if (v.error_status) parts.push('HTTP ' + v.error_status);
if (v.retry_delay_ms != null)
parts.push(Math.round(v.retry_delay_ms) + 'ms');
api.row('note', parts.join(' · '));
return;
}
// plugin_install: coalesced in-place row while the plugin loads.
if (v.subtype === 'api_error') {
const msg = v.error || v.message
|| (v.error_status ? 'HTTP ' + v.error_status : 'unknown');
api.row('note stderr', '✗ api error · ' + msg);
return;
}
// Live thinking-token counter — claude streams many of these per
// turn (a running `estimated_tokens` total while it thinks). Collapse
// consecutive ticks into ONE in-place-updating row instead of a note
// per tick (see `makeCoalescer` above).
if (v.subtype === 'thinking_tokens') {
const n = v.estimated_tokens;
const text = 'thinking … '
+ (n != null ? '~' + Number(n).toLocaleString() + ' tokens' : '');
updateThinkingTokens(api, text);
return;
}
// plugin_install: claude is loading/finishing a plugin (MCP server or
// slash-command provider). Show the status so the operator knows when
// a fresh session is loading its toolset.
if (v.subtype === 'plugin_install') {
updatePluginInstall(api, summary || '⚙ plugin install');
const status = v.status === 'completed' ? '✓ done'
: v.status === 'started' ? 'loading…'
: (v.status || '?');
updatePluginInstall(api, '⚙ plugin install · ' + status);
return;
}
// status: backend provides the base label; when the harness state is
// `compacting` we override with elapsed time from `stateSince` — a
// client-side wall-clock value the backend can't know at emit time.
// commands_changed: the set of available slash commands changed (usually
// right after plugin_install). Show the count in the summary; expand to
// see the full list.
if (v.subtype === 'commands_changed') {
const cmds = Array.isArray(v.commands) ? v.commands : [];
if (!cmds.length) {
api.row('note', '⚙ commands changed · (empty)');
return;
}
const summary = '⚙ commands changed · ' + cmds.length + ' available';
const body = cmds.map((c) => {
const aliases = c.aliases && c.aliases.length
? ' [/' + c.aliases.join(', /') + ']' : '';
return '/' + c.name + aliases;
}).join('\n');
api.details('note', summary, body);
return;
}
// compact_boundary: claude completed a compaction pass. The metadata
// carries pre/post token counts, duration, and the trigger (manual vs
// auto). Show a single summary line so the operator can gauge how much
// context was shed.
if (v.subtype === 'compact_boundary') {
const m = v.compact_metadata || {};
const parts = ['⚙ compact'];
if (m.trigger) parts.push(m.trigger);
if (m.pre_tokens != null && m.post_tokens != null) {
const fmtTok = (n) => n >= 1_000_000 ? (n / 1_000_000).toFixed(1) + 'M'
: n >= 1_000 ? Math.round(n / 1000) + 'k'
: String(n);
parts.push(fmtTok(m.pre_tokens) + '→' + fmtTok(m.post_tokens) + ' tokens');
}
if (m.duration_ms != null) {
const ms = m.duration_ms;
parts.push(ms < 1000 ? ms + 'ms' : (ms / 1000).toFixed(1) + 's');
}
api.row('note', parts.join(' · '));
return;
}
// Bare `status` ticks (claude's own generic "still working" signal,
// no detail beyond the label) — collapse consecutive ticks into one
// updating row instead of a fresh note each (see `makeCoalescer`
// above). When the harness state is `compacting` (set by the
// `turn_state_changed` SSE event), show elapsed time via `stateSince`
// — the same source the state badge uses — so the terminal reflects
// compaction progress; otherwise fall back to the generic label.
if (v.subtype === 'status') {
let label = summary || '⚙ status';
let label = '⚙ status';
if (stateName === 'compacting') {
const elapsed = Math.round((Date.now() - stateSince) / 1000);
label = '⚙ compact · ' + elapsed + 's…';
@ -1713,13 +2024,10 @@ window.marked = marked;
updateStatus(api, label);
return;
}
// details: expandable row with _summary as header, _body as content.
if (cat === 'details') {
api.details('note', summary || '⚙ ' + (v.subtype || ''), v._body || '');
return;
}
// note (and any unknown category): single summary line.
api.row('note', summary || '⚙ ' + (v.subtype || 'system'));
// Other system subtypes (context_window_exceeded, etc.) — render a
// muted note with the subtype label; reserve the loud orange `sys`
// catch-all for truly unrecognised top-level types.
api.row('note', '⚙ ' + (v.subtype || 'system'));
return;
}
// Background-task subagent events (claude's `Task` tool spawns
@ -1744,15 +2052,8 @@ window.marked = marked;
}
else if (c.type === 'tool_use') {
if (c.id && c.name) toolNameById.set(c.id, c.name);
// `_category: "rich"` is stamped by the backend on tools that
// have full-body renderers (Write/Edit diffs, send/ask/answer
// message bodies). Flat-row tools use backend _icon/_summary.
if (c._category === 'rich') {
if (!renderRichToolUse(c, api)) {
api.row('tool-use', c._summary || c.name || '?', c._icon || '🔧');
}
} else {
api.row('tool-use', c._summary || c.name || '?', c._icon || '🔧');
if (!renderRichToolUse(c, api)) {
api.row('tool-use', fmtToolUse(c), toolIcon(c.name));
}
}
}

View file

@ -22,7 +22,6 @@ mod plugins;
mod prompt;
mod serve_common;
mod stats;
mod stream_enrich;
mod turn;
mod turn_stats;
mod vacuum;

View file

@ -1,772 +0,0 @@
//! Enrich raw claude stream-json values before SSE delivery.
//!
//! A single [`enrich`] function stamps `_icon`, `_summary`, `_category`,
//! and optionally `_body` onto [`crate::events::LiveEvent::Stream`] payloads
//! so the frontend can read pre-computed fields instead of duplicating the
//! dispatch logic in JavaScript.
//!
//! The sqlite event log stores raw (un-enriched) events — the DB never needs
//! migration when the enrichment logic changes. Enrichment is applied at
//! SSE-emit time in [`crate::web_ui::stream`] so both the live tail
//! (`events/stream`) and the history replay (`events/history`) endpoints
//! deliver the same enriched shape.
//!
//! # Migration (two-phase)
//!
//! **Phase 1** (this change): backend stamps `_icon`/`_summary`/`_category`
//! fields; the client reads them when present and falls back to its own JS
//! tables when absent. Zero user-visible change — a no-op for clients that
//! haven't yet been updated.
//!
//! **Phase 2** (follow-up): the client drops the JS tables once phase 1 is
//! deployed everywhere.
use serde_json::{Value, json};
/// Stamp enrichment fields onto a raw claude stream-json [`Value`].
///
/// - `type="system"` events get `_category` + `_summary` (and `_body` for
/// expandable detail, e.g. `commands_changed`).
/// - `type="assistant"` events get `_icon`, `_summary`, and optionally
/// `_category: "rich"` stamped onto each `message.content[]` entry that
/// has `type="tool_use"`.
///
/// No-ops for unknown/unhandled top-level types. Existing `_`-prefixed fields
/// are left unchanged so the call is idempotent (history replay may hit
/// already-enriched values if the DB is ever pre-populated by a future phase).
pub fn enrich(v: &mut Value) {
match v.get("type").and_then(Value::as_str).unwrap_or("") {
// Top-level result/rate_limit_event are drop-category noise — stamp
// the same category the frontend uses to silently discard them.
"result" | "rate_limit_event" => {
if let Some(obj) = v.as_object_mut() {
obj.entry("_category").or_insert_with(|| json!("drop"));
}
}
"system" => enrich_system(v),
"assistant" => enrich_assistant(v),
_ => {}
}
}
// ---------------------------------------------------------------------------
// system events
// ---------------------------------------------------------------------------
fn enrich_system(v: &mut Value) {
if v.get("_category").is_some() {
return; // idempotent
}
let subtype = v
.get("subtype")
.and_then(Value::as_str)
.unwrap_or("")
.to_owned();
let (category, summary, body) = system_fields(v, &subtype);
let Some(obj) = v.as_object_mut() else { return };
obj.insert("_category".to_owned(), json!(category));
if let Some(s) = summary {
obj.insert("_summary".to_owned(), json!(s));
}
if let Some(b) = body {
obj.insert("_body".to_owned(), json!(b));
}
}
/// Compute `(category, summary, body)` for a `type="system"` event.
///
/// Categories understood by the frontend:
/// - `"drop"` — client discards without rendering (noise)
/// - `"thinking_tok"` — client updates a single in-place counter row
/// - `"note"` — client renders `_summary` as a terminal note line
/// - `"details"` — client renders a collapsible `<details>` with `_summary`
/// as the header and `_body` as the expanded text
fn system_fields(v: &Value, subtype: &str) -> (&'static str, Option<String>, Option<String>) {
match subtype {
// Silent startup / stream-end noise — drop.
"init" | "result" | "rate_limit_event" => ("drop", None, None),
// Live thinking-token counter.
"thinking_tokens" => {
let n = v.get("estimated_tokens").and_then(Value::as_u64);
let s = n.map_or_else(
|| "thinking…".to_owned(),
|n| format!("thinking… ~{n} tokens"),
);
("thinking_tok", Some(s), None)
}
// API retry back-off.
"api_retry" => {
let mut parts = vec!["⚠ api retry".to_owned()];
if let (Some(a), Some(m)) = (
v.get("attempt").and_then(Value::as_u64),
v.get("max_retries").and_then(Value::as_u64),
) {
parts.push(format!("{a}/{m}"));
}
if let Some(e) = v.get("error").and_then(Value::as_str) {
parts.push(e.to_owned());
} else if let Some(s) = v.get("error_status").and_then(Value::as_u64) {
parts.push(format!("HTTP {s}"));
}
if let Some(ms) = v.get("retry_delay_ms").and_then(Value::as_f64) {
parts.push(format!("{:.0}ms", ms.round()));
}
("note", Some(parts.join(" · ")), None)
}
// Hard API error.
"api_error" => {
let msg = v
.get("error")
.or_else(|| v.get("message"))
.and_then(Value::as_str)
.map_or_else(
|| {
v.get("error_status")
.and_then(Value::as_u64)
.map_or("unknown".to_owned(), |s| format!("HTTP {s}"))
},
str::to_owned,
);
("note", Some(format!("✗ api error · {msg}")), None)
}
// Plugin (MCP server / slash-command provider) load progress.
"plugin_install" => {
let status = v.get("status").and_then(Value::as_str).unwrap_or("?");
let label = match status {
"completed" => "✓ done",
"started" => "loading…",
other => other,
};
("note", Some(format!("⚙ plugin install · {label}")), None)
}
// Available slash-command set changed — expandable list.
"commands_changed" => {
let cmds = v.get("commands").and_then(Value::as_array);
let count = cmds.map_or(0, Vec::len);
let summary = format!("⚙ commands changed · {count} available");
let body = cmds.map(|cmds| {
cmds.iter()
.filter_map(|c| c.get("name").and_then(Value::as_str))
.map(|n| format!("/{n}"))
.collect::<Vec<_>>()
.join("\n")
});
("details", Some(summary), body)
}
// Compaction complete — pre/post token counts + duration.
"compact_boundary" => {
let m = v
.get("compact_metadata")
.cloned()
.unwrap_or_else(|| json!({}));
let mut parts = vec!["⚙ compact".to_owned()];
if let Some(t) = m.get("trigger").and_then(Value::as_str) {
parts.push(t.to_owned());
}
if let (Some(pre), Some(post)) = (
m.get("pre_tokens").and_then(Value::as_u64),
m.get("post_tokens").and_then(Value::as_u64),
) {
parts.push(format!("{}{} tokens", fmt_tok(pre), fmt_tok(post)));
}
if let Some(ms) = m.get("duration_ms").and_then(Value::as_u64) {
let dur = if ms < 1000 {
format!("{ms}ms")
} else {
#[allow(clippy::cast_precision_loss)]
let s_f = ms as f64 / 1_000.0;
format!("{s_f:.1}s")
};
parts.push(dur);
}
("note", Some(parts.join(" · ")), None)
}
// Generic "still working" heartbeat tick.
"status" => ("note", Some("⚙ status".to_owned()), None),
// Unknown subtype — render with the subtype label as a muted note.
other => ("note", Some(format!("{other}")), None),
}
}
// ---------------------------------------------------------------------------
// assistant events
// ---------------------------------------------------------------------------
fn enrich_assistant(v: &mut Value) {
// Navigate message.content[] — absent on text-only turns.
let Some(content) = v
.get_mut("message")
.and_then(|m| m.get_mut("content"))
.and_then(Value::as_array_mut)
else {
return;
};
for entry in content.iter_mut() {
enrich_tool_use_entry(entry);
}
}
fn enrich_tool_use_entry(entry: &mut Value) {
if entry.get("type").and_then(Value::as_str) != Some("tool_use") {
return;
}
if entry.get("_icon").is_some() {
return; // idempotent
}
let name = entry
.get("name")
.and_then(Value::as_str)
.unwrap_or("")
.to_owned();
let input = entry.get("input").cloned().unwrap_or_else(|| json!({}));
let icon = tool_icon(&name);
let summary = fmt_tool_use(&name, &input);
let rich = is_rich_tool(&name);
let Some(obj) = entry.as_object_mut() else {
return;
};
obj.insert("_icon".to_owned(), json!(icon));
obj.insert("_summary".to_owned(), json!(summary));
if rich {
obj.insert("_category".to_owned(), json!("rich"));
}
}
// ---------------------------------------------------------------------------
// tool helpers
// ---------------------------------------------------------------------------
/// Whether this tool is rendered by the frontend's *rich* renderer (diff
/// view, body expansion) rather than the flat `_summary` row. Flagged with
/// `_category: "rich"` so the client can distinguish without re-implementing
/// the tool name list.
fn is_rich_tool(name: &str) -> bool {
matches!(
name,
"Write"
| "Edit"
| "mcp__bash__run"
| "mcp__hyperhive__send"
| "mcp__hyperhive__ask"
| "mcp__hyperhive__answer"
)
}
fn tool_icon(name: &str) -> &'static str {
// Exact-match table first, then prefix/contains fallbacks.
match name {
"mcp__hyperhive__send" => "📤",
"mcp__hyperhive__recv" => "📥",
"mcp__hyperhive__ask" => "",
"mcp__hyperhive__answer" => "✍️",
"mcp__hyperhive__remind" => "",
"mcp__hyperhive__set_status" => "🏷️",
"mcp__hyperhive__get_loose_ends" => "🪢",
"mcp__hyperhive__cancel_loose_end" => "✂️",
"mcp__hyperhive__ack_until" => "",
"mcp__hyperhive__get_agent_meta" => "",
"mcp__hyperhive__request_next_turn" => "",
"mcp__hyperhive__restart" => "",
"mcp__hyperhive__kill" => "⏹️",
"mcp__hyperhive__start" => "▶️",
"mcp__hyperhive__update" => "🔄",
"mcp__hyperhive__list_containers"
| "mcp__matrix__list_rooms"
| "mcp__matrix__list_room_members"
| "mcp__matrix__list_invites" => "📋",
"mcp__hyperhive__get_logs" | "mcp__hyperhive__get_host_journal" => "📜",
"mcp__matrix__read_room" | "Read" => "📖",
"mcp__matrix__mark_read" => "👁️",
"mcp__bash__kill" => "🛑",
"Write" => "💾",
"Edit" => "✏️",
"Glob" | "Grep" => "🔍",
_ => tool_icon_fallback(name),
}
}
fn tool_icon_fallback(name: &str) -> &'static str {
if name.starts_with("mcp__matrix__") {
"💬"
} else if name.starts_with("mcp__bash__") {
"🖥️"
} else if name.contains("schedule") {
"⏱️"
} else if name.starts_with("mcp__hyperhive__request_") {
"📦"
} else {
"🔧"
}
}
// ---------------------------------------------------------------------------
// tool_use summary formatter — dispatch to per-family sub-functions
// ---------------------------------------------------------------------------
fn fmt_tool_use(name: &str, input: &Value) -> String {
// Short name: strip the MCP server prefix for display.
let short = if let Some(rest) = name.strip_prefix("mcp__hyperhive__") {
format!("{rest}*")
} else if let Some(rest) = name.strip_prefix("mcp__bash__") {
format!("{rest}*")
} else if let Some(rest) = name.strip_prefix("mcp__matrix__") {
format!("{rest}*")
} else {
name.to_owned()
};
if name.starts_with("mcp__hyperhive__") {
fmt_hyperhive_tool(name, &short, input)
} else if name.starts_with("mcp__bash__") {
fmt_bash_tool(name, &short, input)
} else if name.starts_with("mcp__matrix__") {
fmt_matrix_tool(name, &short, input)
} else {
fmt_builtin_tool(name, &short, input)
}
}
/// Built-in claude tools: Read/Write/Edit/Glob/Grep/Bash/TodoWrite and any
/// unknown tool that doesn't carry a known MCP server prefix.
fn fmt_builtin_tool(name: &str, short: &str, input: &Value) -> String {
match name {
"Read" | "Write" | "Edit" => format!("{short} {}", sv(input, "file_path")),
"Glob" | "Grep" => format!("{short} {}", sv(input, "pattern")),
"Bash" => {
let bg = if input
.get("run_in_background")
.and_then(Value::as_bool)
.unwrap_or(false)
{
" [bg]"
} else {
""
};
format!("{short}{bg} $ {}", sv(input, "command"))
}
"TodoWrite" => {
let n = input
.get("todos")
.and_then(Value::as_array)
.map_or(0, Vec::len);
format!("{short} ({n} items)")
}
_ => fmt_args_generic(short, input),
}
}
/// `mcp__hyperhive__*` tools.
fn fmt_hyperhive_tool(name: &str, short: &str, input: &Value) -> String {
match name {
"mcp__hyperhive__send" => {
let body = trim_str(&sv(input, "body"), 80);
format!("{short}{}: {}", sv(input, "to"), json_str(&body))
}
"mcp__hyperhive__recv" => {
let mut parts = Vec::new();
if let Some(w) = input.get("wait_seconds").and_then(Value::as_u64) {
parts.push(format!("wait {w}s"));
}
if let Some(m) = input.get("max").and_then(Value::as_u64) {
parts.push(format!("max {m}"));
}
if parts.is_empty() {
format!("{short}()")
} else {
format!("{short} {}", parts.join(" · "))
}
}
"mcp__hyperhive__kill"
| "mcp__hyperhive__restart"
| "mcp__hyperhive__start"
| "mcp__hyperhive__update"
| "mcp__hyperhive__request_init_config" => format!("{short} {}", sv(input, "name")),
"mcp__hyperhive__ack_until" => {
let up_to = input
.get("up_to")
.and_then(Value::as_u64)
.map_or_else(|| "?".to_owned(), |n| n.to_string());
format!("{short}{up_to}")
}
"mcp__hyperhive__get_logs" => {
let lines = input
.get("lines")
.and_then(Value::as_u64)
.map(|n| format!(" · {n}L"))
.unwrap_or_default();
format!("{short} {}{lines}", sv(input, "agent"))
}
"mcp__hyperhive__get_host_journal" => {
let mut parts = Vec::new();
if let Some(c) = input.get("container").and_then(Value::as_str) {
parts.push(c.to_owned());
} else if let Some(u) = input.get("unit").and_then(Value::as_str) {
parts.push(u.to_owned());
}
if let Some(g) = input.get("grep").and_then(Value::as_str) {
parts.push(format!("/{g}/"));
}
if let Some(l) = input.get("lines").and_then(Value::as_u64) {
parts.push(format!("{l}L"));
}
if parts.is_empty() {
format!("{short}()")
} else {
format!("{short} {}", parts.join(" · "))
}
}
"mcp__hyperhive__remind" => fmt_hyperhive_remind(short, input),
"mcp__hyperhive__request_update_meta_inputs"
| "mcp__hyperhive__list_schedules"
| "mcp__hyperhive__cancel_schedule"
| "mcp__hyperhive__fire_schedule_now"
| "mcp__hyperhive__edit_schedule"
| "mcp__hyperhive__request_schedule_prompt" => {
fmt_hyperhive_schedule_tool(name, short, input)
}
"mcp__hyperhive__set_status" => {
format!("{short} \"{}\"", trim_str(&sv(input, "text"), 60))
}
"mcp__hyperhive__get_loose_ends" => {
let agent = input
.get("agent")
.and_then(Value::as_str)
.map_or_else(|| "()".to_owned(), |a| format!(" [{a}]"));
format!("{short}{agent}")
}
"mcp__hyperhive__get_agent_meta" => {
let name_part = input
.get("name")
.and_then(Value::as_str)
.map_or_else(|| "()".to_owned(), |n| format!(" {n}"));
format!("{short}{name_part}")
}
"mcp__hyperhive__cancel_loose_end" => {
let id = input
.get("id")
.and_then(Value::as_u64)
.map_or_else(|| "?".to_owned(), |n| n.to_string());
format!("{short} {} #{id}", sv(input, "kind"))
}
_ => fmt_args_generic(short, input),
}
}
fn fmt_hyperhive_remind(short: &str, input: &Value) -> String {
let when = if let Some(s) = input.get("delay_seconds").and_then(Value::as_u64) {
if s < 60 {
format!("+{s}s")
} else if s < 3_600 {
format!("+{}m", s / 60)
} else {
#[allow(clippy::cast_precision_loss)]
let h_f = s as f64 / 3_600.0;
format!("+{h_f:.1}h")
}
} else if let Some(ts) = input.get("at_unix_timestamp").and_then(Value::as_u64) {
let h = (ts % 86_400) / 3_600;
let m = (ts % 3_600) / 60;
format!("at {h:02}:{m:02}Z")
} else {
String::new()
};
let msg_raw = input
.get("message")
.or_else(|| input.get("file_path"))
.and_then(Value::as_str)
.unwrap_or("");
let msg = trim_str(&msg_raw.replace(char::is_whitespace, " "), 60);
let when_part = if when.is_empty() {
String::new()
} else {
format!(" {when}")
};
let msg_part = if msg.is_empty() {
String::new()
} else {
format!(" \"{msg}\"")
};
format!("{short}{when_part}{msg_part}")
}
fn fmt_hyperhive_edit_schedule(short: &str, input: &Value) -> String {
let id = input
.get("id")
.and_then(Value::as_u64)
.map_or_else(|| "?".to_owned(), |n| n.to_string());
let mut parts = vec![format!("#{id}")];
if input.get("body").is_some() {
parts.push("body".to_owned());
}
if input.get("interval_seconds").is_some() {
parts.push("interval".to_owned());
}
if input.get("next_fire_at_unix").is_some() {
parts.push("next".to_owned());
}
if let Some(a) = input.get("targets_add").and_then(Value::as_array)
&& !a.is_empty()
{
parts.push(format!("+{} tgt", a.len()));
}
if let Some(r) = input.get("targets_remove").and_then(Value::as_array)
&& !r.is_empty()
{
parts.push(format!("-{} tgt", r.len()));
}
format!("{short} {}", parts.join(" · "))
}
/// Schedule-management hyperhive tools (list/cancel/fire/edit/request).
fn fmt_hyperhive_schedule_tool(name: &str, short: &str, input: &Value) -> String {
match name {
"mcp__hyperhive__request_update_meta_inputs" => {
let ins = match input.get("inputs").and_then(Value::as_array) {
Some(arr) if !arr.is_empty() => {
let names: Vec<&str> = arr.iter().filter_map(Value::as_str).take(4).collect();
let tail = if arr.len() > 4 { ", …" } else { "" };
format!("[{}{}]", names.join(", "), tail)
}
_ => "all".to_owned(),
};
format!("{short} {ins}")
}
"mcp__hyperhive__list_schedules" => format!("{short}()"),
"mcp__hyperhive__cancel_schedule" => {
let id = input
.get("id")
.and_then(Value::as_u64)
.map_or_else(|| "?".to_owned(), |n| n.to_string());
let tgts = match input.get("targets").and_then(Value::as_array) {
Some(arr) if !arr.is_empty() => {
let names: Vec<&str> = arr.iter().filter_map(Value::as_str).collect();
format!(" [{}]", names.join(", "))
}
_ => " all".to_owned(),
};
format!("{short} #{id}{tgts}")
}
"mcp__hyperhive__fire_schedule_now" => {
let id = input
.get("id")
.and_then(Value::as_u64)
.map_or_else(|| "?".to_owned(), |n| n.to_string());
format!("{short} #{id}")
}
"mcp__hyperhive__edit_schedule" => fmt_hyperhive_edit_schedule(short, input),
"mcp__hyperhive__request_schedule_prompt" => {
let tgts = match input.get("targets").and_then(Value::as_array) {
Some(arr) => arr
.iter()
.filter_map(Value::as_str)
.collect::<Vec<_>>()
.join(", "),
None => "?".to_owned(),
};
let when = input
.get("first_fire_at_unix")
.and_then(Value::as_u64)
.map_or_else(
|| "?".to_owned(),
|ts| {
let h = (ts % 86_400) / 3_600;
let m = (ts % 3_600) / 60;
format!("{h:02}:{m:02}Z")
},
);
let recur = input
.get("interval_seconds")
.and_then(Value::as_u64)
.map(|s| format!(" +{s}s"))
.unwrap_or_default();
format!("{short}{tgts} at {when}{recur}")
}
_ => fmt_args_generic(short, input),
}
}
/// `mcp__bash__*` tools.
fn fmt_bash_tool(name: &str, short: &str, input: &Value) -> String {
match name {
"mcp__bash__run" => {
let cmd = sv(input, "cmd");
let first = cmd.lines().next().unwrap_or("").trim().to_owned();
format!("{short} $ {}", trim_str(&first, 72))
}
"mcp__bash__status" => {
let wait = input
.get("wait_seconds")
.and_then(Value::as_u64)
.map(|w| format!(" · wait {w}s"))
.unwrap_or_default();
format!("{short} id:{}{wait}", sv(input, "id"))
}
"mcp__bash__kill" => {
let force = if input.get("force").and_then(Value::as_bool).unwrap_or(false) {
" [force]"
} else {
""
};
format!("{short} {}{force}", sv(input, "id"))
}
_ => fmt_args_generic(short, input),
}
}
/// `mcp__matrix__*` tools.
fn fmt_matrix_tool(name: &str, short: &str, input: &Value) -> String {
match name {
"mcp__matrix__read_room" => {
let limit = input
.get("limit")
.and_then(Value::as_u64)
.map(|l| format!(" [{l}]"))
.unwrap_or_default();
format!("{short} {}{limit}", fmt_room(&sv(input, "room")))
}
"mcp__matrix__mark_read" | "mcp__matrix__join_room" | "mcp__matrix__download_file" => {
format!("{short} {}", fmt_room(&sv(input, "room")))
}
"mcp__matrix__send_message" | "mcp__matrix__send_reply" => format!(
"{short} → {}: {}",
fmt_room(&sv(input, "room")),
json_str(&trim_str(&sv(input, "body"), 50))
),
"mcp__matrix__send_dm" => format!(
"{short} → {}: {}",
fmt_user(&sv(input, "user_id")),
json_str(&trim_str(&sv(input, "body"), 50))
),
"mcp__matrix__send_reaction" => format!(
"{short} {} {}",
fmt_room(&sv(input, "room")),
sv(input, "key")
),
"mcp__matrix__open_dm" => format!("{short} {}", fmt_user(&sv(input, "user_id"))),
"mcp__matrix__invite_user" => format!(
"{short} {} → {}",
fmt_user(&sv(input, "user_id")),
fmt_room(&sv(input, "room"))
),
_ => fmt_args_generic(short, input),
}
}
// ---------------------------------------------------------------------------
// generic fallback formatter (port of JS fmtArgsGeneric)
// ---------------------------------------------------------------------------
fn fmt_args_generic(name: &str, input: &Value) -> String {
let Some(obj) = input.as_object() else {
return format!("{name}()");
};
if obj.is_empty() {
return format!("{name}()");
}
let keys: Vec<&String> = obj.keys().collect();
if keys.len() == 1 {
let k = keys[0];
let v = &obj[k];
if let Some(s) = v.as_str() {
return format!("{name} {k}: {}", json_str(&trim_str(s, 100)));
}
if v.is_number() || v.is_boolean() {
return format!("{name} {k}: {v}");
}
}
let pretty: Vec<String> = keys
.iter()
.take(4)
.map(|k| {
let v = &obj[*k];
if v.is_null() {
format!("{k}: null")
} else if let Some(s) = v.as_str() {
format!("{k}: {}", json_str(&trim_str(s, 40)))
} else if v.is_number() || v.is_boolean() {
format!("{k}: {v}")
} else if let Some(arr) = v.as_array() {
format!("{k}: [{}]", arr.len())
} else {
format!("{k}: {{}}")
}
})
.collect();
let tail = if keys.len() > 4 {
format!(" …+{}", keys.len() - 4)
} else {
String::new()
};
format!("{name} {}{tail}", pretty.join(" · "))
}
// ---------------------------------------------------------------------------
// string utilities
// ---------------------------------------------------------------------------
/// Format a token count as "1k", "1.2M", etc. (mirrors JS `fmtTok`).
fn fmt_tok(n: u64) -> String {
if n >= 1_000_000 {
// Divide before casting to keep the value in a range where f64
// has enough precision for one decimal place (≤ ~9×10^12 tokens).
#[allow(clippy::cast_precision_loss)]
let m = (n / 1_000) as f64 / 1_000.0;
format!("{m:.1}M")
} else if n >= 1_000 {
format!("{}k", n / 1_000)
} else {
n.to_string()
}
}
/// Shorten a matrix room id/alias for display.
fn fmt_room(r: &str) -> String {
if r.starts_with('!') {
// Room id: keep only the local part before the colon (up to 9 chars).
// Use chars().take() so we never slice on a non-ASCII byte boundary.
let colon = r.find(':').unwrap_or(r.len());
r.chars().take(colon.min(9)).collect()
} else if r.starts_with('#') {
// Alias: keep `#name` part before the server.
r.split(':').next().unwrap_or(r).to_owned()
} else {
trim_str(r, 20)
}
}
/// Shorten `@user:server` → `@user`.
fn fmt_user(u: &str) -> String {
let end = u.find(':').unwrap_or(u.len());
u[..end].to_owned()
}
/// Extract a string field from a JSON object, defaulting to `""`.
fn sv(v: &Value, key: &str) -> String {
v.get(key).and_then(Value::as_str).unwrap_or("").to_owned()
}
/// Trim to `max` *characters* (not bytes), collapsing whitespace first.
fn trim_str(s: &str, max: usize) -> String {
let collapsed: String = s.split_whitespace().collect::<Vec<_>>().join(" ");
let mut chars = collapsed.chars();
let trimmed: String = chars.by_ref().take(max).collect();
if chars.next().is_some() {
format!("{trimmed}")
} else {
trimmed
}
}
/// Serialize a string as a JSON string literal (e.g. `"hello"`) for display
/// in summary lines. Falls back to the raw string if serialization fails.
fn json_str(s: &str) -> String {
serde_json::to_string(s).unwrap_or_else(|_| s.to_owned())
}

View file

@ -40,17 +40,6 @@ pub(super) async fn events_history(
};
let (events, min_id, has_more) = state.bus.history_page(before, limit);
// Apply the same enrichment as the live SSE path so history replay
// and live tail deliver identical shapes. The DB stores raw events.
let events: Vec<_> = events
.into_iter()
.map(|mut se| {
if let crate::events::LiveEvent::Stream(ref mut v) = se.event {
crate::stream_enrich::enrich(v);
}
se
})
.collect();
let mut resp = serde_json::json!({
"events": events,
"min_id": min_id,
@ -79,15 +68,7 @@ pub(super) async fn events_stream(
.unwrap_or_default(),
);
let live = BroadcastStream::new(rx).filter_map(|res| {
let mut ev = res.ok()?;
// Enrich stream-json values with pre-computed display fields
// (`_icon`, `_summary`, `_category`) so the frontend doesn't need to
// duplicate the dispatch logic. The DB stores raw events; enrichment
// is applied here so both the live tail and the history endpoint
// deliver the same shape (see `events_history` above).
if let crate::events::LiveEvent::Stream(ref mut v) = ev.event {
crate::stream_enrich::enrich(v);
}
let ev = res.ok()?;
let json = serde_json::to_string(&ev).ok()?;
Some(Ok(Event::default().data(json)))
});