import { writable, derived, get } from 'svelte/store' import { liveEvents, subscribeEvents } from './events' import { currentSession, sessions, loadSessions } from './chat' import { fetchPlan, fetchQuestions, type PlanStep, type SessionQuestion } from '$lib/api' // workspace.ts is the live "what is this task doing right now" surface for the // TaskContextPanel: plan progress, the pinned operator question, and entities // the agent is touching or whose health just changed. It is deliberately driven // by the ALWAYS-ON global events stream (subscribeEvents), not the per-turn // chat SSE — the auto-continuation worker and resumeSession run entirely // server-side with no chat turn open, so a chat-bound panel would go stale // exactly when the agent is working autonomously. This also means the panel // keeps updating across a tab reload: hydrate() re-fetches REST state, then // live events carry deltas from there. export const planSteps = writable([]) export const questions = writable([]) export const openQuestion = derived(questions, (qs) => qs.find((q) => q.status === 'open') ?? null) export interface TouchedEntity { slug: string tool: string ts: number } export const touched = writable([]) const TOUCHED_MAX = 12 const TOUCHED_PULSE_MS = 6000 export interface HealthDiff { slug: string from: string to: string ts: number } export const healthDiffs = writable([]) const HEALTH_DIFF_MS = 8000 // The task's own fields (goal/status/outcome/summary) live on the session row. // Rather than a dedicated endpoint, derive from the sessions list (already // fetched for the task board) and keep it fresh here on task-lifecycle events. export const currentTask = derived([sessions, currentSession], ([$sessions, $id]) => $sessions.find((s) => s.id === $id) ?? null ) // Events that can change agent_sessions.status/goal/outcome — see applyEvent. const STATUS_AFFECTING = new Set([ 'goal.set', 'task.status', 'plan.proposed', 'question.raised', 'question.answered' ]) let hydratedFor: string | null = null let unsubStream: (() => void) | null = null let unsubLive: (() => void) | null = null let refreshTimer: ReturnType | null = null let lastSeenId = 0 async function hydrate(sessionId: string) { hydratedFor = sessionId planSteps.set([]) questions.set([]) touched.set([]) healthDiffs.set([]) const [steps, qs] = await Promise.all([fetchPlan(sessionId), fetchQuestions(sessionId)]) if (get(currentSession) !== sessionId) return // switched away while loading planSteps.set(steps) questions.set(qs) } function applyPlanStepEvent(sessionId: string, type: string, data: any) { const stepID = data?.step_id as string | undefined const seq = data?.seq as number | undefined planSteps.update((steps) => { const i = steps.findIndex((s) => (stepID && s.id === stepID) || (seq != null && s.seq === seq)) if (i === -1) return steps const next = [...steps] next[i] = { ...next[i], status: data.status ?? next[i].status, execution_id: data.execution_id ?? next[i].execution_id } return next }) } function applyEvent(ev: { type: string; correlation_id?: string | null; data?: unknown }) { const sid = get(currentSession) if (!sid || ev.correlation_id !== sid) return const data = (ev.data ?? {}) as any // Task fields (status/goal/outcome) live on the session row — refetch the // (cheap) session list so GoalHeader picks up the change without a // dedicated endpoint. Every event that can change agent_sessions.status // (goal.set → planning, propose_plan → executing, ask_operator → // awaiting_input, answerQuestion → executing, complete_task → done/failed) // must trigger this, not just goal.set/task.status — otherwise the status // pill goes stale exactly when resumeSession runs the next turn entirely // server-side, with no client-streaming 'done' event to piggyback a refresh // on (found live: answering a question via the panel left the header stuck // on "Needs your input" after the agent had already resumed). Debounced // since several of these can land in one burst. if (STATUS_AFFECTING.has(ev.type)) { if (refreshTimer) clearTimeout(refreshTimer) refreshTimer = setTimeout(() => loadSessions(), 300) } switch (ev.type) { case 'plan.proposed': if (Array.isArray(data.steps)) { const incoming = data.steps.map((s: any) => ({ id: s.id, seq: s.seq, title: s.title, detail: s.detail ?? '', status: 'pending' as const, target_slug: s.target_slug || undefined })) // The server appends rather than replaces once any step has started // (see store.go proposePlan) — mirror that here so a model that calls // propose_plan once per step still shows the FULL running history in // the panel, not just its latest call's single step. planSteps.update((existing) => (data.appended ? [...existing, ...incoming] : incoming)) } break case 'plan.step.started': case 'plan.step.finished': applyPlanStepEvent(sid, ev.type, data) break case 'question.raised': questions.update((qs) => [ { id: data.question_id, prompt: data.prompt ?? '', context: { why: data.why, options: data.options, entities: data.entities }, status: 'open', created_at: new Date().toISOString() }, ...qs.filter((q) => q.id !== data.question_id) ]) break case 'question.answered': questions.update((qs) => qs.map((q) => (q.id === data.question_id ? { ...q, status: 'answered', answer: data.answer } : q)) ) break case 'entity.touched': if (data.slug) { const now = Date.now() touched.update((t) => [{ slug: data.slug, tool: data.tool ?? '', ts: now }, ...t].slice(0, TOUCHED_MAX)) } break case 'knowledge.recorded': // No dedicated store yet — the outcome/knowledge card reads this task's // digest (fetchSessionDigest) on completion, which already lists it. break } } // health.changed is task-agnostic (fleet-wide), so it's matched separately: // show the diff whenever the changed entity is one this task has touched, not // by correlation_id (health events don't carry one). function applyHealthChanged(ev: { type: string; data?: unknown }) { if (ev.type !== 'health.changed') return const data = (ev.data ?? {}) as any if (!data.slug) return const isRelevant = get(touched).some((t) => t.slug === data.slug) if (!isRelevant) return healthDiffs.update((d) => [{ slug: data.slug, from: data.from, to: data.to, ts: Date.now() }, ...d].slice(0, 8)) } // startWorkspace opens the global event subscription and begins tracking the // active session. Call once from the panel's onMount; call the returned // cleanup on unmount. Safe to call multiple times (ref-counted underneath). export function startWorkspace(): () => void { unsubStream = subscribeEvents() const unsubSession = currentSession.subscribe((sid) => { if (sid && sid !== hydratedFor) hydrate(sid) if (!sid) { hydratedFor = null planSteps.set([]) questions.set([]) touched.set([]) healthDiffs.set([]) } }) unsubLive = liveEvents.subscribe((evs) => { if (evs.length === 0) return const maxId = evs[0].id if (maxId <= lastSeenId) { return } const fresh = evs.filter((e) => e.id > lastSeenId) lastSeenId = maxId // Oldest-first application so ordering (e.g. plan.step.started before // .finished) is preserved. for (const e of fresh.slice().reverse()) { applyEvent(e) applyHealthChanged(e) } }) return () => { unsubSession() unsubLive?.() unsubStream?.() } } // Sweep expired pulses/diffs on an interval so old touches stop glowing. setInterval(() => { const now = Date.now() touched.update((t) => t.filter((e) => now - e.ts < TOUCHED_PULSE_MS)) healthDiffs.update((d) => d.filter((e) => now - e.ts < HEALTH_DIFF_MS)) }, 1000)