diff --git a/internal/httpapi/activity.go b/internal/httpapi/activity.go new file mode 100644 index 0000000..4351c3d --- /dev/null +++ b/internal/httpapi/activity.go @@ -0,0 +1,215 @@ +package httpapi + +import ( + "encoding/json" + "log/slog" + "net/http" + "strconv" + "strings" + + "github.com/go-chi/chi/v5" +) + +// activityItem is one row in the global activity feed — a human-readable +// projection of an execution, independent of the paginated/alphabetically- +// sorted ListExecutions (which orders by target slug for entity-scoped +// browsing, not recency — wrong shape for "what just happened"). +type activityItem struct { + ID string `json:"id"` + Target string `json:"target"` + Verb string `json:"verb"` // e.g. "run", "pct_create", "systemctl" + Summary string `json:"summary"` // human-readable: the command, or purpose, or action detail + RiskClass string `json:"risk_class"` + Status string `json:"status"` + DurationMs *int `json:"duration_ms"` + Error string `json:"error,omitempty"` + CreatedAt string `json:"created_at"` + CompletedAt *string `json:"completed_at"` +} + +// splitAction parses the "verb:params" encoding used throughout executions.action +// (see internal/mcp/server.go) into a verb and a human-readable summary. For +// `run`, params is JSON {command, purpose} — show the purpose if present +// (it's written for a human), falling back to the raw command. For other +// actions (pct_create, systemctl, apt_upgrade, pct_exec), params is either a +// JSON blob or a short flag string — truncate either as a fallback summary. +func splitAction(action string) (verb, summary string) { + idx := strings.IndexByte(action, ':') + if idx < 0 { + return action, "" + } + verb, params := action[:idx], action[idx+1:] + if verb == "run" { + var p struct { + Command string `json:"command"` + Purpose string `json:"purpose"` + } + if json.Unmarshal([]byte(params), &p) == nil { + if p.Purpose != "" { + return verb, p.Purpose + } + return verb, p.Command + } + } + if verb == "pct_create" { + var p struct { + Hostname string `json:"hostname"` + } + if json.Unmarshal([]byte(params), &p) == nil && p.Hostname != "" { + return verb, "provision " + p.Hostname + } + } + if len(params) > 140 { + params = params[:140] + "…" + } + return verb, params +} + +// serveRecentActivity backs the Operations page's live activity feed — the +// global "what is the system doing / what did it just do" view, recency- +// ordered (unlike ListExecutions, which sorts by target for pagination). +// Custom route, same shape/rationale as serveRecentKnowledge. +func (s *Server) serveRecentActivity(w http.ResponseWriter, req *http.Request) { + ctx := req.Context() + limit := 50 + if l := req.URL.Query().Get("limit"); l != "" { + if n, err := strconv.Atoi(l); err == nil && n > 0 && n <= 200 { + limit = n + } + } + + rows, err := s.pool.Query(ctx, ` + SELECT e.entity_id, te.slug, e.action, e.risk_class, e.status, + e.duration_ms, e.result, e.created_at::text, e.completed_at::text + FROM executions e + JOIN entities te ON te.id = e.target_entity_id + ORDER BY e.created_at DESC + LIMIT $1`, limit) + if err != nil { + writeProblem(w, req, http.StatusInternalServerError, "query failed", err.Error()) + return + } + defer rows.Close() + + items := []activityItem{} + for rows.Next() { + var it activityItem + var action string + var resultBytes []byte + var completedAt *string + if err := rows.Scan(&it.ID, &it.Target, &action, &it.RiskClass, &it.Status, + &it.DurationMs, &resultBytes, &it.CreatedAt, &completedAt); err != nil { + slog.Error("httpapi: activity/recent row scan failed", "error", err) + continue + } + it.Verb, it.Summary = splitAction(action) + it.CompletedAt = completedAt + if len(resultBytes) > 0 { + var result map[string]any + if json.Unmarshal(resultBytes, &result) == nil { + if e, ok := result["error"].(string); ok && e != "" { + if len(e) > 200 { + e = e[:200] + "…" + } + it.Error = e + } + } + } + items = append(items, it) + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]any{"items": items}) +} + +// sessionDigestItem summarizes one execution for the session digest. +type sessionDigestItem struct { + Target string `json:"target"` + Verb string `json:"verb"` + Summary string `json:"summary"` + RiskClass string `json:"risk_class"` + Status string `json:"status"` +} + +// serveSessionDigest answers "what did THIS chat session actually do" — +// commands run (grouped by outcome), distinct entities touched, and knowledge +// written during the session's time window. Uses nomos_plan_executions (the +// session<->execution link added for auto-continuation) as the source of +// truth for which executions belong to this session; knowledge correlation is +// a best-effort time-window match since knowledge_entities has no session_id. +func (s *Server) serveSessionDigest(w http.ResponseWriter, req *http.Request) { + ctx := req.Context() + sessionID := chi.URLParam(req, "id") + if sessionID == "" { + writeProblem(w, req, http.StatusBadRequest, "missing session id", "") + return + } + + rows, err := s.pool.Query(ctx, ` + SELECT te.slug, e.action, e.risk_class, e.status + FROM nomos_plan_executions l + JOIN executions e ON e.entity_id = l.execution_id + JOIN entities te ON te.id = e.target_entity_id + WHERE l.session_id = $1 + ORDER BY e.created_at`, sessionID) + if err != nil { + writeProblem(w, req, http.StatusInternalServerError, "query failed", err.Error()) + return + } + defer rows.Close() + + items := []sessionDigestItem{} + byStatus := map[string]int{} + targets := map[string]bool{} + for rows.Next() { + var it sessionDigestItem + var action string + if err := rows.Scan(&it.Target, &action, &it.RiskClass, &it.Status); err != nil { + continue + } + it.Verb, it.Summary = splitAction(action) + items = append(items, it) + byStatus[it.Status]++ + targets[it.Target] = true + } + + entityList := make([]string, 0, len(targets)) + for t := range targets { + entityList = append(entityList, t) + } + + // Best-effort knowledge correlation: notes the agent wrote during this + // session's active window. Not exact (no session_id on knowledge_entities) + // but close enough to show "you learned N things in this session". + var knowledgeTitles []string + krows, err := s.pool.Query(ctx, ` + SELECT ke.title FROM knowledge_entities ke + WHERE ke.source = 'nomos-agent' + AND ke.updated_at BETWEEN + (SELECT COALESCE(MIN(created_at), now()) FROM agent_messages WHERE session_id = $1) + AND + (SELECT COALESCE(MAX(created_at), now()) + interval '2 minutes' FROM agent_messages WHERE session_id = $1) + ORDER BY ke.updated_at`, sessionID) + if err == nil { + defer krows.Close() + for krows.Next() { + var t string + if krows.Scan(&t) == nil { + knowledgeTitles = append(knowledgeTitles, t) + } + } + } + if knowledgeTitles == nil { + knowledgeTitles = []string{} + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]any{ + "session_id": sessionID, + "total_executions": len(items), + "by_status": byStatus, + "entities_touched": entityList, + "executions": items, + "knowledge_created": knowledgeTitles, + }) +} diff --git a/internal/httpapi/server.go b/internal/httpapi/server.go index bbaef88..780ddb3 100644 --- a/internal/httpapi/server.go +++ b/internal/httpapi/server.go @@ -144,6 +144,12 @@ func NewHandler(ctx context.Context, pool *db.Pool, cfg config.Config, uiHandler // HandlerWithOptions so it wins over any generated catch-all. r.With(combinedAuth(cfg)).Get("/api/v1/knowledge/recent", s.serveRecentKnowledge) + // Custom (non-OpenAPI) routes: the global activity feed (recency-ordered, + // unlike ListExecutions which sorts by target for pagination) and the + // per-session "what did this session do" digest. + r.With(combinedAuth(cfg)).Get("/api/v1/activity/recent", s.serveRecentActivity) + r.With(combinedAuth(cfg)).Get("/api/v1/activity/session/{id}", s.serveSessionDigest) + // Mount MCP at /mcp (plan R3-10) nomosAgentID := uuid.Nil if cfg.NomosAgentID != "" { diff --git a/web/src/lib/api.ts b/web/src/lib/api.ts index cda5e7d..5076a19 100644 --- a/web/src/lib/api.ts +++ b/web/src/lib/api.ts @@ -242,6 +242,41 @@ export async function cancelExecution(id: string): Promise { return res.json() } +export interface ActivityItem { + id: string + target: string + verb: string + summary: string + risk_class: string + status: string + duration_ms: number | null + error?: string + created_at: string + completed_at: string | null +} + +export async function fetchRecentActivity(limit = 50): Promise { + const res = await fetch(`${API}/activity/recent?limit=${limit}`) + if (!res.ok) return [] + const data = await res.json() + return data.items ?? [] +} + +export interface SessionDigest { + session_id: string + total_executions: number + by_status: Record + entities_touched: string[] + executions: { target: string; verb: string; summary: string; risk_class: string; status: string }[] + knowledge_created: string[] +} + +export async function fetchSessionDigest(sessionId: string): Promise { + const res = await fetch(`${API}/activity/session/${sessionId}`) + if (!res.ok) return null + return res.json() +} + export interface Signal { id: string slug: string diff --git a/web/src/lib/components/SessionDigest.svelte b/web/src/lib/components/SessionDigest.svelte new file mode 100644 index 0000000..e852661 --- /dev/null +++ b/web/src/lib/components/SessionDigest.svelte @@ -0,0 +1,100 @@ + + +{#if digest && digest.total_executions > 0} +
+ + + {#if open} +
+
+ {#each Object.entries(digest.by_status) as [status, count]} + {status} × {count} + {/each} +
+ + {#if digest.entities_touched.length} +
+
Entities touched
+
+ {#each digest.entities_touched as target} + {target} + {/each} +
+
+ {/if} + +
+ {#each digest.executions as ex} +
+
+
{ex.target}
+
{ex.summary || ex.verb}
+
+ {ex.status} +
+ {/each} +
+ + {#if digest.knowledge_created.length} +
+
+ Learned this session +
+
    + {#each digest.knowledge_created as title} +
  • {title}
  • + {/each} +
+
+ {/if} +
+ {/if} +
+{/if} diff --git a/web/src/pages/Chat.svelte b/web/src/pages/Chat.svelte index e829d05..72f1864 100644 --- a/web/src/pages/Chat.svelte +++ b/web/src/pages/Chat.svelte @@ -2,6 +2,7 @@ import { messages, streaming, sendMessage, cancelStream, error } from '$lib/stores/chat' import SessionRail from '$lib/components/SessionRail.svelte' import SessionGraph from '$lib/components/SessionGraph.svelte' + import SessionDigest from '$lib/components/SessionDigest.svelte' import ToolCallGroup from '$lib/components/ToolCallGroup.svelte' import InlineApproval from '$lib/components/InlineApproval.svelte' import { Button } from '$lib/components/ui/button' @@ -188,8 +189,11 @@ : 'bg-border group-hover/rz:bg-primary/50'}" > -
- +
+ +
+ +
{/if} diff --git a/web/src/pages/Ops.svelte b/web/src/pages/Ops.svelte index 0360a25..6571a75 100644 --- a/web/src/pages/Ops.svelte +++ b/web/src/pages/Ops.svelte @@ -3,10 +3,10 @@ import { fetchApprovals, decideApproval, - fetchExecutions, + fetchRecentActivity, cancelExecution, type Approval, - type Execution + type ActivityItem } from '$lib/api' import { liveEvents, subscribeEvents } from '$lib/stores/events' import * as Tabs from '$lib/components/ui/tabs' @@ -16,30 +16,55 @@ import { toast } from 'svelte-sonner' let approvals = $state([]) - let executions = $state([]) + let activity = $state([]) let deciding = $state(null) async function loadApprovals() { approvals = await fetchApprovals() } - async function loadExecutions() { - executions = await fetchExecutions() + async function loadActivity() { + activity = await fetchRecentActivity() } onMount(() => { loadApprovals() - loadExecutions() + loadActivity() const unsubscribe = subscribeEvents() - return unsubscribe + // The activity feed has no dedicated SSE event type yet — a light poll + // keeps it live without waiting for that wiring. Cheap: one query, only + // while this page is open. + const interval = setInterval(loadActivity, 5000) + return () => { + unsubscribe() + clearInterval(interval) + } }) $effect(() => { const ev = $liveEvents[0] if (!ev) return if (ev.type.startsWith('approval.')) loadApprovals() - if (ev.type.startsWith('execution.')) loadExecutions() + if (ev.type.startsWith('execution.')) loadActivity() }) + function fmtDuration(ms: number | null): string { + if (ms == null) return '—' + if (ms < 1000) return `${ms}ms` + const s = Math.round(ms / 1000) + if (s < 60) return `${s}s` + return `${Math.floor(s / 60)}m ${s % 60}s` + } + + function fmtWhen(iso: string): string { + const d = new Date(iso).getTime() + if (!d) return '' + const s = Math.round((Date.now() - d) / 1000) + if (s < 60) return 'just now' + if (s < 3600) return `${Math.floor(s / 60)}m ago` + if (s < 86400) return `${Math.floor(s / 3600)}h ago` + return `${Math.floor(s / 86400)}d ago` + } + async function decide(id: string, decision: 'approve' | 'deny') { deciding = id const result = await decideApproval(id, decision) @@ -56,22 +81,26 @@ const result = await cancelExecution(id) if (result) { toast.success('Execution cancelled') - loadExecutions() + loadActivity() } else { toast.error('Cancel failed') } } function riskVariant(risk: string): 'default' | 'secondary' | 'destructive' { - if (risk === 'high' || risk === 'critical') return 'destructive' - if (risk === 'medium') return 'secondary' + if (risk === 'destructive') return 'destructive' + if (risk === 'config_mutation') return 'secondary' return 'default' } + // Real status vocabulary (internal/httpapi/phase3.go, cmd/nomos): the + // previous version checked statuses ('proposed', 'auto_approved', + // 'verified', 'executing'...) that don't exist anywhere in the actual + // schema — this table was never actually color-coding correctly. function execStatusVariant(status: string): 'default' | 'secondary' | 'destructive' | 'outline' { - if (['failed', 'timed_out', 'rollback_failed', 'denied'].includes(status)) return 'destructive' - if (['verified', 'auto_approved'].includes(status)) return 'default' - if (['executing', 'verifying'].includes(status)) return 'secondary' + if (['failed', 'denied', 'revoked', 'cancelled'].includes(status)) return 'destructive' + if (status === 'completed') return 'default' + if (['running', 'approved'].includes(status)) return 'secondary' return 'outline' } @@ -87,7 +116,7 @@ Approvals {#if pendingApprovals.length}{pendingApprovals.length}{/if} - Executions + Activity @@ -168,31 +197,39 @@ Target Action + Risk Status - Correlation - Started + Duration + When Actions - {#each executions as execution (execution.id)} + {#each activity as item (item.id)} - {execution.target ?? '—'} - {execution.action} - {execution.status} - {execution.correlation_id} - {execution.started_at ? new Date(execution.started_at).toLocaleString() : '—'} + {item.target ?? '—'} + +
{item.verb}
+ {#if item.summary} +
{item.summary}
+ {/if} + {#if item.error} +
{item.error}
+ {/if} +
+ {item.risk_class} + {item.status} + {fmtDuration(item.duration_ms)} + {fmtWhen(item.created_at)} - {#if ['proposed', 'approved', 'auto_approved', 'executing'].includes(execution.status)} - + {#if ['pending_approval', 'approved', 'running'].includes(item.status)} + {/if}
{:else} - No executions yet. + No activity yet. {/each}