diff --git a/bun.lock b/bun.lock index ec3cc56..1967539 100644 --- a/bun.lock +++ b/bun.lock @@ -5,6 +5,7 @@ "": { "name": "@ellipsis/cli", "dependencies": { + "@ellipsis-dev/sdk": "^0.1.0", "commander": "^12.1.0", "ink": "^5.0.1", "ink-spinner": "^5.0.0", @@ -27,6 +28,8 @@ "packages": { "@alcalzone/ansi-tokenize": ["@alcalzone/ansi-tokenize@0.1.3", "", { "dependencies": { "ansi-styles": "^6.2.1", "is-fullwidth-code-point": "^4.0.0" } }, "sha512-3yWxPTq3UQ/FY9p1ErPxIyfT64elWaMvM9lIHnaqpyft63tkxodF5aUElYHrdisWve5cETkh1+KBw1yJuW0aRw=="], + "@ellipsis-dev/sdk": ["@ellipsis-dev/sdk@0.1.0", "", {}, "sha512-bvIipl/bdDF1v+HUJuJ5YSl1QZM7h4rEpYOs5qKtH/tzay93lMnexqhZ6tQALSrtvMwgY5h7lpYDlYtLBjB3pg=="], + "@esbuild/aix-ppc64": ["@esbuild/aix-ppc64@0.27.7", "", { "os": "aix", "cpu": "ppc64" }, "sha512-EKX3Qwmhz1eMdEJokhALr0YiD0lhQNwDqkPYyPhiSwKrh7/4KRjQc04sZ8db+5DVVnZ1LmbNDI1uAMPEUBnQPg=="], "@esbuild/android-arm": ["@esbuild/android-arm@0.27.7", "", { "os": "android", "cpu": "arm" }, "sha512-jbPXvB4Yj2yBV7HUfE2KHe4GJX51QplCN1pGbYjvsyCZbQmies29EoJbkEc+vYuU5o45AfQn37vZlyXy4YJ8RQ=="], @@ -383,10 +386,6 @@ "yoga-layout": ["yoga-layout@3.2.1", "", {}, "sha512-0LPOt3AxKqMdFBZA3HBAt/t/8vIKq7VaQYbuA8WxCgung+p9TVyKRYdpvCb80HcdTN2NkbIKbhNwKUfm3tQywQ=="], - "@vitest/runner/pathe": ["pathe@1.1.2", "", {}, "sha512-whLdWMYL2TwI08hn8/ZqAbrVemu0LNaNNJZX73O6qaIdCTfXutsLhMkjdENX0qhsQ9uIimo4/aQOmXkoon2nDQ=="], - - "@vitest/snapshot/pathe": ["pathe@1.1.2", "", {}, "sha512-whLdWMYL2TwI08hn8/ZqAbrVemu0LNaNNJZX73O6qaIdCTfXutsLhMkjdENX0qhsQ9uIimo4/aQOmXkoon2nDQ=="], - "cli-truncate/slice-ansi": ["slice-ansi@5.0.0", "", { "dependencies": { "ansi-styles": "^6.0.0", "is-fullwidth-code-point": "^4.0.0" } }, "sha512-FC+lgizVPfie0kkhqUScwRu1O/lF6NOgJmlCgK+/LYxDCTk8sGelYaHDhFcDN+Sn3Cv+3VSa4Byeo+IMCzpMgQ=="], "mlly/pathe": ["pathe@2.0.3", "", {}, "sha512-WUjGcAqP1gQacoQe+OBJsFA7Ld4DyXuUIjZ5cc75cLHvJ7dtNsTugphxIADwspS+AraAUePCKrSVtPLFj/F88w=="], @@ -403,8 +402,6 @@ "vite/esbuild": ["esbuild@0.21.5", "", { "optionalDependencies": { "@esbuild/aix-ppc64": "0.21.5", "@esbuild/android-arm": "0.21.5", "@esbuild/android-arm64": "0.21.5", "@esbuild/android-x64": "0.21.5", "@esbuild/darwin-arm64": "0.21.5", "@esbuild/darwin-x64": "0.21.5", "@esbuild/freebsd-arm64": "0.21.5", "@esbuild/freebsd-x64": "0.21.5", "@esbuild/linux-arm": "0.21.5", "@esbuild/linux-arm64": "0.21.5", "@esbuild/linux-ia32": "0.21.5", "@esbuild/linux-loong64": "0.21.5", "@esbuild/linux-mips64el": "0.21.5", "@esbuild/linux-ppc64": "0.21.5", "@esbuild/linux-riscv64": "0.21.5", "@esbuild/linux-s390x": "0.21.5", "@esbuild/linux-x64": "0.21.5", "@esbuild/netbsd-x64": "0.21.5", "@esbuild/openbsd-x64": "0.21.5", "@esbuild/sunos-x64": "0.21.5", "@esbuild/win32-arm64": "0.21.5", "@esbuild/win32-ia32": "0.21.5", "@esbuild/win32-x64": "0.21.5" }, "bin": "bin/esbuild" }, "sha512-mg3OPMV4hXywwpoDxu3Qda5xCKQi+vCTZq8S9J/EpkhB2HzKXq4SNFZE3+NK93JYxc8VMSep+lOUSC/RVKaBqw=="], - "vite-node/pathe": ["pathe@1.1.2", "", {}, "sha512-whLdWMYL2TwI08hn8/ZqAbrVemu0LNaNNJZX73O6qaIdCTfXutsLhMkjdENX0qhsQ9uIimo4/aQOmXkoon2nDQ=="], - "tsx/esbuild/@esbuild/aix-ppc64": ["@esbuild/aix-ppc64@0.28.1", "", { "os": "aix", "cpu": "ppc64" }, "sha512-Svl7tq8k/08+p6CXPpRjQ1fKX+1odH/BQbb48fV6fj3CWHhsoIOoY87w1oHXm0qEpkIK3ZfVgp0hed3XBXzXMQ=="], "tsx/esbuild/@esbuild/android-arm": ["@esbuild/android-arm@0.28.1", "", { "os": "android", "cpu": "arm" }, "sha512-0k2F129Xdio1TdJfzJ8sy1Q47vUD2NnwdhiAf7drUN1EBTfPf4hsFCtmMgu/6m8JSzsBrlmVjudMBQqOfG8usQ=="], diff --git a/package.json b/package.json index 3680e02..1a4edbb 100644 --- a/package.json +++ b/package.json @@ -20,6 +20,7 @@ "test:watch": "vitest" }, "dependencies": { + "@ellipsis-dev/sdk": "^0.1.0", "commander": "^12.1.0", "ink": "^5.0.1", "ink-spinner": "^5.0.0", diff --git a/src/commands/connect.ts b/src/commands/connect.ts index b4b7d45..7574931 100644 --- a/src/commands/connect.ts +++ b/src/commands/connect.ts @@ -1,12 +1,12 @@ import type { Command } from 'commander' import React from 'react' import { render } from 'ink' +import { SessionTranscriptStore } from '@ellipsis-dev/sdk/store' import { ApiClient } from '../lib/api' import { requireToken, resolveApiBase, resolveAppBase } from '../lib/config' -import { runAction, usdNumberFromMillicents } from '../lib/output' -import { foldCosts, isConnectVisibleRecord, recordToItems, type CCEvent } from '../lib/events' +import { runAction } from '../lib/output' import { sessionUrl } from '../lib/urls' -import { resolveWsBase } from '../lib/ws' +import { makeOpenSocket, resolveWsBase } from '../lib/stream' import { ConnectApp } from '../ui/ConnectApp' import type { AgentSession } from '../lib/types' @@ -93,7 +93,7 @@ export async function runConnect( ): Promise { const api = new ApiClient() const token = requireToken() - const wsBase = resolveWsBase(resolveApiBase()) + const openSocket = makeOpenSocket(token, resolveWsBase(resolveApiBase())) const [session, me] = await Promise.all([api.getAgentSession(sessionId), api.whoami()]) const c = connectability(session) @@ -104,29 +104,20 @@ export async function runConnect( const url = sessionUrl(resolveAppBase(), me.customer_login, sessionId) // No scrollback preamble: the app owns the whole surface, Claude Code-style. - // The banner (brand + version + session link) and the footer carry the - // session identity/status; a watch-only reason surfaces as the app's notice. + // The footer carries the session identity/status; a watch-only reason + // surfaces as the app's notice. - // Fetch the stored records to seed the transcript (unless --no-records), the - // live-refresh cursor (so live updates only append what's new), and the - // opening spend. Records are ordered by feed_seq (the shared transcript + - // lifecycle feed). Lifecycle rows are filtered except the sandbox-ready - // conversation note (isConnectVisibleRecord): connect shows the - // conversation, and the live activity line + footer carry session state. - const records = await api.getAgentSessionRecords(sessionId) - const ordered = [...records].sort((a, b) => a.feed_seq - b.feed_seq) - const initialMaxFeedSeq = ordered.reduce((m, s) => Math.max(m, s.feed_seq), 0) - const initialItems = showRecords - ? ordered - .filter(isConnectVisibleRecord) - .flatMap((st) => recordToItems(st, `s${st.feed_seq}`)) - : [] - const initialCost = foldCosts(ordered.map((st) => st.payload as CCEvent)) - // The server's ledger total at connect time; live updates arrive as - // cost_millicents on status/done frames. - const initialServerCostUsd = usdNumberFromMillicents( - session.cost_tokens + session.cost_sandbox_cpu + session.cost_sandbox_memory + session.cost_fee, - ) + // Seed ONE transcript store with the stored records and the fetched session + // — synthetic frames through the same ingest path the live stream uses, so + // the first paint is instant and streamSession resumes past the seeded + // cursor instead of replaying history. --no-records skips *rendering* the + // seeded history (minRenderFeedSeq), not re-streaming it. + const store = new SessionTranscriptStore() + const page = await api.getAgentSessionRecordsPage(sessionId) + const ordered = [...page.records].sort((a, b) => a.feed_seq - b.feed_seq) + if (ordered.length) store.ingest({ type: 'records_append', records: ordered }) + store.ingest({ type: 'messages', messages: page.messages ?? [] }) + store.ingest({ type: 'session', session }) // Written by the app when it exits because the conversation closed (terminal; // nothing left to reconnect to), so the detach sign-off below stays honest. @@ -134,18 +125,12 @@ export async function runConnect( const app = render( React.createElement(ConnectApp, { api, - token, sessionId, - wsBase, + store, + openSocket, canSend, - initialItems, - // Always advance the cursor past existing records: --no-records skips - // *rendering* history, not re-streaming it live. - initialMaxFeedSeq, - initialStatus: session.surface?.status ?? session.status, + minRenderFeedSeq: showRecords ? 0 : store.cursor, sessionUrl: url, - initialCost, - initialServerCostUsd, initialNotice: reason ?? null, // The session's one model, fixed at creation (backend tokens_model). model: typeof session.tokens_model === 'string' ? session.tokens_model : null, diff --git a/src/commands/session.tsx b/src/commands/session.tsx index ab2fcd4..06442ce 100644 --- a/src/commands/session.tsx +++ b/src/commands/session.tsx @@ -25,16 +25,17 @@ import { } from '../lib/args' import { sessionUrl } from '../lib/urls' import { - resolveWsBase, sessionStatusWord, streamSession, StreamUnavailableError, type StreamFrame, type StreamOutcome, -} from '../lib/ws' -import { isConnectVisibleRecord, recordToItems } from '../lib/events' +} from '@ellipsis-dev/sdk/stream' +import { isConnectVisibleRecord, recordToItems } from '@ellipsis-dev/sdk/store' +import { makeOpenSocket, resolveWsBase } from '../lib/stream' import type { AgentSession, + AgentSessionWire, AgentSessionSource, AgentSessionStatus, GithubAccountSnippet, @@ -808,7 +809,7 @@ export async function watchSessionStreaming( json?: boolean, ): Promise { const token = requireToken() - const wsBase = resolveWsBase(resolveApiBase()) + const openSocket = makeOpenSocket(token, resolveWsBase(resolveApiBase())) // Session frames are LWW snapshots resent on any change (cost ticks // included), so collapse to status-word transitions — both to keep the @@ -819,7 +820,7 @@ export async function watchSessionStreaming( const onFrame = (frame: StreamFrame) => { if (frame.type === 'session' || frame.type === 'snapshot') { const word = sessionStatusWord( - (frame as { session: AgentSession }).session, + (frame as unknown as { session: AgentSessionWire }).session, ) if (word === lastStatus) return lastStatus = word @@ -834,7 +835,7 @@ export async function watchSessionStreaming( let outcome: StreamOutcome try { - outcome = await streamSession({ token, sessionId, wsBase, onFrame }) + outcome = await streamSession({ sessionId, openSocket, onFrame }) } catch (err) { if (err instanceof StreamUnavailableError) { if (!json) { diff --git a/src/lib/api.ts b/src/lib/api.ts index 0799d9f..8a21d5d 100644 --- a/src/lib/api.ts +++ b/src/lib/api.ts @@ -61,10 +61,11 @@ import type { StartSandboxBuildRequest, } from './types' -// Thin REST client over the public `/v1` API. The typed request/response -// surface mirrors ellipsis/src/public_api/routers/v1/v1_router.py and will move -// to @ellipsis/sdk (generated from the backend OpenAPI spec) once that package -// exists; this CLI then imports it instead of hand-rolling types. +// Thin REST client over the public `/v1` API. The session-stream surface +// types come from @ellipsis-dev/sdk (generated from the backend's schema, via +// lib/types re-exports); the rest of the typed surface remains a hand-rolled +// mirror of ellipsis/src/public_api/routers/v1/v1_router.py until the SDK's +// OpenAPI surface widens beyond the protocol endpoints. export class ApiError extends Error { constructor( @@ -208,11 +209,23 @@ export class ApiClient { // The session's full stored transcript as native session_records (transcript // + lifecycle), ordered by feed_seq. async getAgentSessionRecords(sessionId: string): Promise { - const res = await this.request( - 'GET', - `/v1/sessions/${encodeURIComponent(sessionId)}/records`, - ) - return res.records + const res = await this.getAgentSessionRecordsPage(sessionId) + // The OpenAPI response type marks defaulted fields optional; on the wire + // the server always serializes every field (the frames-schema flavor). + return res.records as SessionRecord[] + } + + // The full records response (records + the open inbox slice + + // has_more/earliest_feed_seq), optionally resuming past a feed_seq cursor + // (protocol §4.3) — what the connect UI's REST poll fallback feeds its + // transcript store from. + getAgentSessionRecordsPage( + sessionId: string, + options: { afterSeq?: number } = {}, + ): Promise { + const query = + options.afterSeq != null && options.afterSeq > 0 ? `?after_seq=${options.afterSeq}` : '' + return this.request('GET', `/v1/sessions/${encodeURIComponent(sessionId)}/records${query}`) } // The session's conversation structure — turns and inbox messages, each diff --git a/src/lib/events.ts b/src/lib/events.ts deleted file mode 100644 index 2f7eb2d..0000000 --- a/src/lib/events.ts +++ /dev/null @@ -1,378 +0,0 @@ -// Turn the semantic event stream `agent session connect` receives into -// structured, renderable transcript items — the parsing/shaping layer behind -// the Claude-Code-like connect UI (src/ui/ConnectApp.tsx). -// -// The session WebSocket relays the agent's Claude Code stream-json: one JSON -// object per line, carried in `stdout` frames (see docs/RUN_STREAMING_SPEC.md -// and src/lib/ws.ts). Each object is a `CCEvent` — a discriminated union of -// assistant turns (text / thinking / tool_use blocks), user turns (tool -// results, or a human message), a system init, and a final result. This module -// is pure (no ANSI, no Ink) so it can be unit-tested directly; colours and -// layout live in the UI component. - -import { lifecycleText, oneLine } from './steps' -import type { SessionRecord } from './types' - -// A loosely-typed content block of a Claude Code message. We only read the -// fields we display and never interpret the rest of the payload. -export interface CCContentBlock { - type?: string - text?: string - thinking?: string - id?: string - name?: string - input?: Record - content?: unknown - tool_use_id?: string - is_error?: boolean -} - -// One Claude Code stream-json event. `message.content` is a string or a list -// of blocks; `result` (with cost/duration) caps a turn. Everything the CLI -// doesn't render is passed through untyped. -export interface CCEvent { - type?: string - subtype?: string - message?: { role?: string; content?: unknown } - result?: string - duration_ms?: number - total_cost_usd?: number - num_turns?: number - is_error?: boolean - model?: string - cwd?: string - [key: string]: unknown -} - -// A single line of the rendered transcript. `kind` selects the colour/glyph in -// the UI; `gutter` is the leading marker (● tool, ⎿ result, › you, ✻ thinking). -// `spaceBefore` opens a blank line above to separate message-level blocks; -// grouped sub-lines (a tool's result under its call) set it false so they hug. -export type ItemKind = - | 'assistant' - | 'thinking' - | 'tool' - | 'tool_result' - | 'summary' - | 'system' - | 'user' - | 'notice' - | 'error' - -export interface TranscriptItem { - key: string - kind: ItemKind - text: string - // Secondary, dimmed text shown after `text` on the same logical block (a - // tool call's argument summary, a result's body). - detail?: string - gutter?: string - spaceBefore?: boolean - isError?: boolean -} - -// Reassembles the byte chunks of `stdout`/`stderr` frames into whole lines: a -// single JSON event can be split across frames, so we only emit a line once its -// terminating newline has arrived. `flush()` releases any trailing partial at -// stream end. -export class LineBuffer { - private buf = '' - - push(chunk: string): string[] { - this.buf += chunk - const parts = this.buf.split('\n') - this.buf = parts.pop() ?? '' - return parts - } - - flush(): string[] { - const rest = this.buf.trim() - this.buf = '' - return rest ? [rest] : [] - } -} - -// Parse one relayed line into a CCEvent, or null if it isn't a JSON object (a -// blank keepalive line, or plain non-event text the caller renders verbatim). -export function parseEventLine(line: string): CCEvent | null { - const trimmed = line.trim() - if (!trimmed) return null - try { - const value = JSON.parse(trimmed) as unknown - if (value && typeof value === 'object' && !Array.isArray(value)) { - return value as CCEvent - } - } catch { - // Not JSON — the caller shows it as raw text. - } - return null -} - -// Claude-Code-style one-line summary of a tool call's arguments: the salient -// field for the common tools (a path, a command, a pattern), else compact JSON. -export function summarizeToolInput(name: string, input: unknown): string { - const args = (input && typeof input === 'object' ? input : {}) as Record - const str = (v: unknown): string | undefined => (typeof v === 'string' ? v : undefined) - const tool = name.toLowerCase() - - const path = str(args.file_path) ?? str(args.path) ?? str(args.notebook_path) - if (['read', 'write', 'edit', 'multiedit', 'notebookedit'].includes(tool) && path) { - return oneLine(path, 100) - } - if (tool === 'bash' && str(args.command)) return oneLine(str(args.command)!, 100) - if ((tool === 'grep' || tool === 'glob') && str(args.pattern)) { - const where = str(args.path) ?? str(args.glob) - return oneLine(str(args.pattern)! + (where ? ` in ${where}` : ''), 100) - } - if ((tool === 'task' || tool === 'agent') && str(args.description)) { - return oneLine(str(args.description)!, 100) - } - if (tool === 'webfetch' && str(args.url)) return oneLine(str(args.url)!, 100) - if (tool === 'websearch' && str(args.query)) return oneLine(str(args.query)!, 100) - - const keys = Object.keys(args) - if (keys.length === 0) return '' - return oneLine(JSON.stringify(args), 100) -} - -// Human-readable duration from seconds, Claude Code style: "42s" under a -// minute, then "3m 21s". Whole seconds only ("200.7s" reads as noise). -export function formatDuration(seconds: number): string { - const s = Math.round(seconds) - return s >= 60 ? `${Math.floor(s / 60)}m ${s % 60}s` : `${s}s` -} - -// Flatten a message's `content` to a block list (a bare string becomes one -// text block), so assistant and user turns share one iteration path. -function blocksOf(content: unknown): CCContentBlock[] { - if (typeof content === 'string') return [{ type: 'text', text: content }] - if (Array.isArray(content)) return content as CCContentBlock[] - return [] -} - -// Best-effort display text for a tool_result block: its content is a string, a -// list of text blocks, or arbitrary JSON. Whitespace is preserved (results are -// shown as a small indented body), only trimmed at the ends. -function toolResultText(block: CCContentBlock): string { - const c = block.content - if (typeof c === 'string') return c.trim() - if (Array.isArray(c)) { - const parts: string[] = [] - for (const inner of c as CCContentBlock[]) { - if (typeof inner.text === 'string') parts.push(inner.text) - else parts.push(JSON.stringify(inner)) - } - return parts.join('\n').trim() - } - if (c === undefined || c === null) return '' - return JSON.stringify(c) -} - -// Render a whole event into transcript items (usually one, sometimes several -// for a multi-block assistant turn). `keyBase` makes React keys unique across -// the stream; blank when the event carries nothing worth showing. -export function eventToItems(event: CCEvent, keyBase: string): TranscriptItem[] { - const items: TranscriptItem[] = [] - const push = (item: Omit): void => { - items.push({ ...item, key: `${keyBase}:${items.length}` }) - } - const type = event.type - - // A final result event: a dim one-liner capping the turn (cost + duration). - if (type === 'result') { - const bits: string[] = [] - if (typeof event.duration_ms === 'number') bits.push(formatDuration(event.duration_ms / 1000)) - if (typeof event.total_cost_usd === 'number') bits.push(`$${event.total_cost_usd.toFixed(2)}`) - const label = event.is_error ? 'turn ended with an error' : 'turn complete' - push({ - kind: 'summary', - text: bits.length ? `${label} · ${bits.join(' · ')}` : label, - spaceBefore: true, - isError: event.is_error, - }) - return items - } - - // System events (the per-query init with model/cwd, and any informational - // notices) render NOTHING in connect: Claude Code emits an init for every - // user message it processes in stream-json mode, so rendering it printed - // "session started" once per turn — pure noise. The model lives in the - // banner; the raw records stay visible via `agent session records`. - if (type === 'system') { - return items - } - - const blocks = blocksOf(event.message?.content) - - // A user turn is either tool results (grouped under the calls above) or a - // human message injected into the conversation. - if (type === 'user') { - for (const block of blocks) { - if (block.type === 'tool_result') { - const body = toolResultText(block) - push({ - kind: 'tool_result', - gutter: '⎿', - text: body || '(no output)', - spaceBefore: false, - isError: block.is_error, - }) - } else if (typeof block.text === 'string' && block.text.trim()) { - push({ kind: 'user', gutter: '›', text: block.text.trim(), spaceBefore: true }) - } - } - return items - } - - // Assistant turn (and any other event carrying message content): text, - // thinking, and tool calls, in order. - for (const block of blocks) { - if (block.type === 'thinking' && typeof block.thinking === 'string' && block.thinking.trim()) { - push({ kind: 'thinking', gutter: '✻', text: block.thinking.trim(), spaceBefore: true }) - } else if (block.type === 'tool_use') { - const name = block.name ?? 'tool' - const summary = summarizeToolInput(name, block.input) - push({ - kind: 'tool', - gutter: '●', - text: name, - detail: summary ? `(${summary})` : undefined, - spaceBefore: true, - }) - } else if (typeof block.text === 'string' && block.text.trim()) { - push({ kind: 'assistant', text: block.text.trim(), spaceBefore: true }) - } - } - return items -} - -// Render one native session_record into transcript items. A claude_code -// record's payload is a CCEvent (assistant/user/result — expanded by -// eventToItems); a lifecycle record becomes a single dim notice line (the -// spawn/respawn/idle notifications). `keyBase` makes React keys unique. -export function recordToItems(record: SessionRecord, keyBase: string): TranscriptItem[] { - if (record.source === 'lifecycle') { - const text = lifecycleText(record.record_type, record.payload) - return text ? [{ key: keyBase, kind: 'notice', text, spaceBefore: true }] : [] - } - return eventToItems(record.payload as CCEvent, keyBase) -} - -// Which records the connect transcript renders. Claude Code records are the -// conversation; lifecycle records are filtered EXCEPT sandbox_ready — the one -// moment worth a conversation note (the box is up, work can start). The other -// lifecycle rows (starting, paused, closed, resumed) are carried by the live -// activity line / footer / exit notice instead. -export function isConnectVisibleRecord(record: SessionRecord): boolean { - return record.source !== 'lifecycle' || record.record_type === 'sandbox_ready' -} - -// The tool calls that are executing RIGHT NOW, inferred from the committed -// transcript: a `tool` item whose `tool_result` hasn't arrived yet is a tool -// in flight (CC's headless stream emits nothing between the call committing -// and its result landing — this inference is the only live signal). Matching -// is FIFO within the current burst; any non-tool item (prose, thinking, a -// turn's result summary, a user message) means earlier calls resolved, so the -// pending set resets — a stale unmatched call from an errored old turn can -// never read as "running" forever. -export function pendingToolCalls(items: TranscriptItem[]): TranscriptItem[] { - let pending: TranscriptItem[] = [] - for (const item of items) { - if (item.kind === 'tool') pending.push(item) - else if (item.kind === 'tool_result') pending.shift() - else pending = [] - } - return pending -} - -// Collapse each maximal run of consecutive tool activity (tool calls + their -// results) into one dim summary line — the Claude-Code-app treatment ("Ran 8 -// shell commands") — so a burst of shell work reads as one beat of the -// conversation. ctrl+r (the caller's `expanded` state) renders the original -// items instead, restoring the full ● call / ⎿ result blocks. Labels: all-Bash -// runs count shell commands, all-Read runs count files read, mixed runs name -// the tools involved. -export function collapseToolRuns(items: TranscriptItem[]): TranscriptItem[] { - const out: TranscriptItem[] = [] - let group: TranscriptItem[] = [] - const flush = (): void => { - if (group.length === 0) return - const names = [...new Set(group.filter((i) => i.kind === 'tool').map((i) => i.text))] - const n = group.filter((i) => i.kind === 'tool').length || group.length - const plural = n === 1 ? '' : 's' - let label: string - if (names.length === 1 && names[0] === 'Bash') label = `Ran ${n} shell command${plural}` - else if (names.length === 1 && names[0] === 'Read') label = `Read ${n} file${plural}` - else if (names.length === 0) label = `Ran ${n} tool call${plural}` - else { - const shown = names.slice(0, 3).join(', ') + (names.length > 3 ? ', …' : '') - label = `Ran ${n} tool call${plural} (${shown})` - } - out.push({ key: `grp:${group[0].key}`, kind: 'notice', text: label, spaceBefore: true }) - group = [] - } - for (const item of items) { - if (item.kind === 'tool' || item.kind === 'tool_result') group.push(item) - else { - flush() - out.push(item) - } - } - flush() - return out -} - -// Clamp a multi-line body to `maxLines`, appending a dim "+N lines" marker when -// truncated — used for tool-result bodies so a huge file read stays compact. -export function clampLines(text: string, maxLines: number): { body: string; more: number } { - const lines = text.split('\n') - if (lines.length <= maxLines) return { body: text, more: 0 } - return { body: lines.slice(0, maxLines).join('\n'), more: lines.length - maxLines } -} - -// The cumulative session cost (USD) a Claude Code `result` event carries: each -// result caps a turn and reports the running total for the whole session (not -// the turn alone). null for non-result events or a result without a cost. -export function resultCostUsd(event: CCEvent): number | null { - if (event.type !== 'result') return null - return typeof event.total_cost_usd === 'number' ? event.total_cost_usd : null -} - -// Fold a chronological event list into the spend the footer shows: `total` is -// the latest result's cumulative cost; `lastStep` is the delta from the result -// before it — i.e. the cost of the most recent completed turn. Both null until -// the first result lands; for the first result, lastStep equals total. -export function foldCosts(events: CCEvent[]): { - total: number | null - lastStep: number | null -} { - let prev: number | null = null - let total: number | null = null - for (const event of events) { - const cost = resultCostUsd(event) - if (cost == null) continue - prev = total - total = cost - } - const lastStep = total != null && prev != null ? Math.max(0, total - prev) : total - return { total, lastStep } -} - -// The concrete label for the ✻ activity line during INFRASTRUCTURE phases — -// the sandbox spawning/waking — where nothing else on screen moves. It -// re-renders in place as the status changes; statuses never append transcript -// lines (that was the old, noisy model). null for `working` (the UI shows a -// whimsical Claude-Code-style gerund there instead) and for calm states -// (waiting, sleeping, terminal), where the line hides entirely. -export function statusActivityText(status: string): string | null { - switch (status) { - case 'scheduled': - return 'Waiting for a worker' - case 'starting': - return 'Starting sandbox' - case 'retrying': - return 'Retrying after a transient error' - default: - return null - } -} diff --git a/src/lib/steps.ts b/src/lib/steps.ts index 6d6b656..28679f3 100644 --- a/src/lib/steps.ts +++ b/src/lib/steps.ts @@ -1,83 +1,15 @@ +import { lifecycleText, oneLine } from '@ellipsis-dev/sdk/store' import { formatTs } from './output' import type { SessionRecord } from './types' +// Re-exported for the record-view callers below and their historical +// importers; the implementations live in the SDK's store layer now. +export { lifecycleText, oneLine, setupOutputHook, setupOutputLine } from '@ellipsis-dev/sdk/store' + // Record-rendering helpers shared by `session records` and `session connect` // (moved out of commands/session.tsx so connect.ts can use them without an // import cycle; session.tsx re-exports them for compatibility). -// Human one-liner for a source==='lifecycle' record — the spawn/respawn/idle -// notifications the transcript itself doesn't carry. null for a type we don't -// surface (falls back to the raw record_type). -export function lifecycleText( - recordType: string, - payload: Record, -): string | null { - switch (recordType) { - case 'sandbox_starting': - return 'Starting sandbox…' - case 'sandbox_setup_output': { - // One chunk of setup-script output ({hook, chunk, lines}): show the - // script's latest line so the record view reads as install progress. - const line = setupOutputLine(payload) - return line ? `${setupOutputHook(payload)} · ${line}` : null - } - case 'sandbox_ready': { - const repos = Array.isArray(payload.repositories) - ? (payload.repositories as unknown[]).filter((r): r is string => typeof r === 'string') - : [] - const parts = ['Sandbox ready'] - if (repos.length) parts.push(repos.join(', ')) - const tier = cacheTierLabel(payload.cache_tier) - if (tier) parts.push(tier) - return parts.join(' · ') - } - case 'session_resumed': - return 'Resumed the conversation' - case 'session_paused': - return 'Sleeping — your next message wakes it' - case 'session_closed': - return 'Conversation closed' - case 'session_cancelled': { - const reason = typeof payload.reason === 'string' ? payload.reason : null - return reason ? `Session cancelled · ${reason}` : 'Session cancelled' - } - default: - return null - } -} - -// The customer script a sandbox_setup_output chunk came from (image.setup / -// post_start / post_clone). -export function setupOutputHook(payload: Record): string { - return typeof payload.hook === 'string' ? payload.hook : 'setup' -} - -// The last non-empty output line of a sandbox_setup_output chunk — what the -// live "Starting sandbox" sub-line and the record view both show. -export function setupOutputLine(payload: Record): string | null { - const lines = Array.isArray(payload.lines) - ? (payload.lines as unknown[]).filter( - (l): l is string => typeof l === 'string' && l.trim().length > 0, - ) - : [] - return lines.length ? lines[lines.length - 1].trim() : null -} - -// Customer-facing wording for sandbox_ready's cache_tier, explaining why the -// start was fast or slow. -function cacheTierLabel(tier: unknown): string | null { - switch (tier) { - case 'exact': - return 'cached image' - case 'incremental': - return 'incremental build' - case 'full': - return 'full build' - default: - return null - } -} - // A content block of a Claude Code stream event, typed loosely: the CLI only // extracts display text and names, never interprets the payload. interface StepContentBlock { @@ -136,9 +68,3 @@ function contentText(content: unknown): string { } return parts.join(' ') } - -// Collapse whitespace/newlines to one displayable line, truncated to `max`. -export function oneLine(text: string, max: number): string { - const collapsed = text.replace(/\s+/g, ' ').trim() - return collapsed.length <= max ? collapsed : `${collapsed.slice(0, max - 3)}...` -} diff --git a/src/lib/stream.ts b/src/lib/stream.ts new file mode 100644 index 0000000..2044913 --- /dev/null +++ b/src/lib/stream.ts @@ -0,0 +1,51 @@ +// The CLI's transport adapter for @ellipsis-dev/sdk/stream: the SDK owns the +// session-stream machinery (frames, reconnect/backoff, after_seq resume, +// heartbeat liveness); this module owns what is CLI-specific — resolving the +// WebSocket base URL from the environment/API base, building the bearer-door +// URL, and adapting the `ws` package to the SDK's injected socket surface. + +import WebSocket from 'ws' +import type { OpenSocket, StreamSocket } from '@ellipsis-dev/sdk/stream' +import { resolveApiBase } from './config' +import { DEFAULT_WS_BASE, USER_AGENT } from './constants' + +// Resolve the WebSocket base URL: explicit env wins, else derive from the +// resolved API base (so http://localhost -> ws://localhost), else the default. +export function resolveWsBase(apiBase?: string): string { + const explicit = process.env.ELLIPSIS_WS_BASE + if (explicit) return explicit.replace(/\/+$/, '') + const base = (apiBase ?? resolveApiBase()).replace(/\/+$/, '') + if (base.startsWith('https://')) return 'wss://' + base.slice('https://'.length) + if (base.startsWith('http://')) return 'ws://' + base.slice('http://'.length) + if (base.startsWith('ws://') || base.startsWith('wss://')) return base + return DEFAULT_WS_BASE +} + +// The bearer door's stream URL: /v1/sessions/{id}/stream plus the SDK's +// handshake query (`protocol=2`, `after_seq` when resuming). +export function buildStreamUrl(wsBase: string, sessionId: string, query: string): string { + return `${wsBase}/v1/sessions/${encodeURIComponent(sessionId)}/stream?${query}` +} + +// An OpenSocket over the `ws` package with bearer auth — what every CLI +// stream consumer injects into the SDK's streamSession. +export function makeOpenSocket(token: string, wsBase?: string): OpenSocket { + const base = wsBase ?? resolveWsBase() + return ({ sessionId, query }): StreamSocket => { + const ws = new WebSocket(buildStreamUrl(base, sessionId, query), { + headers: { authorization: `Bearer ${token}`, 'user-agent': USER_AGENT }, + }) + return { + onOpen: (cb) => ws.on('open', cb), + onMessage: (cb) => ws.on('message', (raw: WebSocket.RawData) => cb(raw.toString())), + onClose: (cb) => ws.on('close', (code: number) => cb(code)), + // The `ws` package emits a real Error here; surface its message rather + // than String(err) (which produced "[object ErrorEvent]" once). + onError: (cb) => + ws.on('error', (err: unknown) => + cb(err instanceof Error ? err : new Error(String(err))), + ), + close: () => ws.close(), + } + } +} diff --git a/src/lib/types.ts b/src/lib/types.ts index d4536ac..70d97c9 100644 --- a/src/lib/types.ts +++ b/src/lib/types.ts @@ -1,8 +1,38 @@ -// TypeScript mirror of the backend `/v1` request/response models. These are -// hand-rolled for now and will be replaced by the generated @ellipsis/sdk -// package once it exists. Nested config/input/output payloads are typed loosely -// (the CLI only displays summary fields); the rest mirror the Pydantic models -// in ellipsis/src/public_api/routers/v1/v1_router.py and the shared services. +// TypeScript types for the backend `/v1` request/response models. +// +// The session-stream surface (records, inbox messages, turns, the enriched +// session wire shape, and their request/response DTOs) comes from +// @ellipsis-dev/sdk — generated from the server's schema, never hand-written — +// re-exported below under the CLI's historical names. Everything else (the +// endpoints outside the SDK's REST surface: session list/start, configs, +// sandboxes, integrations, …) remains a hand-rolled mirror of the Pydantic +// models in ellipsis's v1_router until the SDK's OpenAPI surface widens. +// Nested config/input/output payloads are typed loosely (the CLI only +// displays summary fields). + +import type { + AgentSessionSource, + AgentSessionStatus, + SessionMessageWire, + SessionRecordWire, + SessionState, + SessionSurface, +} from '@ellipsis-dev/sdk' + +export type { + AgentSessionSource, + AgentSessionStatus, + AgentSessionWire, + ListSessionRecordsResponse, + ListSessionTurnsResponse, + SendSessionMessageRequest, + SessionState, + SessionSurface, +} from '@ellipsis-dev/sdk' + +// The CLI's historical names for the SDK's wire models. +export type SessionRecord = SessionRecordWire +export type SessionMessage = SessionMessageWire // ------------------------------- identity ------------------------------- @@ -76,26 +106,6 @@ export interface UsageDashboard { // ----------------------------- agent sessions ---------------------------- -export type AgentSessionSource = - | 'react' - | 'manual' - | 'api' - | 'cli' - | 'mention' - | 'cron' - // A Claude Code session ingested from a developer laptop via `agent session sync`. - | 'laptop' - -export type AgentSessionStatus = - | 'scheduled' - | 'creating_sandbox' - | 'running' - | 'retrying' - | 'completed' - | 'error' - | 'cancelled' - | 'stopped' - // Loosely typed: the CLI reads a handful of summary fields and otherwise treats // the session as opaque JSON. See AgentSession in the backend for the full shape. export interface AgentSession { @@ -156,15 +166,6 @@ export type AgentConfig = Record // --------------------------- request / response ------------------------- -// Body of POST /v1/sessions/{id}/messages (`agent session connect`): a human -// message appended to a durable session's inbox. -export interface SendSessionMessageRequest { - message: string - // Retry-safety key, unique per (session, key): a retried POST returns the - // original message instead of double-queueing a turn (protocol v2 §4.2). - idempotency_key?: string | null -} - // Laptop -> cloud handoff params: start a fresh session on the built-in // handoff config, chained to the handed-off session (parent_kind=handoff). // Mutually exclusive with config_id / config / template_id. @@ -337,90 +338,6 @@ export interface ListAgentSessionsQuery { // The render switch on a session_record (session_records.source). The CLI // renders claude_code records natively and lifecycle records as system lines. -export type RecordSource = 'claude_code' | 'lifecycle' - -// session_records.record_type for source==='lifecycle' rows: the -// spawn/respawn/idle notifications the transcript itself does not carry. -export type LifecycleRecordType = - | 'sandbox_starting' - | 'sandbox_ready' - | 'session_resumed' - | 'session_paused' - | 'session_closed' - | 'session_cancelled' - -// One native session_record from GET /v1/sessions/{id}/records. `payload` is the -// harness's verbatim line — for claude_code, a Claude Code stream event -// (assistant turn, tool call/result, result); for lifecycle, a small blob. The -// CLI switches on `source` and only extracts display text. `feed_seq` is the -// shared per-session order (transcript + lifecycle merged); `stream_seq` is the -// per-execution native index (NEGATIVE for lifecycle records). -export interface SessionRecord { - id: string - agent_session_id: string - created_at: string - feed_seq: number - stream_seq: number - // Free TEXT server-side; unknown sources must be ignored, not crashed on. - source: RecordSource | string - record_type: string - record_format: string - // The turn this record belongs to (stream protocol v2 §3.3). - agent_turn_id?: string | null - // Correlates a user-echo record back to the inbox message it echoes, so - // queued chips retire by id (§4.2). - session_message_id?: string | null - payload: Record - [key: string]: unknown -} - -export interface ListSessionRecordsResponse { - records: SessionRecord[] - // The OPEN inbox slice (pending rows only) — the server-side queued signal. - messages?: SessionMessage[] - // Whether records past this page remain (only meaningful with ?limit=). - has_more?: boolean - // The retention head: lowest stored feed_seq, or null when no records. - earliest_feed_seq?: number | null -} - -// One exchange within a keyed session (a single Claude Code turn), from -// GET /v1/sessions/{id}/turns. -export interface SessionTurn { - id: string - agent_session_id: string - turn_index: number - status: string - created_at: string - completed_at: string | null - [key: string]: unknown -} - -// One entry in a keyed session's inbox: a message posted into the conversation, -// `pending` until a turn consumes it (the server-side "queued" signal), -// `delivered` once it reached the agent. `body` is the raw text as sent. -export interface SessionMessage { - id: string - agent_session_id: string - body: string - status: 'pending' | 'delivered' - feed_seq: number | null - // The sender's display name for single-human-utterance messages; null for - // event-rendered messages. - author?: string | null - created_at: string - delivered_at: string | null - delivered_turn_id: string | null - [key: string]: unknown -} - -// GET /v1/sessions/{id}/turns: the conversation structure of a keyed session — -// empty lists for single-shot sessions. -export interface ListSessionTurnsResponse { - turns: SessionTurn[] - messages: SessionMessage[] -} - // One process's raw transcript from GET /v1/sessions/{id}/transcripts: the // pointer metadata plus a short-lived presigned S3 GET for the .jsonl.gz // object itself (expires after expires_in seconds — fetch it immediately). diff --git a/src/lib/ws.ts b/src/lib/ws.ts deleted file mode 100644 index 2897bb1..0000000 --- a/src/lib/ws.ts +++ /dev/null @@ -1,361 +0,0 @@ -import WebSocket from 'ws' -import { resolveApiBase } from './config' -import { DEFAULT_WS_BASE, USER_AGENT } from './constants' -import type { AgentSession, SessionMessage, SessionRecord } from './types' - -// The session stream protocol v2 (server -> client), one JSON object per -// message. Mirrors the frame DTOs in the backend's -// public_api/realtime/session_stream.py (documents/eng/SESSION_STREAM_PROTOCOL.md): -// -// snapshot: { type, protocol, earliest_feed_seq, session, messages } -// records_append: { type, records } — raw session_records, by feed_seq -// messages: { type, messages } — the OPEN inbox slice (LWW) -// session: { type, session } — the lean session, resent on change -// delta: { type, agent_turn_id, kind, text?, output_tokens? } -// heartbeat: { type, ts } -// done: { type } — conversation over; close 1000 follows -// error: { type, message } — curated copy; close 1011 follows -// -// Delivery classes: records_append is cursored append-only — the ONE resume -// cursor (`?after_seq=` = the highest feed_seq seen in records_append frames; -// messages/session frames NEVER advance it). session/messages are -// last-writer-wins snapshots; delta/heartbeat are fire-and-forget; done/error -// are terminal. Unknown frame types (and unknown source/record_type/kind -// values) MUST be ignored — additive server changes are not a protocol break. -export const SESSION_STREAM_PROTOCOL_VERSION = 2 - -export type StreamFrame = - | { - type: 'snapshot' - protocol: number - // The retention head: the lowest feed_seq still stored, or null when the - // session has no stored records. Replaying from before it means history - // was truncated under the customer's log retention. - earliest_feed_seq: number | null - session: AgentSession - messages: SessionMessage[] - } - | { type: 'records_append'; records: SessionRecord[] } - | { type: 'messages'; messages: SessionMessage[] } - | { type: 'session'; session: AgentSession } - | { - type: 'delta' - agent_turn_id?: string | null - // "text" today; "thinking" reserved — ignore unknown kinds. - kind?: string - text?: string | null - output_tokens?: number | null - } - | { type: 'heartbeat'; ts?: string } - | { type: 'done' } - | { type: 'error'; message?: string } - // Forward compatibility: an unrecognized frame type parses, is handed to - // onFrame, and must be ignored by renderers. - | { type: string; [key: string]: unknown } - -// The display word for a session: the server-derived surface status -// (working/waiting/sleeping/starting/closed/…) when present, else the raw -// per-execution status. Shared by every frame consumer so the stream and the -// REST poll can never disagree about the word. -export function sessionStatusWord(session: AgentSession): string { - return session.surface?.status ?? session.status -} - -// How streamSession() finished. `done`/`error` are normal terminal outcomes -// (`status` is the last session frame's derived word); `aborted` means the -// caller cancelled via the AbortSignal. -export type StreamOutcome = - | { type: 'done'; status: string; exitStatus?: string | null } - | { type: 'error'; message: string } - | { type: 'aborted' } - -// Thrown when streaming isn't usable (no endpoint, unsupported protocol, or -// reconnects exhausted) — the caller should fall back to REST polling. -export class StreamUnavailableError extends Error { - constructor(message: string) { - super(message) - this.name = 'StreamUnavailableError' - } -} - -// Thrown when the server rejects the credential for this session (close 1008). -// Polling would fail the same way, so this is not a fallback case. -export class StreamAuthError extends Error { - constructor(message: string) { - super(message) - this.name = 'StreamAuthError' - } -} - -export interface StreamOptions { - token: string - sessionId: string - onFrame: (frame: StreamFrame) => void - wsBase?: string - afterSeq?: number - signal?: AbortSignal - maxReconnects?: number - // Socket factory, injectable for tests. Defaults to a real `ws` WebSocket. - connect?: SocketFactory -} - -// Minimal socket surface streamSession() depends on, so tests can drive the -// reconnect/resume logic without a real server. -export interface StreamSocket { - onOpen(cb: () => void): void - onMessage(cb: (data: string) => void): void - onClose(cb: (code: number) => void): void - onError(cb: (err: Error) => void): void - close(): void -} -export type SocketFactory = (url: string, token: string) => StreamSocket - -// Server heartbeat cadence is 20s; if we hear nothing for ~2x that the socket -// is presumed dead and we reconnect (the protocol's own liveness rule). -const HEARTBEAT_TIMEOUT_MS = 45_000 -const DEFAULT_MAX_RECONNECTS = 5 - -const defaultFactory: SocketFactory = (url, token) => { - const ws = new WebSocket(url, { - headers: { authorization: `Bearer ${token}`, 'user-agent': USER_AGENT }, - }) - return { - onOpen: (cb) => ws.on('open', cb), - onMessage: (cb) => ws.on('message', (raw: WebSocket.RawData) => cb(raw.toString())), - onClose: (cb) => ws.on('close', (code: number) => cb(code)), - // The `ws` package emits a real Error here; surface its message rather than - // String(err) (which produced "[object ErrorEvent]" against the old stub). - onError: (cb) => - ws.on('error', (err: unknown) => - cb(err instanceof Error ? err : new Error(String(err))), - ), - close: () => ws.close(), - } -} - -// ----------------------------- pure helpers -------------------------------- - -export type CloseKind = 'normal' | 'auth' | 'unsupported' | 'retry' - -export function classifyCloseCode(code: number): CloseKind { - switch (code) { - case 1000: - return 'normal' - case 1008: - return 'auth' - case 1002: // unsupported protocol version requested - case 1003: // no ?protocol= param (shouldn't happen — we always send it) - return 'unsupported' - default: - // 1011 server error, 1013 over capacity, 1006 abnormal, etc. - return 'retry' - } -} - -// Exponential backoff, capped. Deterministic (no jitter) so it's easy to test. -export function nextReconnectDelayMs(attempt: number): number { - const base = 500 - const max = 8_000 - return Math.min(max, base * 2 ** Math.max(0, attempt - 1)) -} - -export interface ReconnectDecision { - action: 'reconnect' | 'fallback' | 'fail-auth' - delayMs?: number -} - -// Decide what to do after a connection ends without a terminal frame. Pure, so -// the reconnect/resume/fallback policy is unit-tested directly. -export function decideReconnect(params: { - closeKind?: CloseKind - errored?: boolean - everReceivedFrame: boolean - attempt: number - maxReconnects: number -}): ReconnectDecision { - const { closeKind, everReceivedFrame, attempt, maxReconnects } = params - if (closeKind === 'auth') return { action: 'fail-auth' } - if (closeKind === 'unsupported') return { action: 'fallback' } - // Retryable: a transport error, an abnormal close, or a normal close that - // arrived before the `done` frame. Be persistent once we've seen the server - // actually stream (it clearly supports it); bail out fast otherwise so a - // backend without the endpoint falls back to polling promptly. - const cap = everReceivedFrame ? maxReconnects : Math.min(2, maxReconnects) - if (attempt >= cap) return { action: 'fallback' } - return { action: 'reconnect', delayMs: nextReconnectDelayMs(attempt) } -} - -// Resolve the WebSocket base URL: explicit env wins, else derive from the -// resolved API base (so http://localhost -> ws://localhost), else the default. -export function resolveWsBase(apiBase?: string): string { - const explicit = process.env.ELLIPSIS_WS_BASE - if (explicit) return explicit.replace(/\/+$/, '') - const base = (apiBase ?? resolveApiBase()).replace(/\/+$/, '') - if (base.startsWith('https://')) return 'wss://' + base.slice('https://'.length) - if (base.startsWith('http://')) return 'ws://' + base.slice('http://'.length) - if (base.startsWith('ws://') || base.startsWith('wss://')) return base - return DEFAULT_WS_BASE -} - -export function buildStreamUrl(wsBase: string, sessionId: string, afterSeq: number): string { - // ?protocol= is REQUIRED by the v2 handshake: a server that doesn't see it - // closes 1003 (how pre-v2 binaries degrade to polling); an unknown version - // closes 1002 with the supported list in the reason. - const url = `${wsBase}/v1/sessions/${encodeURIComponent(sessionId)}/stream?protocol=${SESSION_STREAM_PROTOCOL_VERSION}` - return afterSeq > 0 ? `${url}&after_seq=${afterSeq}` : url -} - -// ------------------------------ connection --------------------------------- - -type ConnResult = - | { kind: 'done' } - | { kind: 'frameError'; message: string } - | { kind: 'closed'; code: number } - | { kind: 'error'; err: Error } - | { kind: 'aborted' } - -// One WebSocket connection. Resolves when the stream reaches a terminal frame, -// the socket closes/errors, the heartbeat lapses, or the signal aborts. -function connectOnce( - url: string, - token: string, - factory: SocketFactory, - emit: (frame: StreamFrame) => void, - signal?: AbortSignal, -): Promise { - return new Promise((resolve) => { - const sock = factory(url, token) - let settled = false - let heartbeat: ReturnType | undefined - - const finish = (result: ConnResult) => { - if (settled) return - settled = true - if (heartbeat) clearTimeout(heartbeat) - if (signal) signal.removeEventListener('abort', onAbort) - sock.close() - resolve(result) - } - const onAbort = () => finish({ kind: 'aborted' }) - const bumpHeartbeat = () => { - if (heartbeat) clearTimeout(heartbeat) - heartbeat = setTimeout( - () => finish({ kind: 'error', err: new Error('heartbeat timeout') }), - HEARTBEAT_TIMEOUT_MS, - ) - } - - if (signal) { - if (signal.aborted) { - finish({ kind: 'aborted' }) - return - } - signal.addEventListener('abort', onAbort) - } - - sock.onOpen(() => bumpHeartbeat()) - sock.onMessage((data) => { - bumpHeartbeat() - let frame: StreamFrame - try { - frame = JSON.parse(data) as StreamFrame - } catch { - return // ignore non-JSON keepalives / garbage - } - emit(frame) - if (frame.type === 'done') { - finish({ kind: 'done' }) - } else if (frame.type === 'error') { - finish({ - kind: 'frameError', - message: (frame as { message?: string }).message ?? 'stream error', - }) - } - }) - sock.onClose((code) => finish({ kind: 'closed', code })) - sock.onError((err) => finish({ kind: 'error', err })) - }) -} - -function sleep(ms: number, signal?: AbortSignal): Promise { - return new Promise((resolve) => { - const timer = setTimeout(resolve, ms) - signal?.addEventListener( - 'abort', - () => { - clearTimeout(timer) - resolve() - }, - { once: true }, - ) - }) -} - -// ------------------------------ public API --------------------------------- - -// Stream an agent session to completion, reconnecting with backoff and -// resuming from the last records_append feed_seq so a dropped socket loses no -// records (§3.4: ONLY records advance the cursor — session/messages snapshots -// are re-sent fresh on reconnect). Calls `onFrame` for every frame received. -// Resolves with the terminal outcome — `done`'s status/exitStatus come from -// the last session frame, which the server guarantees carries the end state -// before `done`. Throws StreamUnavailableError (caller should poll instead) / -// StreamAuthError. -export async function streamSession(opts: StreamOptions): Promise { - const wsBase = opts.wsBase ?? resolveWsBase() - const factory = opts.connect ?? defaultFactory - const maxReconnects = opts.maxReconnects ?? DEFAULT_MAX_RECONNECTS - - let afterSeq = opts.afterSeq ?? 0 - let everReceivedFrame = false - let attempt = 0 - let lastStatusWord = '' - let lastExitStatus: string | null | undefined - - const emit = (frame: StreamFrame) => { - everReceivedFrame = true - if (frame.type === 'records_append') { - const records = (frame as { records: SessionRecord[] }).records - for (const record of records) { - if (typeof record.feed_seq === 'number') { - afterSeq = Math.max(afterSeq, record.feed_seq) - } - } - } else if (frame.type === 'snapshot' || frame.type === 'session') { - const session = (frame as { session: AgentSession }).session - lastStatusWord = sessionStatusWord(session) - lastExitStatus = (session.exit_status as string | null | undefined) ?? null - } - opts.onFrame(frame) - } - - for (;;) { - if (opts.signal?.aborted) return { type: 'aborted' } - const url = buildStreamUrl(wsBase, opts.sessionId, afterSeq) - const res = await connectOnce(url, opts.token, factory, emit, opts.signal) - - if (res.kind === 'done') { - return { type: 'done', status: lastStatusWord, exitStatus: lastExitStatus } - } - if (res.kind === 'frameError') return { type: 'error', message: res.message } - if (res.kind === 'aborted') return { type: 'aborted' } - - attempt++ - const decision = decideReconnect({ - closeKind: res.kind === 'closed' ? classifyCloseCode(res.code) : undefined, - errored: res.kind === 'error', - everReceivedFrame, - attempt, - maxReconnects, - }) - if (decision.action === 'fail-auth') { - throw new StreamAuthError('not authorized to stream this session') - } - if (decision.action === 'fallback') { - const why = - res.kind === 'error' ? res.err.message : `stream closed (code ${res.code})` - throw new StreamUnavailableError(why) - } - await sleep(decision.delayMs ?? 0, opts.signal) - } -} diff --git a/src/ui/ConnectApp.tsx b/src/ui/ConnectApp.tsx index 1558623..afaf904 100644 --- a/src/ui/ConnectApp.tsx +++ b/src/ui/ConnectApp.tsx @@ -1,13 +1,18 @@ -import React, { useCallback, useEffect, useMemo, useRef, useState } from 'react' +import React, { + useCallback, + useEffect, + useMemo, + useRef, + useState, + useSyncExternalStore, +} from 'react' import { Box, Text, useApp, useInput, useStdin, useStdout } from 'ink' -import { ApiClient, ApiError } from '../lib/api' import { - sessionStatusWord, streamSession, + sessionStatusWord, StreamUnavailableError, - type StreamFrame, -} from '../lib/ws' -import type { AgentSession, SessionMessage, SessionRecord } from '../lib/types' + type OpenSocket, +} from '@ellipsis-dev/sdk/stream' import { clampLines, collapseToolRuns, @@ -16,61 +21,58 @@ import { isConnectVisibleRecord, pendingToolCalls, recordToItems, - resultCostUsd, + setupOutputHook, + setupOutputLine, statusActivityText, + oneLine, type CCEvent, type ItemKind, + type SessionTranscriptStore, type TranscriptItem, -} from '../lib/events' +} from '@ellipsis-dev/sdk/store' +import { ApiClient, ApiError } from '../lib/api' import { hyperlink } from '../lib/urls' import { usdNumberFromMillicents } from '../lib/output' -import { oneLine, setupOutputHook, setupOutputLine } from '../lib/steps' import { VERSION } from '../lib/constants' // The interactive `agent session connect` UI, modelled on Claude Code: a // committed transcript that groups tool calls with their results and spaces // messages apart — live activity (✻ Running/Generating) rendered on the // transcript block it describes — above a footer with a composer that echoes -// what you send. Rendering shape lives in lib/events.ts -// (pure); this component owns the data flow, the composer, and the colours. +// what you send. Rendering shape lives in @ellipsis-dev/sdk/store (pure); this +// component owns the data flow, the composer, and the colours. // -// Data flow: the committed transcript comes from the structured records API -// (GET /v1/sessions/{id}/records, whose payload is the full native event) — -// grouped into tool calls / results — with the socket as a low-latency -// "something changed" wake plus status source, backed by a slow poll. On top of -// that, the socket also carries EPHEMERAL `delta` frames (partial assistant text -// + a running output-token count) that render as a live, in-progress line and a -// footer token counter — the token-by-token feel of local Claude Code — until -// the committed assistant step lands and supersedes it. +// Data flow: ONE SessionTranscriptStore (pre-seeded by the caller with the +// stored records + session, so the first paint is instant) is fed by the +// SDK's streamSession — records arrive PUSHED as records_append frames, the +// session/messages frames carry status + the open inbox, and ephemeral +// `delta` frames overlay the in-progress response token-by-token. Everything +// on screen derives from the store snapshot; there is no REST refresh loop. +// If the stream is unavailable (old backend, blocked socket), a REST poll +// feeds the SAME store through synthetic frames, so the UI is identical +// either way. export interface ConnectAppProps { api: ApiClient - token: string sessionId: string - wsBase: string + // The one transcript store, pre-seeded with the fetched records + session. + store: SessionTranscriptStore + // The bearer-door socket factory (lib/stream.ts makeOpenSocket). + openSocket: OpenSocket // Keyed, open sessions accept messages (show the composer); single-shot / // closed / --no-input sessions follow read-only and exit when the stream ends. canSend: boolean - initialItems: TranscriptItem[] - // The highest feed_seq already rendered into initialItems, so live refreshes - // only append records newer than what's on screen (feed_seq is the shared - // per-session order across transcript + lifecycle records). - initialMaxFeedSeq: number - initialStatus: string + // Records at or below this feed_seq are not RENDERED (--no-records skips + // replaying history on screen without re-streaming it). 0 renders everything. + minRenderFeedSeq: number // The clickable dashboard link for this session (app.ellipsis.dev/…/sessions/{id}), // shown in the footer status line. sessionUrl: string - // Spend seeded from the stored steps: cumulative total + the last turn's cost. - initialCost: { total: number | null; lastStep: number | null } - // The server-reported session total (USD) at connect time, from the session's - // cost columns — the billing authority. Live updates arrive as - // `cost_millicents` on status/done frames; null against older backends. - initialServerCostUsd?: number | null // A one-line caveat shown as the app's opening notice (e.g. "watch-only: // this conversation is closed"). null for the normal connect. initialNotice?: string | null // The session's model (backend tokens_model, fixed at creation), shown in - // the banner under the dashboard link. + // the footer meta line. model?: string | null // Written (not read) by the app: set true when the app exits because the // conversation closed, so the caller skips the "detached — still running" @@ -88,14 +90,13 @@ function isWorkingStatus(status: string): boolean { // One local send awaiting server acknowledgement: messageId is null while the // POST is in flight, then the created SessionMessage's id (protocol v2 §4.2) — -// the chip retires the moment the server acknowledges that id (a messages -// frame, an inbox fetch, or the transcript user-echo record's -// session_message_id back-reference), and the server's own pending row takes -// over as its representation. +// the chip retires the moment the store acknowledges that id (a messages +// frame, or the transcript user-echo record's session_message_id +// back-reference), and the server's own pending row takes over. type QueuedSend = { text: string; messageId: string | null } export function ConnectApp(props: ConnectAppProps): React.ReactElement { - const { api, token, sessionId, wsBase, canSend } = props + const { api, sessionId, store, openSocket, canSend } = props const { exit } = useApp() const { isRawModeSupported } = useStdin() const { stdout } = useStdout() @@ -114,352 +115,170 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { } }, [stdout]) - const [items, setItems] = useState(props.initialItems) - const [status, setStatus] = useState(props.initialStatus) - const [working, setWorking] = useState(isWorkingStatus(props.initialStatus)) + // The store snapshot is the single source of truth for everything streamed. + const snapshot = useSyncExternalStore(store.subscribe, store.getSnapshot, store.getSnapshot) + const statusWord = snapshot.session ? sessionStatusWord(snapshot.session) : 'starting' + // Bridge the gap between a send and the server's status flip: treat the + // session as working until the next status transition lands. + const [sendPending, setSendPending] = useState(false) + useEffect(() => { + setSendPending(false) + }, [statusWord]) + const working = isWorkingStatus(statusWord) || sendPending + const [elapsed, setElapsed] = useState(0) - // The setup script's latest output line while the sandbox spawns (streamed - // sandbox_setup_output lifecycle chunks) — the sub-line under "Starting - // sandbox…" that says WHAT the box is doing. Cleared by sandbox_ready. - const [setupLine, setSetupLine] = useState(null) const [notice, setNotice] = useState(props.initialNotice ?? null) const [input, setInput] = useState('') // ctrl+r toggles full vs. collapsed tool output across the whole transcript. const [expanded, setExpanded] = useState(false) - // Messages you've sent that the agent hasn't picked up yet — shown as a - // queued region below the composer, exactly like Claude Code. A send renders - // as a LOCAL chip only until a records/turns refresh that started after its - // POST resolved confirms the server has the message (postSeq, see - // refreshSteps); from then on the server's own pending rows (serverQueued) - // are the queued truth. That keeps the region honest even when the agent - // consumes messages in ways local text-matching can't follow (the server once - // coalesced two queued sends into one "a\nb" user turn, stranding both chips - // forever). Chips also leave via the transcript: when the backend relays a - // send as a user turn, or when the turn idles (whichever first), so a send is - // never lost. + // Messages you've sent that the server hasn't acknowledged yet — shown as a + // queued region below the composer, exactly like Claude Code. From the first + // acknowledgement on (a messages frame or the user-echo record carrying the + // id), the server's own pending rows (serverQueued) are the queued truth. const [queued, setQueued] = useState([]) - // Bodies of the server's PENDING inbox messages — the durable queued signal. - const [serverQueued, setServerQueued] = useState([]) - // Ephemeral live-streaming overlay for the CURRENT assistant response, driven - // by `delta` frames: the prose so far and the running output-token count. Both - // clear when the committed assistant step lands (it supersedes the overlay) or - // the turn settles. - const [liveText, setLiveText] = useState('') - const [liveTokens, setLiveTokens] = useState(null) - // Running spend for the footer: `total` is the cumulative session cost from - // the latest Claude Code result; `lastStep` is the cost of the most recent - // turn (the delta between the last two results). - const [cost, setCost] = useState(props.initialCost) - // The server's authoritative running total (USD), from `cost_millicents` on - // status/done frames (and the status poll as backstop). Preferred over the - // CC-derived total when present; kept monotonic so the footer never dips. - const [serverCostUsd, setServerCostUsd] = useState( - props.initialServerCostUsd ?? null, - ) - const applyServerCostUsd = useCallback((usd: number): void => { - setServerCostUsd((prev) => (prev == null ? usd : Math.max(prev, usd))) - }, []) - - // Live-flow state that must survive re-renders without triggering them. - // The stream resume cursor (protocol v2 §3.4: the highest feed_seq seen in - // records_append frames) — seeded from the REST-rendered transcript so the - // first attach doesn't replay history the screen already shows. - const afterSeq = useRef(props.initialMaxFeedSeq) + + // Whether the sandbox ever reached a connectable state, so a terminal status + // *before* that (a preflight/budget gate) is reported as a failure, not idle. + const everRunning = useRef(isWorkingStatus(statusWord) && statusWord !== 'scheduled') + useEffect(() => { + if (['working', 'waiting'].includes(statusWord)) everRunning.current = true + }, [statusWord]) + const streaming = useRef(false) + const polling = useRef(false) const abort = useRef(new AbortController()) - const lastStatus = useRef(props.initialStatus) - const keyCounter = useRef(0) - // The highest feed_seq committed to the transcript, and the refresh - // in-flight / re-run guards so overlapping wakes don't double-append. - const maxFeed = useRef(props.initialMaxFeedSeq) - // The cumulative cost last committed, so a new result's delta = the turn cost. - const costTotal = useRef(props.initialCost.total) - // Whether the sandbox ever reached `running`, so a terminal status *before* - // that (a preflight/budget gate) is reported as a failure, not idle. - const everRunning = useRef(props.initialStatus === 'running') - const refreshing = useRef(false) - const pendingRefresh = useRef(false) - // Guard so the closed-conversation teardown (final refresh + exit) runs once - // no matter which signal lands first (status frame, poll, stream end). + // Guard so the closed-conversation teardown runs once no matter which signal + // lands first (session frame, poll, stream outcome). const closingDown = useRef(false) - const refreshTimer = useRef | null>(null) - // A mirror of `queued` for async callbacks, and the texts already committed - // locally so the backend's later echo of the same send is dropped. - const queuedRef = useRef([]) - const flushed = useRef([]) - useEffect(() => { - queuedRef.current = queued - }, [queued]) - // Mirror of serverQueued for the refresh gate. - const serverQueuedRef = useRef([]) - useEffect(() => { - serverQueuedRef.current = serverQueued - }, [serverQueued]) - const fetchedTurnsOnce = useRef(false) - - // Retire local chips the server has acknowledged (by SessionMessage id): - // once an id shows up in a messages frame, an inbox fetch, or a transcript - // user-echo record, the server's own rows are the truth for that send. - const retireAcknowledgedChips = useCallback((ids: Set): void => { - if (ids.size === 0) return - setQueued((prev) => - prev.filter((q) => q.messageId === null || !ids.has(q.messageId)), - ) - }, []) - - // Append items, promoting a queued send when its user turn arrives and - // dropping the backend's echo of a send we already committed locally. - const append = useCallback((incoming: TranscriptItem[]): void => { - const commit: TranscriptItem[] = [] - for (const it of incoming) { - if (it.kind === 'user') { - if (queuedRef.current.some((q) => q.text === it.text)) { - setQueued((prev) => { - const j = prev.findIndex((q) => q.text === it.text) - return j < 0 ? prev : [...prev.slice(0, j), ...prev.slice(j + 1)] - }) - } else { - const fi = flushed.current.indexOf(it.text) - if (fi >= 0) { - flushed.current.splice(fi, 1) - continue - } - } - } - commit.push(it) - } - if (commit.length) setItems((prev) => [...prev, ...commit]) - }, []) - - // Fold a committed event into the footer spend: a result carries the running - // total, so the turn cost is its delta from the previous total. - const applyCost = useCallback((event: CCEvent): void => { - const total = resultCostUsd(event) - if (total == null) return - const prev = costTotal.current - costTotal.current = total - setCost({ - total, - lastStep: prev != null ? Math.max(0, total - prev) : total, - }) - }, []) - - // Pull the structured steps and append any newer than what's on screen. This - // is the sole source of transcript content — the socket only decides *when* - // to call it. Guarded so concurrent wakes coalesce into one trailing refresh. - const refreshSteps = useCallback(async (): Promise => { - if (refreshing.current) { - pendingRefresh.current = true - return - } - refreshing.current = true - try { - // Reconcile the queued region against the server's inbox whenever there - // is anything to reconcile (plus once at startup, to pick up messages - // already pending on reconnect). Chips retire by SessionMessage id: a - // fetched inbox that contains a chip's id means the server owns the - // send — delivered messages vanish, still-pending ones re-render from - // serverQueued with identical text. The v2 `messages` frames keep this - // current between fetches. - const wantTurns = - !fetchedTurnsOnce.current || - queuedRef.current.length > 0 || - serverQueuedRef.current.length > 0 - const turnsPromise = wantTurns - ? api.getAgentSessionTurns(sessionId).catch(() => null) - : Promise.resolve(null) - const records = await api.getAgentSessionRecords(sessionId) - const turns = await turnsPromise - if (turns !== null) { - fetchedTurnsOnce.current = true - setServerQueued( - turns.messages.filter((m) => m.status === 'pending').map((m) => m.body), - ) - retireAcknowledgedChips(new Set(turns.messages.map((m) => m.id))) - } - const ordered = [...records].sort((a, b) => a.feed_seq - b.feed_seq) - const fresh = ordered.filter((st) => st.feed_seq > maxFeed.current) - retireAcknowledgedChips( - new Set( - fresh - .map((st) => st.session_message_id) - .filter((id): id is string => typeof id === 'string'), - ), + + // The committed transcript, derived from the store's record log. Keys ride + // feed_seq (the shared per-session order), so items are stable across + // re-derivations. + const items = useMemo( + () => + snapshot.records + .filter((r) => r.feed_seq > props.minRenderFeedSeq) + .filter(isConnectVisibleRecord) + .flatMap((r) => recordToItems(r, `s${r.feed_seq}`)), + [snapshot.records, props.minRenderFeedSeq], + ) + + // Footer spend: the server's ledger total (the session frame's four cost + // columns — the billing authority, resent on every cost tick) with the + // CC-result fold as the last-turn readout and older-backend fallback (§6: + // record folding is display-only). + const cost = useMemo( + () => + foldCosts( + snapshot.records + .filter((r) => r.source === 'claude_code') + .map((r) => r.payload as CCEvent), + ), + [snapshot.records], + ) + const serverCostUsd = snapshot.session + ? usdNumberFromMillicents( + snapshot.session.cost_tokens + + snapshot.session.cost_sandbox_cpu + + snapshot.session.cost_sandbox_memory + + snapshot.session.cost_fee, ) - if (fresh.length) { - for (const st of fresh) { - maxFeed.current = Math.max(maxFeed.current, st.feed_seq) - applyCost(st.payload as CCEvent) - // Setup-script output chunks drive the "Starting sandbox" sub-line; - // sandbox_ready retires it (the spawn is over). - if (st.source === 'lifecycle') { - if (st.record_type === 'sandbox_setup_output') { - const line = setupOutputLine(st.payload) - if (line) setSetupLine(`${setupOutputHook(st.payload)} · ${line}`) - } else if (st.record_type === 'sandbox_ready') { - setSetupLine(null) - } - } - } - // Lifecycle rows stay off the transcript — the activity line + footer - // carry session state, closing surfaces as the exit notice — except - // the sandbox-ready conversation note (isConnectVisibleRecord). - append( - fresh.filter(isConnectVisibleRecord).flatMap((st) => recordToItems(st, `s${st.feed_seq}`)), - ) - // Drop the live overlay only when a committed ASSISTANT step (or the - // turn's result) lands — those supersede the streamed prose. Other - // records (lifecycle rows, your own relayed message) arrive mid- - // generation; clearing on them would lose the streamed prefix for - // good, since later deltas append to the now-empty overlay. - const supersedes = fresh.some((st) => { - if (st.source === 'lifecycle') return false - const t = (st.payload as CCEvent).type - return t === 'assistant' || t === 'result' - }) - if (supersedes) { - setLiveText('') - setLiveTokens(null) - } - } - } catch { - // Transient fetch failure — the next wake/poll retries. - } finally { - refreshing.current = false - if (pendingRefresh.current) { - pendingRefresh.current = false - void refreshSteps() + : null + + // The setup script's latest output line while the sandbox spawns (streamed + // sandbox_setup_output lifecycle chunks) — the sub-line under "Starting + // sandbox…" that says WHAT the box is doing. sandbox_ready retires it. + const setupLine = useMemo(() => { + for (let i = snapshot.records.length - 1; i >= 0; i--) { + const record = snapshot.records[i] + if (record.source !== 'lifecycle') continue + if (record.record_type === 'sandbox_ready') return null + if (record.record_type === 'sandbox_setup_output') { + const line = setupOutputLine(record.payload) + return line ? `${setupOutputHook(record.payload)} · ${line}` : null } } - }, [api, append, applyCost, retireAcknowledgedChips, sessionId]) + return null + }, [snapshot.records]) + + // Bodies of the server's PENDING inbox messages — the durable queued signal. + const serverQueued = useMemo( + () => snapshot.messages.filter((m) => m.status === 'pending').map((m) => m.body), + [snapshot.messages], + ) + + // Retire local chips the store has acknowledged (by SessionMessage id): once + // an id shows up in a messages frame or a transcript user-echo record, the + // server's own rows are the truth for that send. + useEffect(() => { + setQueued((prev) => { + const remaining = prev.filter( + (q) => q.messageId === null || !snapshot.acknowledgedMessageIds.has(q.messageId), + ) + return remaining.length === prev.length ? prev : remaining + }) + }, [snapshot.acknowledgedMessageIds]) // A closed conversation is over — nothing can ever be sent or received - // again (a send would 409) — so pull the final records, leave one dim - // notice as the sign-off, and exit instead of sitting at the composer. + // again (a send would 409) — so leave one dim notice as the sign-off and + // exit instead of sitting at the composer. The server flushes the final + // records before the closing session frame, so there is nothing to fetch. const finishClosed = useCallback((): void => { if (closingDown.current) return closingDown.current = true if (props.exitState) props.exitState.closed = true - setWorking(false) setNotice('conversation closed') - void refreshSteps().finally(() => exit()) - }, [exit, props.exitState, refreshSteps]) - - // Coalesce bursts of wakes into one refresh a beat later. - const scheduleRefresh = useCallback((): void => { - if (refreshTimer.current) return - refreshTimer.current = setTimeout(() => { - refreshTimer.current = null - void refreshSteps() - }, 250) - }, [refreshSteps]) - - // The socket speaks stream protocol v2: session/snapshot frames carry the - // status word + the server's running cost; messages frames carry the open - // inbox slice (queued chips); records_append is the committed-steps wake - // (the REST refresh stays the transcript source in this port — the SDK swap - // rebases rendering on pushed records); deltas are the ephemeral overlay. - const handleFrame = useCallback( - (frame: StreamFrame): void => { - if (frame.type === 'snapshot' || frame.type === 'session') { - const session = (frame as { session: AgentSession }).session - // The server's authoritative running total; the session frame is - // resent on every cost change, so the footer climbs mid-turn. - applyServerCostUsd( - usdNumberFromMillicents( - session.cost_tokens + - session.cost_sandbox_cpu + - session.cost_sandbox_memory + - session.cost_fee, - ), - ) - const word = sessionStatusWord(session) - if (word && word !== lastStatus.current) { - lastStatus.current = word - setStatus(word) - setWorking(isWorkingStatus(word)) - // Box-up states (working/waiting) mean the session became - // connectable; used to tell a preflight failure from a - // mid-conversation one. - if (['working', 'waiting'].includes(word)) everRunning.current = true - if (word === 'closed') { - finishClosed() - return - } - } - if (frame.type === 'snapshot') { - // The snapshot's inbox is the open slice — seed the queued region. - const messages = (frame as { messages: SessionMessage[] }).messages - setServerQueued( - messages.filter((m) => m.status === 'pending').map((m) => m.body), - ) - retireAcknowledgedChips(new Set(messages.map((m) => m.id))) - } - return - } - if (frame.type === 'messages') { - // The open inbox slice, resent on any change: pending rows are the - // queued truth; a row that flipped delivered retires its chip. - const messages = (frame as { messages: SessionMessage[] }).messages - setServerQueued( - messages.filter((m) => m.status === 'pending').map((m) => m.body), - ) - retireAcknowledgedChips(new Set(messages.map((m) => m.id))) - return - } - if (frame.type === 'error') { - setNotice((frame as { message?: string }).message ?? 'stream error') - return - } - // Ephemeral streaming delta: update the live overlay in place. NOT a - // committed-steps wake, so it must not trigger a steps refresh (that would - // poll the API several times a second during generation). Unknown kinds - // ("thinking" is reserved) are ignored per §3.6. - if (frame.type === 'delta') { - const delta = frame as { kind?: string; text?: string | null; output_tokens?: number | null } - if (delta.kind !== undefined && delta.kind !== 'text') return - if (delta.text) setLiveText((t) => t + delta.text) - if (typeof delta.output_tokens === 'number') setLiveTokens(delta.output_tokens) - return - } - if (frame.type === 'records_append') { - // Track the resume cursor from the pushed records; the debounced REST - // refresh commits them to the transcript. - for (const record of (frame as { records: SessionRecord[] }).records) { - if (typeof record.feed_seq === 'number') { - afterSeq.current = Math.max(afterSeq.current, record.feed_seq) - } + exit() + }, [exit, props.exitState]) + useEffect(() => { + if (statusWord === 'closed') finishClosed() + }, [statusWord, finishClosed]) + + // REST fallback when the stream is unavailable: poll the records + session + // and feed the SAME store through synthetic frames (the REST rows are the + // same wire shapes) — the cursor dedupes, the UI can't tell the difference. + const startPollFallback = useCallback((): void => { + if (polling.current || abort.current.signal.aborted) return + polling.current = true + const tick = async (): Promise => { + try { + const [page, session] = await Promise.all([ + api.getAgentSessionRecordsPage(sessionId, { afterSeq: store.cursor }), + api.getAgentSession(sessionId), + ]) + if (page.records.length) { + store.ingest({ type: 'records_append', records: page.records }) } - scheduleRefresh() - return + store.ingest({ type: 'messages', messages: page.messages ?? [] }) + store.ingest({ type: 'session', session }) + } catch { + // Transient fetch failure — the next tick retries. } - if (frame.type === 'heartbeat' || frame.type === 'done') return - // Unknown frame types are ignored (protocol §3.6). - }, - [applyServerCostUsd, finishClosed, retireAcknowledgedChips, scheduleRefresh], - ) + } + void tick() + const timer = setInterval(() => void tick(), 3000) + abort.current.signal.addEventListener('abort', () => clearInterval(timer), { + once: true, + }) + }, [api, sessionId, store]) // Keep the socket attached across reconnects/resume. A keyed session going // terminal is not the end of the conversation — it idles between turns — so - // we refresh once more, report idle, and re-attach on the next send. - // Watch-only sessions exit when the stream ends. + // we report idle and re-attach on the next send. Watch-only sessions exit + // when the stream ends. const pump = useCallback((): void => { if (streaming.current || abort.current.signal.aborted) return streaming.current = true streamSession({ - token, sessionId, - wsBase, - afterSeq: afterSeq.current, - onFrame: handleFrame, + openSocket, + afterSeq: store.cursor, + onFrame: store.ingest, signal: abort.current.signal, }) - .then(async (outcome) => { + .then((outcome) => { if (abort.current.signal.aborted) return - await refreshSteps() // pull the final steps of the turn before settling - setLiveText('') - setLiveTokens(null) - setWorking(false) + setSendPending(false) if (outcome.type === 'error') setNotice(`stream error: ${outcome.message}`) // A terminal failure before the sandbox ever ran is a preflight/budget // gate — there's no conversation to attend, so report it and exit. @@ -476,7 +295,7 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { // The warm loop can end by closing the conversation (a closing event's // final turn, or an ephemeral session finishing): that's terminal, not // an idle nap — tear down instead of inviting a doomed send. - if (outcome.type === 'done' && lastStatus.current === 'closed') { + if (outcome.type === 'done' && outcome.status === 'closed') { finishClosed() return } @@ -488,67 +307,26 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { }) .catch((err: unknown) => { if (abort.current.signal.aborted) return - // A dropped/unreachable socket is not worth announcing: the poll below - // keeps the transcript current, so the fallback is invisible. Only a - // genuine stream error gets a notice. - if (!(err instanceof StreamUnavailableError)) { - setNotice(`stream error: ${(err as Error).message}`) + if (err instanceof StreamUnavailableError) { + // No socket (old backend, blocked port): fall back to REST polling, + // silently — the store keeps filling, so the fallback is invisible. + startPollFallback() + return } - // No socket: only a watch-only session (which would never poll long) - // exits here. + setNotice(`stream error: ${(err as Error).message}`) if (!canSend) exit() }) .finally(() => { streaming.current = false }) - }, [canSend, exit, finishClosed, handleFrame, refreshSteps, sessionId, token, wsBase]) - - // A status backstop that doesn't depend on the socket: the lifecycle - // transitions may arrive before the socket attaches, so poll the session too - // and drive the same transition handling. Use the derived surface word so this - // matches the stream's `frame.status` exactly (the raw `status` is a different - // vocabulary); fall back to raw for un-keyed sessions. `lastStatus` dedupes - // against socket frames so a transition is only announced once. - const pollStatus = useCallback(async (): Promise => { - try { - const s = await api.getAgentSession(sessionId) - // Same figure the stream's cost_millicents carries, so the footer keeps - // climbing even when the socket never attached (polling fallback). - applyServerCostUsd( - usdNumberFromMillicents( - s.cost_tokens + s.cost_sandbox_cpu + s.cost_sandbox_memory + s.cost_fee, - ), - ) - const word = s.surface?.status ?? s.status - if (word !== lastStatus.current) { - lastStatus.current = word - setStatus(word) - setWorking(isWorkingStatus(word)) - if (['working', 'waiting'].includes(word)) everRunning.current = true - if (word === 'closed') finishClosed() - } - } catch { - // Transient fetch failure — the next tick retries. - } - }, [api, applyServerCostUsd, finishClosed, sessionId]) + }, [canSend, exit, finishClosed, openSocket, sessionId, startPollFallback, store]) useEffect(() => { pump() - scheduleRefresh() // catch steps created between the initial fetch and connect - // A slow poll backs up the socket wake (and covers a socket that never - // connected, so following still updates). The socket wake is the fast path; - // this is just a safety net, so it can be gentle. - const poll = setInterval(scheduleRefresh, 3000) - const statusPoll = setInterval(() => void pollStatus(), 3000) const controller = abort.current - return () => { - controller.abort() - clearInterval(poll) - clearInterval(statusPoll) - if (refreshTimer.current) clearTimeout(refreshTimer.current) - } + return () => controller.abort() // eslint-disable-next-line react-hooks/exhaustive-deps - }, [pump, scheduleRefresh, pollStatus]) + }, [pump]) // Tick an elapsed-seconds counter while the agent works — a steady progress // read alongside the live token counter, and the sole liveness cue during a @@ -587,24 +365,6 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { return () => clearInterval(t) }, [pendingToolKey]) - // When the turn goes idle, commit any still-queued messages into the - // transcript (they weren't relayed) so they never vanish; remember them so a - // late backend echo is de-duplicated. A message sent while already idle - // flushes immediately — it starts a fresh turn rather than queueing. - useEffect(() => { - if (working || queued.length === 0) return - flushed.current.push(...queued.map((q) => q.text)) - const items = queued.map((q) => ({ - key: `q${keyCounter.current++}`, - kind: 'user', - gutter: '›', - text: q.text, - spaceBefore: true, - })) - setItems((prev) => [...prev, ...items]) - setQueued([]) - }, [working, queued]) - const submit = useCallback( (raw: string): void => { const text = raw.trim() @@ -621,11 +381,10 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { setNotice(`stop requested (${s.status}) — the conversation survives`) return } - // Show the message as queued (or, if idle, it flushes to the - // transcript at once), then post it. The POST returns the created - // SessionMessage (protocol v2 §4.2): stamp the chip with its id so - // the first messages frame / inbox fetch / user-echo record carrying - // that id retires it in favour of the server's own row. + // Show the message as queued, then post it. The POST returns the + // created SessionMessage (protocol v2 §4.2): stamp the chip with its + // id so the first messages frame / user-echo record carrying that id + // retires it in favour of the server's own row. setQueued((prev) => [...prev, { text, messageId: null }]) setNotice(null) const created = await api.sendSessionMessage(sessionId, text) @@ -639,7 +398,7 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { return q }) }) - setWorking(true) + setSendPending(true) pump() } catch (err) { setQueued((prev) => { @@ -709,12 +468,12 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { // (cumulative total + the last turn's cost), and the session identity — the // dashboard link rendered as the session id, the model, and the CLI version. // Command hints live in --help. The total prefers the server's ledger figure - // (live via cost_millicents frames, climbing mid-turn); the CC-derived - // result total is the fallback against older backends. + // (live via the session frames' cost columns, climbing mid-turn); the + // CC-derived result total is the fallback against older backends. const totalStr = `$${(serverCostUsd ?? cost.total ?? 0).toFixed(2)}` const lastStepStr = cost.lastStep != null ? ` (Last step: $${cost.lastStep.toFixed(2)})` : '' const metaLine = [ - `${status} · ${totalStr} total${lastStepStr}`, + `${statusWord} · ${totalStr} total${lastStepStr}`, hyperlink(props.sessionUrl, sessionId), ...(props.model ? [props.model] : []), `v${VERSION}`, @@ -731,14 +490,16 @@ export function ConnectApp(props: ConnectAppProps): React.ReactElement { // the live third), naming the tool and ticking its own timer (generating // wins if both somehow read true; the model can't stream past an // unresolved call). - const infraActivity = statusActivityText(status) - const generating = status === 'working' && (liveText !== '' || liveTokens != null) + const liveText = snapshot.liveText + const liveTokens = snapshot.liveOutputTokens + const infraActivity = statusActivityText(statusWord) + const generating = statusWord === 'working' && (liveText !== '' || liveTokens != null) const generatingBits = [ formatDuration(elapsed), ...(liveTokens != null ? [`↓ ${formatTokens(liveTokens)} tokens`] : []), ...(inputActive ? ['esc to interrupt'] : []), ].join(' · ') - const runningTool = status === 'working' && !generating && pendingTools.length > 0 + const runningTool = statusWord === 'working' && !generating && pendingTools.length > 0 const runningToolLabel = pendingTools.length === 1 ? `Running ${pendingTools[0].text}${pendingTools[0].detail ?? ''}` diff --git a/src/ui/SessionView.tsx b/src/ui/SessionView.tsx deleted file mode 100644 index f8fea4d..0000000 --- a/src/ui/SessionView.tsx +++ /dev/null @@ -1,69 +0,0 @@ -import React, { useEffect, useState } from 'react' -import { Box, Text, useApp } from 'ink' -import Spinner from 'ink-spinner' -import { sessionStatusWord, streamSession, type StreamFrame } from '../lib/ws' -import { isConnectVisibleRecord, recordToItems } from '../lib/events' -import type { AgentSession, SessionRecord } from '../lib/types' - -interface Props { - sessionId: string - token: string -} - -export function SessionView({ sessionId, token }: Props): React.ReactElement { - const { exit } = useApp() - const [lines, setLines] = useState([]) - const [status, setStatus] = useState('connecting') - - useEffect(() => { - const controller = new AbortController() - const onFrame = (frame: StreamFrame) => { - switch (frame.type) { - case 'records_append': { - const records = (frame as { records: SessionRecord[] }).records - const rendered = records - .filter(isConnectVisibleRecord) - .flatMap((record) => recordToItems(record, `v${record.feed_seq}`)) - .map((item) => (item.detail ? `${item.text} ${item.detail}` : item.text)) - .filter((line) => line.trim().length > 0) - if (rendered.length) setLines((prev) => [...prev, ...rendered]) - break - } - case 'snapshot': - case 'session': - setStatus(sessionStatusWord((frame as { session: AgentSession }).session)) - break - default: - break // heartbeat/messages/delta/unknown: nothing to draw here - } - } - streamSession({ token, sessionId, onFrame, signal: controller.signal }) - .then((outcome) => { - if (outcome.type === 'error') setStatus(`error: ${outcome.message}`) - else if (outcome.type === 'done') setStatus('done') - exit() - }) - .catch((err: Error) => { - setStatus(`error: ${err.message}`) - exit() - }) - return () => controller.abort() - }, [sessionId, token, exit]) - - const done = status === 'done' - - return ( - - - {done ? : } - - {' '} - session {sessionId}: {status} - - - {lines.map((line, i) => ( - {line} - ))} - - ) -} diff --git a/test/events.test.ts b/test/events.test.ts deleted file mode 100644 index 7096b90..0000000 --- a/test/events.test.ts +++ /dev/null @@ -1,361 +0,0 @@ -import { describe, expect, it } from 'vitest' -import { - clampLines, - collapseToolRuns, - pendingToolCalls, - eventToItems, - foldCosts, - formatDuration, - LineBuffer, - parseEventLine, - resultCostUsd, - statusActivityText, - summarizeToolInput, - type CCEvent, -} from '../src/lib/events' - -describe('LineBuffer', () => { - it('emits only complete lines and buffers the trailing partial', () => { - const buf = new LineBuffer() - expect(buf.push('{"a":1}\n{"b":2}')).toEqual(['{"a":1}']) - expect(buf.push('\n{"c"')).toEqual(['{"b":2}']) - expect(buf.push(':3}\n')).toEqual(['{"c":3}']) - }) - - it('reassembles a single event split across frames', () => { - const buf = new LineBuffer() - expect(buf.push('{"ty')).toEqual([]) - expect(buf.push('pe":"x"}')).toEqual([]) - expect(buf.push('\n')).toEqual(['{"type":"x"}']) - }) - - it('flush releases a trailing partial once', () => { - const buf = new LineBuffer() - buf.push('tail') - expect(buf.flush()).toEqual(['tail']) - expect(buf.flush()).toEqual([]) - }) -}) - -describe('parseEventLine', () => { - it('parses a JSON object', () => { - expect(parseEventLine('{"type":"result"}')).toEqual({ type: 'result' }) - }) - - it('returns null for blank, non-JSON, or non-object lines', () => { - expect(parseEventLine(' ')).toBeNull() - expect(parseEventLine('not json')).toBeNull() - expect(parseEventLine('[1,2,3]')).toBeNull() - expect(parseEventLine('"a string"')).toBeNull() - }) -}) - -describe('summarizeToolInput', () => { - it('shows the file path for file tools', () => { - expect(summarizeToolInput('Read', { file_path: 'src/auth.ts' })).toBe('src/auth.ts') - expect(summarizeToolInput('Edit', { file_path: 'a.ts', old_string: 'x' })).toBe('a.ts') - }) - - it('shows the command for Bash and the pattern for Grep', () => { - expect(summarizeToolInput('Bash', { command: 'ls -la' })).toBe('ls -la') - expect(summarizeToolInput('Grep', { pattern: 'foo', path: 'src' })).toBe('foo in src') - }) - - it('falls back to compact JSON for unknown tools', () => { - expect(summarizeToolInput('Mystery', { a: 1 })).toBe('{"a":1}') - expect(summarizeToolInput('Mystery', {})).toBe('') - }) -}) - -describe('eventToItems', () => { - it('renders assistant text, thinking, and tool calls as grouped items', () => { - const event: CCEvent = { - type: 'assistant', - message: { - content: [ - { type: 'thinking', thinking: 'checking the auth flow' }, - { type: 'text', text: 'Reading the file.' }, - { - type: 'tool_use', - name: 'Read', - input: { file_path: 'src/auth.ts' }, - }, - ], - }, - } - const items = eventToItems(event, 's1') - expect(items.map((i) => i.kind)).toEqual(['thinking', 'assistant', 'tool']) - expect(items[0].gutter).toBe('✻') - expect(items[2]).toMatchObject({ - kind: 'tool', - gutter: '●', - text: 'Read', - detail: '(src/auth.ts)', - }) - // Unique keys within one event. - expect(new Set(items.map((i) => i.key)).size).toBe(items.length) - }) - - it('groups a tool_result under the call with a ⎿ gutter', () => { - const event: CCEvent = { - type: 'user', - message: { - content: [{ type: 'tool_result', tool_use_id: 't1', content: 'file contents' }], - }, - } - const [item] = eventToItems(event, 's2') - expect(item).toMatchObject({ - kind: 'tool_result', - gutter: '⎿', - text: 'file contents', - spaceBefore: false, - }) - }) - - it('unwraps nested text blocks in a tool_result and marks errors', () => { - const event: CCEvent = { - type: 'user', - message: { - content: [ - { - type: 'tool_result', - is_error: true, - content: [{ type: 'text', text: 'boom' }], - }, - ], - }, - } - const [item] = eventToItems(event, 's3') - expect(item.text).toBe('boom') - expect(item.isError).toBe(true) - }) - - it('renders a human user message with a › gutter (not as a tool result)', () => { - const event: CCEvent = { - type: 'user', - message: { content: 'please add tests' }, - } - const [item] = eventToItems(event, 's4') - expect(item).toMatchObject({ - kind: 'user', - gutter: '›', - text: 'please add tests', - spaceBefore: true, - }) - }) - - it('summarizes a result event with duration and cost', () => { - const event: CCEvent = { - type: 'result', - duration_ms: 12300, - total_cost_usd: 0.042, - } - const [item] = eventToItems(event, 's5') - expect(item.kind).toBe('summary') - expect(item.text).toBe('turn complete · 12s · $0.04') - }) - - it('formats long turn durations human-readably (minutes, not 200.7s)', () => { - const [item] = eventToItems({ type: 'result', duration_ms: 200700 }, 's5b') - expect(item.text).toBe('turn complete · 3m 21s') - }) - - it('renders nothing for system events (CC emits an init per query — noise)', () => { - expect(eventToItems({ type: 'system', subtype: 'init', model: 'claude-opus-4-8' }, 's6')).toEqual([]) - expect(eventToItems({ type: 'system', subtype: 'status' }, 's6b')).toEqual([]) - }) - - it('produces no items for an empty assistant turn', () => { - expect(eventToItems({ type: 'assistant', message: { content: [] } }, 's7')).toEqual([]) - }) -}) - -describe('formatDuration', () => { - it('shows whole seconds under a minute, minutes past it', () => { - expect(formatDuration(42.3)).toBe('42s') - expect(formatDuration(200.7)).toBe('3m 21s') - expect(formatDuration(0)).toBe('0s') - }) - - it('never shows 60s (rounds up into the minute)', () => { - expect(formatDuration(59.7)).toBe('1m 0s') - }) -}) - -describe('clampLines', () => { - it('leaves short bodies untouched', () => { - expect(clampLines('a\nb', 6)).toEqual({ body: 'a\nb', more: 0 }) - }) - - it('clamps long bodies and counts the remainder', () => { - const text = Array.from({ length: 10 }, (_, i) => `line ${i}`).join('\n') - const { body, more } = clampLines(text, 6) - expect(body.split('\n')).toHaveLength(6) - expect(more).toBe(4) - }) -}) - -describe('resultCostUsd', () => { - it('returns the total_cost_usd of a result event', () => { - expect(resultCostUsd({ type: 'result', total_cost_usd: 0.4381 })).toBe(0.4381) - }) - - it('is null for non-result events and results without a cost', () => { - expect(resultCostUsd({ type: 'assistant' })).toBeNull() - expect(resultCostUsd({ type: 'result' })).toBeNull() - }) -}) - -describe('foldCosts', () => { - it('is all-null with no result events', () => { - expect(foldCosts([{ type: 'assistant' }, { type: 'user' }])).toEqual({ - total: null, - lastStep: null, - }) - }) - - it('makes the first turn cost the whole cumulative total', () => { - expect(foldCosts([{ type: 'result', total_cost_usd: 0.32 }])).toEqual({ - total: 0.32, - lastStep: 0.32, - }) - }) - - it('takes the latest total and the delta from the previous result as the last step', () => { - // The cumulative result totals from the "u up" prod session: 0.41 then 0.44 - // → total 0.44, last turn 0.03. - const events: CCEvent[] = [ - { type: 'result', total_cost_usd: 0.3224 }, - { type: 'result', total_cost_usd: 0.41 }, - { type: 'result', total_cost_usd: 0.4381 }, - ] - const { total, lastStep } = foldCosts(events) - expect(total).toBeCloseTo(0.4381, 4) - expect(lastStep).toBeCloseTo(0.0281, 4) - }) - - it('never reports a negative last step if a total regresses', () => { - const events: CCEvent[] = [ - { type: 'result', total_cost_usd: 0.5 }, - { type: 'result', total_cost_usd: 0.4 }, - ] - expect(foldCosts(events).lastStep).toBe(0) - }) -}) - -describe('collapseToolRuns', () => { - const tool = (key: string, name: string): Parameters[0][number] => ({ - key, - kind: 'tool', - text: name, - gutter: '●', - }) - const result = (key: string): Parameters[0][number] => ({ - key, - kind: 'tool_result', - text: 'ok', - gutter: '⎿', - }) - const prose = (key: string, text: string): Parameters[0][number] => ({ - key, - kind: 'assistant', - text, - }) - - it('folds a run of Bash calls into one shell-command summary', () => { - const out = collapseToolRuns([ - prose('a', 'Checking.'), - tool('t1', 'Bash'), - result('r1'), - tool('t2', 'Bash'), - result('r2'), - prose('b', 'Done.'), - ]) - expect(out.map((i) => i.kind)).toEqual(['assistant', 'notice', 'assistant']) - expect(out[1].text).toBe('Ran 2 shell commands') - }) - - it('singular for one call, and names mixed tools', () => { - expect(collapseToolRuns([tool('t1', 'Bash'), result('r1')])[0].text).toBe( - 'Ran 1 shell command', - ) - const mixed = collapseToolRuns([tool('t1', 'Bash'), result('r1'), tool('t2', 'Grep')]) - expect(mixed[0].text).toBe('Ran 2 tool calls (Bash, Grep)') - }) - - it('counts files for all-Read runs and splits groups on prose between them', () => { - const out = collapseToolRuns([ - tool('t1', 'Read'), - result('r1'), - prose('a', 'Now the fix.'), - tool('t2', 'Bash'), - result('r2'), - ]) - expect(out.map((i) => i.text)).toEqual(['Read 1 file', 'Now the fix.', 'Ran 1 shell command']) - }) - - it('passes non-tool items through untouched', () => { - const items = [prose('a', 'Hello.')] - expect(collapseToolRuns(items)).toEqual(items) - }) -}) - -describe('pendingToolCalls', () => { - const tool = (key: string, name: string, detail?: string): Parameters[0][number] => ({ - key, - kind: 'tool', - text: name, - detail, - gutter: '●', - }) - const result = (key: string): Parameters[0][number] => ({ - key, - kind: 'tool_result', - text: 'ok', - gutter: '⎿', - }) - const prose = (key: string): Parameters[0][number] => ({ - key, - kind: 'assistant', - text: 'hi', - }) - const summary = (key: string): Parameters[0][number] => ({ - key, - kind: 'summary', - text: 'turn complete', - }) - - it('reports a call whose result has not arrived', () => { - const pending = pendingToolCalls([prose('a'), tool('t1', 'Bash', '(pytest -q)')]) - expect(pending).toHaveLength(1) - expect(pending[0].text).toBe('Bash') - }) - - it('clears once the result lands, FIFO for parallel calls', () => { - expect(pendingToolCalls([tool('t1', 'Bash'), result('r1')])).toHaveLength(0) - const two = pendingToolCalls([tool('t1', 'Bash'), tool('t2', 'Grep'), result('r1')]) - expect(two.map((t) => t.text)).toEqual(['Grep']) - }) - - it('resets on any non-tool item, so an errored old turn never reads as running', () => { - expect(pendingToolCalls([tool('t1', 'Bash'), summary('s1')])).toHaveLength(0) - expect(pendingToolCalls([tool('t1', 'Bash'), prose('a')])).toHaveLength(0) - }) -}) - -describe('statusActivityText', () => { - it('labels the infrastructure phases for the ✻ activity line', () => { - expect(statusActivityText('scheduled')).toBe('Waiting for a worker') - expect(statusActivityText('starting')).toBe('Starting sandbox') - expect(statusActivityText('retrying')).toBe('Retrying after a transient error') - }) - - it('is null for working (the whimsy label takes over) and calm states', () => { - expect(statusActivityText('working')).toBeNull() - expect(statusActivityText('waiting')).toBeNull() - expect(statusActivityText('sleeping')).toBeNull() - expect(statusActivityText('closed')).toBeNull() - expect(statusActivityText('failed')).toBeNull() - expect(statusActivityText('nonsense')).toBeNull() - }) -}) diff --git a/test/stream.test.ts b/test/stream.test.ts new file mode 100644 index 0000000..8b132dc --- /dev/null +++ b/test/stream.test.ts @@ -0,0 +1,31 @@ +import { afterEach, describe, expect, it } from 'vitest' +import { buildStreamUrl, resolveWsBase } from '../src/lib/stream' + +// The CLI-side transport adapter only: the stream machinery itself +// (reconnect/resume/backoff, frame parsing, close-code policy) lives in +// @ellipsis-dev/sdk/stream and is tested there. + +describe('resolveWsBase', () => { + afterEach(() => { + delete process.env.ELLIPSIS_WS_BASE + }) + + it('prefers the explicit env, then derives from the api base', () => { + process.env.ELLIPSIS_WS_BASE = 'wss://explicit.example' + expect(resolveWsBase('https://ignored')).toBe('wss://explicit.example') + delete process.env.ELLIPSIS_WS_BASE + expect(resolveWsBase('https://api.example')).toBe('wss://api.example') + expect(resolveWsBase('http://localhost:5000')).toBe('ws://localhost:5000') + }) +}) + +describe('buildStreamUrl', () => { + it('builds the bearer-door URL with the SDK handshake query', () => { + expect(buildStreamUrl('wss://h', 'session 1', 'protocol=2')).toBe( + 'wss://h/v1/sessions/session%201/stream?protocol=2', + ) + expect(buildStreamUrl('wss://h', 's', 'protocol=2&after_seq=7')).toBe( + 'wss://h/v1/sessions/s/stream?protocol=2&after_seq=7', + ) + }) +}) diff --git a/test/ws.test.ts b/test/ws.test.ts deleted file mode 100644 index c2781c5..0000000 --- a/test/ws.test.ts +++ /dev/null @@ -1,256 +0,0 @@ -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -import { - buildStreamUrl, - classifyCloseCode, - decideReconnect, - nextReconnectDelayMs, - resolveWsBase, - streamSession, - StreamAuthError, - StreamUnavailableError, - type SocketFactory, - type StreamFrame, - type StreamSocket, -} from '../src/lib/ws' - -// A controllable in-memory socket so the reconnect/resume/fallback machinery can -// be driven deterministically with fake timers — no real server. -class FakeSocket implements StreamSocket { - private openCb?: () => void - private msgCb?: (data: string) => void - private closeCb?: (code: number) => void - private errCb?: (err: Error) => void - closed = false - - onOpen(cb: () => void): void { - this.openCb = cb - } - onMessage(cb: (data: string) => void): void { - this.msgCb = cb - } - onClose(cb: (code: number) => void): void { - this.closeCb = cb - } - onError(cb: (err: Error) => void): void { - this.errCb = cb - } - close(): void { - this.closed = true - } - - emitOpen(): void { - this.openCb?.() - } - emitFrame(frame: StreamFrame): void { - this.msgCb?.(JSON.stringify(frame)) - } - emitClose(code: number): void { - this.closeCb?.(code) - } - emitError(err: Error): void { - this.errCb?.(err) - } -} - -function makeFactory(): { factory: SocketFactory; sockets: { url: string; sock: FakeSocket }[] } { - const sockets: { url: string; sock: FakeSocket }[] = [] - const factory: SocketFactory = (url) => { - const sock = new FakeSocket() - sockets.push({ url, sock }) - return sock - } - return { factory, sockets } -} - -// ------------------------------ pure helpers -------------------------------- - -describe('pure helpers', () => { - it('classifies WebSocket close codes', () => { - expect(classifyCloseCode(1000)).toBe('normal') - expect(classifyCloseCode(1008)).toBe('auth') - expect(classifyCloseCode(1002)).toBe('unsupported') // unknown protocol version - expect(classifyCloseCode(1003)).toBe('unsupported') // no ?protocol= param - expect(classifyCloseCode(1011)).toBe('retry') - expect(classifyCloseCode(1013)).toBe('retry') // over capacity - expect(classifyCloseCode(1006)).toBe('retry') - }) - - it('backs off exponentially up to a cap', () => { - expect(nextReconnectDelayMs(1)).toBe(500) - expect(nextReconnectDelayMs(2)).toBe(1000) - expect(nextReconnectDelayMs(3)).toBe(2000) - expect(nextReconnectDelayMs(99)).toBe(8000) - }) - - it('decides auth failures and unsupported closes without retrying', () => { - const base = { everReceivedFrame: true, attempt: 1, maxReconnects: 5 } - expect(decideReconnect({ ...base, closeKind: 'auth' }).action).toBe('fail-auth') - expect(decideReconnect({ ...base, closeKind: 'unsupported' }).action).toBe('fallback') - }) - - it('retries persistently once streaming has worked, briefly otherwise', () => { - // Seen a frame: retry up to maxReconnects, then fall back. - expect( - decideReconnect({ closeKind: 'retry', everReceivedFrame: true, attempt: 4, maxReconnects: 5 }) - .action, - ).toBe('reconnect') - expect( - decideReconnect({ closeKind: 'retry', everReceivedFrame: true, attempt: 5, maxReconnects: 5 }) - .action, - ).toBe('fallback') - // Never connected: bail out after a couple of attempts so polling kicks in. - expect( - decideReconnect({ errored: true, everReceivedFrame: false, attempt: 1, maxReconnects: 5 }) - .action, - ).toBe('reconnect') - expect( - decideReconnect({ errored: true, everReceivedFrame: false, attempt: 2, maxReconnects: 5 }) - .action, - ).toBe('fallback') - }) - - it('builds the stream URL with the required protocol + optional cursor', () => { - expect(buildStreamUrl('wss://h', 'session 1', 0)).toBe( - 'wss://h/v1/sessions/session%201/stream?protocol=2', - ) - expect(buildStreamUrl('wss://h', 's', 7)).toBe( - 'wss://h/v1/sessions/s/stream?protocol=2&after_seq=7', - ) - }) - - it('resolves the ws base from env, then derives it from the api base', () => { - process.env.ELLIPSIS_WS_BASE = 'wss://explicit.example' - expect(resolveWsBase('https://ignored')).toBe('wss://explicit.example') - delete process.env.ELLIPSIS_WS_BASE - expect(resolveWsBase('https://api.example')).toBe('wss://api.example') - expect(resolveWsBase('http://localhost:5000')).toBe('ws://localhost:5000') - }) -}) - -// ---------------------------- streamSession --------------------------------- - -// Minimal v2 frame builders (protocol §3.3). The session frame carries the -// lean wire session; only the fields streamSession reads are populated. -function sessionFrame(word: string, exitStatus: string | null = null): StreamFrame { - return { - type: 'session', - session: { - id: 's', - status: word, - exit_status: exitStatus, - surface: { session: null, run: null, status: word }, - }, - } as unknown as StreamFrame -} - -function recordsFrame(...feedSeqs: number[]): StreamFrame { - return { - type: 'records_append', - records: feedSeqs.map((feed_seq) => ({ - id: `rec-${feed_seq}`, - agent_session_id: 's', - feed_seq, - stream_seq: feed_seq, - source: 'claude_code', - record_type: 'assistant', - record_format: 'claude_stream_json@2.0', - payload: {}, - created_at: 'now', - })), - } as unknown as StreamFrame -} - -describe('streamSession', () => { - beforeEach(() => { - vi.useFakeTimers() - }) - afterEach(() => { - vi.useRealTimers() - delete process.env.ELLIPSIS_WS_BASE - }) - - it('emits every frame and resolves on done with the last session word', async () => { - const { factory, sockets } = makeFactory() - const frames: StreamFrame[] = [] - const p = streamSession({ - token: 't', - sessionId: 's', - wsBase: 'ws://x', - connect: factory, - onFrame: (f) => frames.push(f), - }) - sockets[0].sock.emitOpen() - sockets[0].sock.emitFrame(sessionFrame('working')) - sockets[0].sock.emitFrame(recordsFrame(1)) - // The final session frame carries the end state, then done (§3.3). - sockets[0].sock.emitFrame(sessionFrame('closed', 'completed')) - sockets[0].sock.emitFrame({ type: 'done' }) - - const outcome = await p - expect(outcome).toEqual({ type: 'done', status: 'closed', exitStatus: 'completed' }) - expect(frames.map((f) => f.type)).toEqual([ - 'session', - 'records_append', - 'session', - 'done', - ]) - expect(sockets[0].sock.closed).toBe(true) - }) - - it('reconnects after a drop and resumes from the last record feed_seq', async () => { - const { factory, sockets } = makeFactory() - const p = streamSession({ - token: 't', - sessionId: 's', - wsBase: 'ws://x', - connect: factory, - onFrame: () => {}, - }) - sockets[0].sock.emitOpen() - sockets[0].sock.emitFrame(recordsFrame(1, 2)) - // Only records advance the cursor (§3.4) — a session frame never does. - sockets[0].sock.emitFrame(sessionFrame('working')) - sockets[0].sock.emitClose(1011) // server error: retryable drop - - await vi.advanceTimersByTimeAsync(500) // backoff for attempt 1 - expect(sockets).toHaveLength(2) - expect(sockets[1].url).toBe('ws://x/v1/sessions/s/stream?protocol=2&after_seq=2') - - sockets[1].sock.emitFrame(sessionFrame('closed', null)) - sockets[1].sock.emitFrame({ type: 'done' }) - const outcome = await p - expect(outcome).toEqual({ type: 'done', status: 'closed', exitStatus: null }) - }) - - it('surfaces a server error frame as an error outcome (not a fallback)', async () => { - const { factory, sockets } = makeFactory() - const p = streamSession({ token: 't', sessionId: 's', wsBase: 'ws://x', connect: factory, onFrame: () => {} }) - sockets[0].sock.emitFrame({ type: 'error', message: 'boom' }) - expect(await p).toEqual({ type: 'error', message: 'boom' }) - }) - - it('falls back (throws StreamUnavailableError) on an unsupported close', async () => { - const { factory, sockets } = makeFactory() - const p = streamSession({ token: 't', sessionId: 's', wsBase: 'ws://x', connect: factory, onFrame: () => {} }) - sockets[0].sock.emitClose(1003) - await expect(p).rejects.toBeInstanceOf(StreamUnavailableError) - }) - - it('falls back when the socket never connects after a couple of tries', async () => { - const { factory, sockets } = makeFactory() - const p = streamSession({ token: 't', sessionId: 's', wsBase: 'ws://x', connect: factory, onFrame: () => {} }) - const rejection = expect(p).rejects.toBeInstanceOf(StreamUnavailableError) - sockets[0].sock.emitError(new Error('ECONNREFUSED')) - await vi.advanceTimersByTimeAsync(500) - expect(sockets).toHaveLength(2) - sockets[1].sock.emitError(new Error('ECONNREFUSED')) - await rejection - }) - - it('fails hard (StreamAuthError) on an auth-rejected close', async () => { - const { factory, sockets } = makeFactory() - const p = streamSession({ token: 't', sessionId: 's', wsBase: 'ws://x', connect: factory, onFrame: () => {} }) - sockets[0].sock.emitClose(1008) - await expect(p).rejects.toBeInstanceOf(StreamAuthError) - }) -})