Files
oikos/web/src/lib/stores/chat.ts
dtoro 873b00ac42
Some checks failed
ci / build-test (push) Has been cancelled
ci / docker-build (push) Has been cancelled
ci / web (push) Has been cancelled
Desktop App / Build Linux (amd64) (push) Has been cancelled
Desktop App / Attach to Release (push) Has been cancelled
style(web): fix prettier config, format entire web/ tree
.prettierrc.json was missing "semi": false, so prettier wanted to add
semicolons to a codebase written without them (763 semicolon-free
statements vs. 150 with, in hand-written .ts; zero hand-written .svelte
files use them at all). That's why prettier --check failed on 249 files
— not because the code was unformatted, but because the config didn't
match the actual house style. Added "semi": false; left printWidth/etc
as configured (printWidth barely moves the failure count: 218/213/212
files at 100/120/140).

Ran `prettier --write .` with the corrected config. Verified
semantics-preserving before and after:
- eslint: 142 problems both before and after, byte-identical
- build passes, 38/38 tests pass
- token-stream diff (whitespace/semicolons/quotes normalized) on all
  218 changed files: only 52 had any remaining token change, all either
  trailing-comma removal (matching trailingComma: "none") or import/
  ternary reflow — no semantic changes
- live smoke test: Knowledge, Tasks, Fleet map, and a chat window
  (AgentTrace, markdown, Scope graph, activity rail) all render
  correctly, no console errors

Most of the diff is shadcn/ui vendor files (lib/components/ui/) moving
from the CLI's own style (double quotes, tabs, semicolons) to house
style; re-running `shadcn-svelte add` on a component will need a
follow-up format pass.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-27 12:56:07 +02:00

843 lines
31 KiB
TypeScript

