import { writable, get } from 'svelte/store' import { streamChat, fetchSessions, fetchMessages, deleteSession as apiDeleteSession } from '$lib/api' import type { ChatEvent, Session, Message } from '$lib/api' export interface PendingApproval { executionId: string action: string target: string destructive: boolean command?: string purpose?: string } export interface ChatMessage { id: string role: 'user' | 'assistant' text: string tools: ToolCallResult[] pendingApprovals: PendingApproval[] } const APPROVAL_RE = /execution\s+([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})/i // Deliberately NOT filtered by tool name. There is no fixed set of gated // tools — `run` can execute anything, and any future tool that queues an // approval should surface a card the same way. A prior version hardcoded // `t.name === 'request_execution'`, so approvals raised by the newer `run` // tool were silently invisible in chat: no card, no feedback, nothing to // self-heal, forcing the operator to the Ops page with zero acknowledgement // back in the conversation. Matching on the response shape (not the tool // name) is what makes this robust to new gated tools without another // silent breakage. function extractApprovals(tools: ToolCallResult[]): PendingApproval[] { const out: PendingApproval[] = [] for (const t of tools) { if (t.type !== 'tool_result') continue const text = typeof t.result === 'string' ? t.result : JSON.stringify(t.result ?? '') if (!text.includes('requires approval')) continue const m = text.match(APPROVAL_RE) if (m) { out.push({ executionId: m[1], action: t.args?.action ?? t.args?.purpose ?? t.name ?? 'unknown', target: t.args?.target ?? 'unknown', destructive: /\bDESTRUCTIVE\b/.test(text), command: t.args?.command, purpose: t.args?.purpose }) } } return out } export interface ToolCallResult { type: 'tool_use' | 'tool_result' name: string id?: string args?: any result?: any error?: string } function mid(): string { return crypto.randomUUID() } export const messages = writable([]) export const streaming = writable(false) export const currentSession = writable(null) export const sessions = writable([]) export const sessionMessages = writable([]) export const error = writable(null) // Per-session controller tracking. Multiple tasks can stream concurrently // (see sendMessage's session guard above this used to be a single global // `activeController`, which meant cancelStream()/newChat() always aborted // whichever stream happened to be the MOST RECENTLY started one, regardless // of what the operator was currently viewing — starting Task A, switching to // (already-loaded) Task B, then clicking "New task" would silently abort // Task A's still-running turn even though the operator was never looking at // it and never asked to cancel it. Keyed by session id once known; // pendingController covers the brief window for a brand-new task between // streamChat() starting and its 'session' event assigning a real id. const activeControllers = new Map() let pendingController: AbortController | null = null export async function loadSessions() { const list = await fetchSessions() sessions.set(list) } // mergeToolCalls collapses a persisted tool_calls array into one entry per // call id. Nomos persists the tool_use and tool_result as two separate // entries sharing the same id (matching the SSE event pair); Chat.svelte // renders tools in a keyed {#each ... (tool.id)}, which throws on duplicate // keys and silently aborts the whole message list. Live-streamed messages // never hit this because sendMessage() merges tool_result into the existing // tool_use entry in place rather than appending a second one. function mergeToolCalls(raw: ToolCallResult[] | undefined): ToolCallResult[] { const byId = new Map() for (const tc of raw ?? []) { const key = tc.id ?? crypto.randomUUID() const existing = byId.get(key) byId.set(key, existing ? { ...existing, ...tc, id: key } : { ...tc, id: key }) } return Array.from(byId.values()) } function toChatMessages(msgs: Message[]): ChatMessage[] { return msgs.map((m) => { const tools = mergeToolCalls(m.content?.tool_calls) return { id: m.id, role: m.role as 'user' | 'assistant', text: m.content?.text ?? (typeof m.content === 'string' ? m.content : ''), tools, pendingApprovals: extractApprovals(tools) } }) } export async function loadSessionMessages(sessionId: string) { currentSession.set(sessionId) // This is a fresh view of sessionId's current (REST-loaded) state — reset // streaming regardless of whether some OTHER task's stream happens to still // be in flight in the background. Without this, switching to a task while // a different one is mid-turn could leave `streaming` stuck true here (that // other stream's completion callback now correctly skips touching it, per // sendMessage's session guard) — which would disable the input AND silently // stop startPolling's loop from ever applying updates (it bails while // $streaming is true), making the newly-opened task look frozen. streaming.set(false) const msgs = await fetchMessages(sessionId) sessionMessages.set(msgs) messages.set(toChatMessages(msgs)) startPolling(sessionId) } // Live visibility for autonomous work: the auto-continuation worker (see // cmd/nomos/continue.go) runs entirely server-side and has no live push — // previously the only way to see its result was to manually reload the // session, so approving a plan and then waiting felt like nothing was // happening even while the agent was actively working. This polls the // session's persisted messages every few seconds and merges in anything new // (an auto-continuation's result, a fresh pending approval it queued, etc.) // so the transcript updates on its own. Only runs between turns — never // while a live streaming turn owns the message list, to avoid clobbering the // in-progress optimistic UI. let pollTimer: ReturnType | null = null let pollingSessionId: string | null = null function startPolling(sessionId: string) { stopPolling() pollingSessionId = sessionId pollTimer = setInterval(async () => { if (get(streaming)) return if (pollingSessionId !== sessionId || get(currentSession) !== sessionId) return const msgs = await fetchMessages(sessionId) if (get(streaming) || pollingSessionId !== sessionId) return // re-check: the fetch itself takes time // No cheap "anything new?" check: the auto-continuation worker updates a // placeholder message IN PLACE as each tool call lands (see // cmd/nomos/continue.go), so the message COUNT stays the same while the // content changes — a length-only diff (the previous version of this // code) never detected those updates and progress looked frozen even // though the backend was actively working. Just re-set every tick; // Svelte's own diffing keeps the actual re-render cheap. sessionMessages.set(msgs) messages.set(toChatMessages(msgs)) }, 3000) } export function stopPolling() { if (pollTimer) { clearInterval(pollTimer) pollTimer = null } pollingSessionId = null } export function sendMessage(text: string) { error.set(null) streaming.set(true) const userMsg: ChatMessage = { id: mid(), role: 'user', text, tools: [], pendingApprovals: [] } messages.update((ms) => [...ms, userMsg]) const assistantMsg: ChatMessage = { id: mid(), role: 'assistant', text: '', tools: [], pendingApprovals: [] } messages.update((ms) => [...ms, assistantMsg]) let activeTools: Map = new Map() // Multiple tasks can stream concurrently (the backend runs each turn as its // own goroutine — nothing serializes them), but `messages`/`currentSession` // are a single global view. Without this guard, switching to a different // task while this stream is still open lets its later events (tool_use, // text_delta, ..., and worst of all the 'done' handler's // currentSession.set) get applied to whatever the operator is NOW looking // at — silently corrupting another task's transcript, or yanking the view // back to this one. openedFor is the session this call started for (null // for a brand-new task, until the 'session' event assigns the real id); // every branch below checks the CURRENT $currentSession still matches // before touching `messages`. The task itself keeps running server-side // regardless — dropped events just mean the live view isn't watching it; // navigating back re-hydrates via REST/poll same as it already does for // auto-continuation. const openedFor = get(currentSession) let streamSessionID = openedFor const controller = streamChat( text, get(currentSession), // continue the active session so the agent keeps context (ev: ChatEvent) => { if (ev.type === 'session') { streamSessionID = ev.data // Move this stream's controller into the per-session map now that its // real id is known, so a later cancelStream()/newChat() from THIS // session's view can find and abort it — and, just as importantly, // so cancelling/leaving a DIFFERENT session never reaches this one. // For a continued (non-new) session, openedFor already equals ev.data // and the controller was stored under that key at creation below; // this only does real work for a brand-new task's first assignment. if (pendingController === controller) pendingController = null activeControllers.set(ev.data, controller) // Only claim currentSession if the operator hasn't already navigated // to something else since this call started (openedFor covers both // "still on the task I was on" and "still hadn't opened one yet"). if (get(currentSession) === openedFor) currentSession.set(ev.data) return } if (get(currentSession) !== streamSessionID) return // stream's task isn't the one on screen — drop if (ev.type === 'tool_use') { const tr: ToolCallResult = { type: 'tool_use', name: ev.data.name, id: ev.data.id, args: ev.data.args } activeTools.set(ev.data.id, tr) messages.update((ms) => { const last = ms[ms.length - 1] if (last && last.role === 'assistant') { last.tools = [...last.tools, tr] } return [...ms] }) } else if (ev.type === 'tool_result') { const existing = activeTools.get(ev.data.id) if (existing) { const updated: ToolCallResult = { ...existing, type: 'tool_result', result: ev.data.result, error: ev.data.error } activeTools.set(ev.data.id, updated) messages.update((ms) => { const last = ms[ms.length - 1] if (last && last.role === 'assistant') { last.tools = last.tools.map((t) => t.id === ev.data.id ? updated : t ) } return [...ms] }) } } else if (ev.type === 'text_delta') { messages.update((ms) => { const last = ms[ms.length - 1] if (last && last.role === 'assistant') { last.text += ev.data } return [...ms] }) } else if (ev.type === 'text') { // Final authoritative content for the turn; replaces accumulated deltas. messages.update((ms) => { const last = ms[ms.length - 1] if (last && last.role === 'assistant') { last.text = ev.data } return [...ms] }) } else if (ev.type === 'done') { messages.update((ms) => { const last = ms[ms.length - 1] if (last && last.role === 'assistant') { last.pendingApprovals = extractApprovals(last.tools) } return [...ms] }) const sid = ev.data?.session_id ?? ev.session_id // Start polling for auto-continuation results now that the live turn // is over — this is what makes an approved plan's later steps show up // on their own instead of requiring a manual reload. (startPolling's // own loop already re-checks $currentSession before applying results, // so this is safe to call even if the operator has since navigated // elsewhere — it just won't visibly do anything until/unless they // come back.) if (sid) startPolling(sid) } else if (ev.type === 'error') { error.set(ev.data) } }, (err: string) => { if (get(currentSession) === streamSessionID) error.set(err) }, () => { if (get(currentSession) === streamSessionID) streaming.set(false) // Clean up whichever slot this controller ended up in — normally // activeControllers[streamSessionID] once the 'session' event has // fired, but fall back to pendingController for the (rare) case where // the stream errored/completed before ever getting one. if (streamSessionID && activeControllers.get(streamSessionID) === controller) { activeControllers.delete(streamSessionID) } if (pendingController === controller) pendingController = null loadSessions() } ) // Register immediately (not just inside the 'session' handler above) so a // cancelStream() during the brief pre-'session' window for a CONTINUED // session (openedFor already known) can find it right away. if (openedFor) { activeControllers.set(openedFor, controller) } else { pendingController = controller } } export function newChat() { cancelStream() stopPolling() currentSession.set(null) messages.set([]) error.set(null) streaming.set(false) // fresh view — see loadSessionMessages for why this must not depend on cancelStream's own reset } // Cancels the stream for whatever the operator is CURRENTLY VIEWING — never // some other, unrelated task's background stream. Before per-session // tracking, this aborted a single global `activeController`, which meant it // always targeted the MOST RECENTLY STARTED stream regardless of what was on // screen: start Task A, switch to already-loaded Task B, click "New task" — // newChat()'s cancelStream() would silently abort Task A's still-running // turn, even though the operator was never looking at it and never asked to // cancel it. Now it looks up by $currentSession (or pendingController for // the brief pre-'session'-event window of a just-started new task) so it can // only ever touch the stream that belongs to the view being left. export function cancelStream() { const sid = get(currentSession) const controller = sid ? activeControllers.get(sid) : pendingController if (!controller) return controller.abort() if (sid) activeControllers.delete(sid) if (pendingController === controller) pendingController = null streaming.set(false) } export async function deleteSession(sessionId: string) { const ok = await apiDeleteSession(sessionId) if (!ok) return if (get(currentSession) === sessionId) { newChat() } loadSessions() }