import { writable, get, type Writable } from 'svelte/store'
import {
streamChat,
fetchSessions,
fetchMessages,
fetchMessagesOrNotFound,
deleteSession as apiDeleteSession
} from '$lib/api'
import type { ChatEvent, Session, Message } from '$lib/api'
import type { ToolCallResult } from '$lib/types'
export type { ToolCallResult }
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[]
created_at?: string
}
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) {
const args = t.args ?? {}
const purpose = typeof args.purpose === 'string' ? args.purpose : undefined
out.push({
executionId: m[1],
action: purpose
? purpose.slice(0, 60)
: typeof args.action === 'string'
? args.action
: t.name,
target: typeof args.target === 'string' ? args.target : 'unknown',
destructive: /\bDESTRUCTIVE\b/.test(text),
command: typeof args.command === 'string' ? args.command : undefined,
purpose
})
}
}
return out
}
function mid(): string {
return crypto.randomUUID()
}
export const messages = writable<ChatMessage[]>([])
export const streaming = writable(false)
export const connectionState = writable<'connected' | 'disconnected' | 'reconnecting'>('connected')
export const currentSession = writable<string | null>(null)
export const sessions = writable<Session[]>([])
export const sessionMessages = writable<Message[]>([])
export const error = writable<string | null>(null)
export const chatErrors = writable<{ id: string; message: string; action?: string }[]>([])
export function dismissError(id: string) {
chatErrors.update((e) => e.filter((x) => x.id !== id))
}
export function addChatError(message: string, action?: string) {
chatErrors.update((e) => [...e, { id: crypto.randomUUID(), message, action }])
}
// 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<string, AbortController>()
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<string, ToolCallResult>()
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 content = typeof m.content === 'string' ? { text: m.content } : m.content
const tools = mergeToolCalls(content?.tool_calls)
return {
id: m.id,
role: m.role as 'user' | 'assistant',
text: content?.text ?? '',
tools,
pendingApprovals: extractApprovals(tools),
created_at: m.created_at
}
})
}
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<typeof setInterval> | null = null
let pollingSessionId: string | null = null
function startPolling(sessionId: string) {
stopPolling()
pollingSessionId = sessionId
pollTimer = setInterval(async () => {
// Allow polling while disconnected — the agent is still working
// server-side and the poller is the only way to see it.
if (get(streaming) && get(connectionState) === 'connected') 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])
const activeTools: Map<string, ToolCallResult> = 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
let receivedDone = false
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') {
ms[ms.length - 1] = { ...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') {
const tools = last.tools.map((t) => (t.id === ev.data.id ? updated : t))
ms[ms.length - 1] = { ...last, tools, pendingApprovals: extractApprovals(tools) }
}
return [...ms]
})
}
} else if (ev.type === 'text_delta') {
messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, text: 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') {
ms[ms.length - 1] = { ...last, text: ev.data }
}
return [...ms]
})
} else if (ev.type === 'done') {
receivedDone = true
connectionState.set('connected')
messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...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) => {
// Distinguish user abort from network drop.
if (err === 'AbortError' || err.includes('aborted')) {
if (get(currentSession) === streamSessionID) streaming.set(false)
return
}
// Network blip / server restart — initiate reconnect.
if (get(currentSession) === streamSessionID) {
error.set(err)
if (!receivedDone && streamSessionID) {
handleDisconnect(streamSessionID)
} else {
streaming.set(false)
}
}
},
() => {
// SSE stream completed without error. If we never received 'done',
// the connection was severed mid-turn — treat as disconnect.
if (!receivedDone && streamSessionID && get(currentSession) === streamSessionID) {
handleDisconnect(streamSessionID)
} else if (get(currentSession) === streamSessionID) {
streaming.set(false)
}
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
}
}
// handleDisconnect is called when the SSE stream drops mid-turn without
// receiving a 'done' event. Falls back to polling and attempts reconnection.
function handleDisconnect(sessionId: string) {
const MAX_RECONNECT = 3
connectionState.set('disconnected')
startPolling(sessionId)
addChatError('Agent connection lost. The task is still running — retrying…', 'Dismiss')
let attempts = 0
let delay = 1000
const attemptReconnect = () => {
if (get(currentSession) !== sessionId || attempts >= MAX_RECONNECT) {
connectionState.set('disconnected')
streaming.set(false)
return
}
if (attempts > 0) {
connectionState.set('reconnecting')
addChatError(`Reconnecting to agent (attempt ${attempts + 1}/${MAX_RECONNECT})…`, 'Dismiss')
}
attempts++
const controller = streamChat(
'',
sessionId,
(_ev: ChatEvent) => {},
(_err: string) => {
delay = Math.min(delay * 2, 8000)
setTimeout(attemptReconnect, delay)
},
() => {
if (get(currentSession) === sessionId) {
connectionState.set('connected')
streaming.set(false)
loadSessionMessages(sessionId)
}
}
)
if (activeControllers.get(sessionId)) {
activeControllers.get(sessionId)?.abort()
}
activeControllers.set(sessionId, controller)
}
setTimeout(attemptReconnect, delay)
}
export function reconnect() {
const sid = get(currentSession)
if (!sid) return
connectionState.set('reconnecting')
const controller = streamChat(
'',
sid,
(_ev: ChatEvent) => {},
(_err: string) => {
connectionState.set('disconnected')
addChatError(
'Reconnect failed. The task may still be running — try sending a message to wake the agent.',
'Dismiss'
)
},
() => {
if (get(currentSession) === sid) {
connectionState.set('connected')
streaming.set(false)
loadSessionMessages(sid)
}
}
)
if (activeControllers.get(sid)) {
activeControllers.get(sid)?.abort()
}
activeControllers.set(sid, controller)
}
export function newChat() {
cancelStream()
stopPolling()
connectionState.set('connected')
currentSession.set(null)
messages.set([])
error.set(null)
chatErrors.set([])
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()
}
// ─── per-session chat state, for floating task windows ─────────────────────
//
// Everything above this point is the single "whatever's on screen" view used
// by the main Chat page and the chat drawer — one global `currentSession`,
// one `messages` array, guarded so a background stream never clobbers the
// view. Floating task windows break that assumption: several sessions can be
// open and legitimately streaming at once, each wanting its own live
// transcript. Rather than retrofit the guard-heavy logic above (streamed
// events checking `get(currentSession) === streamSessionID` before applying),
// each window gets its own isolated store bundle keyed by session id, so
// there's nothing to guard — events for session X always land in X's own
// bundle regardless of what else is open or on screen.
export interface SessionChatState {
messages: Writable<ChatMessage[]>
streaming: Writable<boolean>
connectionState: Writable<'connected' | 'disconnected' | 'reconnecting'>
error: Writable<string | null>
// Set by loadSessionChat when the backend 404s the session outright
// (deleted, or an id that was never valid — a stale persisted window, a
// bad deep link). Distinct from a merely-empty transcript, which is the
// normal state for a session that exists but hasn't sent a message yet.
notFound: Writable<boolean>
}
const sessionChats = new Map<string, SessionChatState>()
const sessionPollers = new Map<string, ReturnType<typeof setInterval>>()
// Lazily creates (and memoizes) the store bundle for a session — call this to
// get the stores to subscribe to; it does not fetch anything.
export function chatFor(sessionId: string): SessionChatState {
let c = sessionChats.get(sessionId)
if (!c) {
c = {
messages: writable([]),
streaming: writable(false),
connectionState: writable('connected'),
error: writable(null),
notFound: writable(false)
}
sessionChats.set(sessionId, c)
}
return c
}
function startSessionPolling(sessionId: string) {
const existing = sessionPollers.get(sessionId)
if (existing) clearInterval(existing)
const chat = chatFor(sessionId)
sessionPollers.set(
sessionId,
setInterval(async () => {
if (get(chat.streaming) && get(chat.connectionState) === 'connected') return
const msgs = await fetchMessages(sessionId)
if (get(chat.streaming)) return // re-check: the fetch itself takes time
chat.messages.set(toChatMessages(msgs))
}, 3000)
)
}
export function stopSessionPolling(sessionId: string) {
const t = sessionPollers.get(sessionId)
if (t) {
clearInterval(t)
sessionPollers.delete(sessionId)
}
}
// Fetches sessionId's current transcript into its own store bundle and
// starts polling it for auto-continuation updates — the per-session
// equivalent of loadSessionMessages, for a window rather than the main view.
export async function loadSessionChat(sessionId: string): Promise<void> {
const chat = chatFor(sessionId)
chat.streaming.set(false)
const msgs = await fetchMessagesOrNotFound(sessionId)
if (msgs === null) {
chat.notFound.set(true)
return // nothing to poll — the session doesn't exist
}
chat.messages.set(toChatMessages(msgs))
startSessionPolling(sessionId)
}
// Per-session equivalent of sendMessage — writes into sessionId's own store
// bundle unconditionally (no "is this still on screen" guard needed, since
// the bundle IS the screen for this session's window) and shares
// `activeControllers` with the singleton path above so cancelStream() from
// either a window or the main view (if the same session happens to be open
// in both) finds the same in-flight call.
export function sendSessionMessage(sessionId: string, text: string) {
const chat = chatFor(sessionId)
chat.error.set(null)
chat.streaming.set(true)
const userMsg: ChatMessage = { id: mid(), role: 'user', text, tools: [], pendingApprovals: [] }
chat.messages.update((ms) => [...ms, userMsg])
const assistantMsg: ChatMessage = {
id: mid(),
role: 'assistant',
text: '',
tools: [],
pendingApprovals: []
}
chat.messages.update((ms) => [...ms, assistantMsg])
const activeTools: Map<string, ToolCallResult> = new Map()
let receivedDone = false
const controller = streamChat(
text,
sessionId,
(ev: ChatEvent) => {
if (ev.type === 'session') return // sessionId is already known for a window
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)
chat.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...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)
chat.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
const tools = last.tools.map((t) => (t.id === ev.data.id ? updated : t))
ms[ms.length - 1] = { ...last, tools, pendingApprovals: extractApprovals(tools) }
}
return [...ms]
})
}
} else if (ev.type === 'text_delta') {
chat.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, text: last.text + ev.data }
}
return [...ms]
})
} else if (ev.type === 'text') {
chat.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, text: ev.data }
}
return [...ms]
})
} else if (ev.type === 'done') {
receivedDone = true
chat.connectionState.set('connected')
chat.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, pendingApprovals: extractApprovals(last.tools) }
}
return [...ms]
})
startSessionPolling(sessionId)
} else if (ev.type === 'error') {
chat.error.set(ev.data)
}
},
(err: string) => {
if (err === 'AbortError' || err.includes('aborted')) {
chat.streaming.set(false)
return
}
chat.error.set(err)
if (!receivedDone) {
chat.connectionState.set('disconnected')
startSessionPolling(sessionId)
addChatError('Agent connection lost. The task is still running — retrying…', 'Dismiss')
} else {
chat.streaming.set(false)
}
},
() => {
chat.streaming.set(false)
if (activeControllers.get(sessionId) === controller) activeControllers.delete(sessionId)
loadSessions()
}
)
activeControllers.set(sessionId, controller)
}
export function cancelSessionStream(sessionId: string) {
const controller = activeControllers.get(sessionId)
if (!controller) return
controller.abort()
activeControllers.delete(sessionId)
chatFor(sessionId).streaming.set(false)
}
// ─── new-task launcher (desktop center input / Tasks app) ───────────────────
//
// Starting a brand-new task has no session id to hang a window off of until
// the stream's own 'session' event assigns one (see the 'session' branch in
// sendMessage above) — the desktop launcher needs to open that task's window
// the moment an id exists, not before. startTask begins the stream
// immediately, buffers any events that arrive before 'session' (defensive:
// in practice 'session' always arrives first), then seeds that session's own
// chatFor() bundle exactly like sendSessionMessage does and hands the id back
// via onSession so the caller can open its window. From that point on the
// window behaves exactly like any other task window.
export function startTask(text: string, onSession: (sessionId: string) => void): void {
const userMsg: ChatMessage = { id: mid(), role: 'user', text, tools: [], pendingApprovals: [] }
const assistantMsg: ChatMessage = {
id: mid(),
role: 'assistant',
text: '',
tools: [],
pendingApprovals: []
}
const activeTools: Map<string, ToolCallResult> = new Map()
let receivedDone = false
let sessionId: string | null = null
let chat: SessionChatState | null = null
const buffered: ChatEvent[] = []
function apply(ev: ChatEvent) {
const c = chat
if (!c || !sessionId) return
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)
c.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...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)
c.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
const tools = last.tools.map((t) => (t.id === ev.data.id ? updated : t))
ms[ms.length - 1] = { ...last, tools, pendingApprovals: extractApprovals(tools) }
}
return [...ms]
})
}
} else if (ev.type === 'text_delta') {
c.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, text: last.text + ev.data }
}
return [...ms]
})
} else if (ev.type === 'text') {
c.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, text: ev.data }
}
return [...ms]
})
} else if (ev.type === 'done') {
receivedDone = true
c.connectionState.set('connected')
c.messages.update((ms) => {
const last = ms[ms.length - 1]
if (last && last.role === 'assistant') {
ms[ms.length - 1] = { ...last, pendingApprovals: extractApprovals(last.tools) }
}
return [...ms]
})
startSessionPolling(sessionId)
} else if (ev.type === 'error') {
c.error.set(ev.data)
}
}
const controller = streamChat(
text,
null,
(ev: ChatEvent) => {
if (ev.type === 'session') {
sessionId = ev.data
activeControllers.set(sessionId, controller)
chat = chatFor(sessionId)
chat.streaming.set(true)
chat.messages.update((ms) => [...ms, userMsg, assistantMsg])
onSession(sessionId)
for (const b of buffered.splice(0)) apply(b)
return
}
if (!chat) {
buffered.push(ev)
return
}
apply(ev)
},
(err: string) => {
if (!chat) return // never got a session id — nothing to show the error in
if (err === 'AbortError' || err.includes('aborted')) {
chat.streaming.set(false)
return
}
chat.error.set(err)
if (!receivedDone && sessionId) {
chat.connectionState.set('disconnected')
startSessionPolling(sessionId)
addChatError('Agent connection lost. The task is still running — retrying…', 'Dismiss')
} else {
chat.streaming.set(false)
}
},
() => {
if (chat) chat.streaming.set(false)
if (sessionId && activeControllers.get(sessionId) === controller)
activeControllers.delete(sessionId)
loadSessions()
}
)
}