Untangles the interleaved work in a Claude Code session (prompts, mid-turn asides, subagents, loops) into colour-coded streams with a live navigator pane, a…


<h1>claudeflow</h1>
<a href="LICENSE"><img alt="License: Proprietary" src="https://img.shields.io/badge/license-proprietary-7c5cff?style=for-the-badge"></a> <a href="https://github.com/macleodlabs-ai/claudeflow/releases"><img alt="Version 0.3.4" src="https://img.shields.io/badge/version-0.3.4-22d3ee?style=for-the-badge"></a> <img alt="Claude Code 2.1.287+" src="https://img.shields.io/badge/Claude%20Code-2.1.287%2B-d97757?style=for-the-badge&logo=claude&logoColor=white"> <img alt="Tests 74 passing" src="https://img.shields.io/badge/tests-74%20passing-2ea043?style=for-the-badge&logo=checkmarx&logoColor=white"> <a href="#-install"><img alt="Install: /plugin marketplace add macleodlabs-ai/claudeflow" src="https://img.shields.io/badge/%2Fplugin%20marketplace%20add-macleodlabs--ai%2Fclaudeflow-0d1117?style=for-the-badge&logo=gnubash&logoColor=white&labelColor=7c5cff"></a> <a href="#-install"><img alt="Installs" src="https://img.shields.io/endpoint?url=https%3A%2F%2Fraw.githubusercontent.com%2Fmacleodlabs-ai%2Fclaudeflow%2Fstats%2Finstalls.json&style=for-the-badge&logo=download&logoColor=white"></a> <a href="https://github.com/macleodlabs-ai/claudeflow/stargazers"><img alt="GitHub stars" src="https://img.shields.io/github/stars/macleodlabs-ai/claudeflow?style=for-the-badge&logo=github&color=fcc2d7&label=stars"></a> <a href="https://github.com/macleodlabs-ai/claudeflow/network/members"><img alt="Forks" src="https://img.shields.io/github/forks/macleodlabs-ai/claudeflow?style=for-the-badge&logo=github&color=ffd8a8"></a> <a href="https://github.com/macleodlabs-ai/claudeflow/watchers"><img alt="Watchers" src="https://img.shields.io/github/watchers/macleodlabs-ai/claudeflow?style=for-the-badge&logo=github&color=b2f2bb"></a> <img alt="Claude Code mod" src="https://img.shields.io/badge/Claude%20Code-mod-a78bfa?style=flat-square"> <img alt="Platform macOS | Linux" src="https://img.shields.io/badge/platform-macOS%20%7C%20Linux-34d399?style=flat-square&logo=apple&logoColor=white"> <img alt="TypeScript" src="https://img.shields.io/badge/TypeScript-strict-3178c6?style=flat-square&logo=typescript&logoColor=white"> <img alt="Classifier Haiku" src="https://img.shields.io/badge/classifier-Haiku-febc2e?style=flat-square"> <a href="https://github.com/macleodlabs-ai/claudeflow/issues"><img alt="Open issues" src="https://img.shields.io/github/issues/macleodlabs-ai/claudeflow?style=flat-square&color=ffd8a8"></a> <a href="https://github.com/macleodlabs-ai/claudeflow/commits/main"><img alt="Last commit" src="https://img.shields.io/github/last-commit/macleodlabs-ai/claudeflow?style=flat-square&color=8b90c4"></a> <a href="https://macleodlabs.ai"><img alt="Made by Macleod Labs" src="https://img.shields.io/badge/made%20by-Macleod%20Labs-12164a?style=flat-square"></a> <a href="https://macleodlabs.ai/?utm_source=github&utm_medium=readme&utm_campaign=claudeflow"><img alt="Hire Macleod Labs" src="https://img.shields.io/badge/Hire%20us-Claude%20Code%20%26%20agent%20workflows%20built%20for%20your%20team%20%E2%86%92-ff7b72?style=for-the-badge&labelColor=12164a"></a>
<a href="#-install">Install</a> · <a href="#-what-you-get">What you get</a> · <a href="#-use">Use</a> · <a href="#%EF%B8%8F-configure">Configure</a> · <a href="#-troubleshooting">Troubleshooting</a> · <a href="#-work-with-macleod-labs"><b>Hire us</b></a>
A real Claude Code session is never one task. You start a feature, ask a side question mid-turn, leave a /loop watching a deploy, and fan out three subagents to chase a bug. It all lands in one transcript, interleaved.
streams, the mod in this repository, sorts that work into semantic streams as it happens, and gives each one its own colour, status and history.

<table> <tr> <td width="33%" valign="top">
Every stream beside the transcript: its main turn and subagents with live status, what each is doing right now, elapsed time, tool count, and loop countdowns.
</td> <td width="33%" valign="top">
One pill per stream above the prompt, coloured by health. Number keys jump between them; 0 shows everything.
</td> <td width="33%" valign="top">
A thin pastel line down the left of every prompt, reply and tool row. Focus a stream and the rest fold to one-line stubs.
</td> </tr> <tr> <td valign="top">
Structure first (subagents, loops, follow-ups, #tags with autocomplete), then one small Haiku call for anything new. Prompts sent mid-turn go to the right stream.
</td> <td valign="top">
Existing sessions are filed into streams in the background, with batched, parallel classification and a live progress line.
</td> <td valign="top">
Hide finished streams with ✕ and bring them back later. Fold any stream to its last 10 rows, last row, or header. Idle streams fold themselves.
</td> </tr> </table>
Type status, or press status in the bar, and a card opens above the prompt with every piece of work in the session and where it stands. It is answered locally, so it costs no model call and works while a turn is running.

TL-260, ENG-1042) gets its own row: running agents on it say what they are doing and how long they have been quiet; otherwise its latest news, or that its agent failed.✕ close (top right) or your next prompt hides the card.status button, or type status) and scrolls by touch.Driving Claude Code from the Claude app over Remote Control? When your phone connects, streams opens there as an accordion made for touch: one colour-bordered card per stream, a summary row of chips, and what needs you first.
Once a session starts, and every six hours after, streams checks every plugin you have installed against its marketplace. When one has a newer release, an ⬆ update button appears in the bar, the pane and on your phone, and a toast says so. Pressing it (or <kbd>u</kbd>, or /streams update) installs the updates on your Mac and reloads plugins into the running session: no restart.
| Colour | Status | Meaning |
|---|---|---|
| 🟨 | RUNNING | A turn, subagent or loop is working now |
| 🟩 | DONE | Finished cleanly |
| 🟦 | WAITING FOR YOU | On the status card: the stream's last reply asked you something |
| 🟥 | ERROR | A subagent or turn failed |
| 🟧 | STALLED | No activity for longer than expected |
| ↻ | LOOP | A /loop or cron is armed, with a countdown to the next tick. A loop that stops, or lapses 10 minutes without re-arming, drops back to its stream's status |
claude plugin marketplace add macleodlabs-ai/claudeflow
claude plugin install streams@claudeflow
/plugin marketplace add macleodlabs-ai/claudeflow
/plugin install streams@claudeflow
New sessions load it automatically; in a running session, run /reload-plugins. The pane opens by itself in fullscreen terminals at least 144 columns wide. Anywhere else, type /streams.
[!TIP] Several Claude Code accounts? Plugins install per account, so sign in to each one and run the install there (or, if you use claude-sessions, once from each
cc-<client>launcher).
Add the folder as the marketplace. Claude Code reads it straight from disk, so /reload-plugins picks up your changes:
git clone https://github.com/macleodlabs-ai/claudeflow
claude plugin marketplace add ./claudeflow
claude plugin install streams@claudeflow
Or for one session only, without installing:
claude --plugin-dir ./claudeflow/plugins/streams
| Action | Command |
|---|---|
| Update | claude plugin update streams@claudeflow |
| Uninstall | claude plugin uninstall streams@claudeflow |
| Command | What it does |
|---|---|
status or status? | Show the status card above the prompt: every stream's state (running, loop, waiting for you, done) and what it is doing, plus git branch and uncommitted files. Answered locally: no model call, works mid-turn |
/streams | Open the navigator pane |
/streams status | Same as typing status |
/streams update | Check every installed plugin for a newer release, install them, and reload plugins into this session, no restart |
/stream <name> | Focus one stream; others fold to stubs |
/stream off | Show every stream again |
/stream move <name> | Refile the last prompt, and everything after it, under another stream (created if new) when it was sorted wrongly |
/stream | List streams with their summaries |
/streams import | List this project's past sessions |
/streams import <id> | File a past session into streams (re-importing replaces, never duplicates) |
| To | Do |
|---|---|
| See the status of all work | Press status in the bar (<kbd>t</kbd> when the bar has focus), or type status. Click a stream's name on the card to open it; ✕ close or your next prompt hides it |
| Focus a stream | Click its name in the pane, or press its pill |
| Show everything | ← all streams, or the all pill |
| File a prompt by hand | Start it with #name, e.g. #billing why is the total off? Type # and the bar lists matching streams; <kbd>tab</kbd> completes the first |
| Fold a stream | ▾ all cycles all → last 10 → last 1 → header only |
| Archive or restore | ✕ beside a stream; ▸ archived (N) lists them |
| Collapse everything | collapse all / expand all at the top of the pane |
| Fold the pane away | ⇥ hide (<kbd>h</kbd>) folds the docked pane to a ◂ streams tab at the right of the bar; the tab (<kbd>s</kbd>) brings it back at the width it had. It stays folded in new sessions until you open it |
| Switch a stream's chat view | The stream header shows view ◉ full ○ compact: full draws the chat as the session does, with markdown, syntax-highlighted code and diffs; compact is one line per row. Click either, or press <kbd>v</kbd> in the pane |
| Keys | Action |
|---|---|
| <kbd>ctrl</kbd>+<kbd>x</kbd> then <kbd>tab</kbd> | Move focus to the pill bar |
| <kbd>0</kbd> | All streams |
| <kbd>1</kbd>–<kbd>9</kbd> | Jump to stream n |
| <kbd>s</kbd> | Open the pane |
Streams persist per project directory, across sessions.
Settings live in /config under streams, or in settings.json:
{
"pluginConfigs": {
"streams@claudeflow": {
"options": { "chatStyle": "full", "diagnostics": false }
}
}
}
| Setting | Default | Description |
|---|---|---|
chatStyle | full | How a stream's own view draws its chat: full, as the session draws it, with markdown, syntax-highlighted commands and file contents, and edits as coloured diffs; or compact, one line per row. The view switch in a stream's header changes it for the session. |
diagnostics | false | Writes debug.json into the plugin folder every few seconds: what the pane last drew, rows it could not place, and the last background error. Turn on only when troubleshooting. |
| Event | Haiku requests |
|---|---|
| New prompt | 1 |
Follow-up (yes, continue), slash command, #tagged prompt | 0 |
| Filing history | 1 per 25 prompts, plus 1 merge pass, plus 1 per turn that had a mid-turn prompt |
| Symptom | Fix |
|---|---|
| No pill bar | It appears once the first prompt is sorted. Check /plugin lists streams as enabled. |
| Pane doesn't open | The terminal is under 144 columns or not fullscreen. Type /streams. |
Dim streams: … line in the transcript | Claude Code is reporting a failed hook; the line names it. Include it in an issue. |
| Older rows have no stripe | History is still filing; watch the progress line at the top of the pane. |
<table> <tr> <td>
claudeflow is built by Macleod Labs. We build Claude Code mods, agent workflows and AI tooling for engineering teams: custom plugins, multi-agent pipelines, and getting your team productive with Claude Code.
<a href="https://macleodlabs.ai/?utm_source=github&utm_medium=readme&utm_campaign=claudeflow"><img alt="Hire Macleod Labs" src="https://img.shields.io/badge/Hire%20Macleod%20Labs-%E2%86%92-7c5cff?style=for-the-badge&labelColor=12164a"></a>
</td> </tr> </table>
Proprietary · © 2026 Macleod Labs · See LICENSE
hooks/register.tsx 1837 lines1import { atom, read, update } from 'claude-code'
2import type { EngineInterface, Register, RenderElement } from 'claude-code'
3
4import type { AgentRun, ChatStyle, Folded, Health, Stream, StreamRow, StreamRowKind } from '../types'
5import type { Installed, MarketEntry, Update } from './updates'
6import { CHECK_EVERY_MS, manifestPathOf, pluginsDirOf, updatesOf } from './updates'
7import type { BadgeKind, Fold, HistoryItem, HistoryTurn, Proposal, StatusKind, StatusLine } from './classify'
8import {
9 BATCH_SYSTEM,
10 MERGE_SYSTEM,
11 buildBatchPrompt,
12 buildMergePrompt,
13 inParallel,
14 parseBatch,
15 parseMerge,
16 BADGE_BG,
17 BADGE_FG,
18 FOLD_LABEL,
19 HEALTH_TEXT,
20 NEXT_FOLD,
21 clockOf,
22 rowKey,
23 REPLIES_SYSTEM,
24 buildRepliesPrompt,
25 pickReplyStreams,
26 turnsOf,
27 nextPastel,
28 pastelOf,
29 partialTag,
30 gitStatus,
31 sortStatus,
32 statusOf,
33 ticketLines,
34 limitView,
35 codeOf,
36 toolLine,
37 tagMatches,
38 completeTag,
39 REPLY_SYSTEM,
40 buildReplyPrompt,
41 pickReplyStream,
42 readTranscript,
43 HEALTH_COLOR,
44 HEALTH_GLYPH,
45 SYSTEM,
46 healthOf,
47 ago,
48 buildPrompt,
49 fallbackName,
50 isFollowUp,
51 loopKey,
52 oneLine,
53 parseTag,
54 parseVerdict,
55 slug,
56 textKey,
57} from './classify'
58
59const PANE = 'streams'
60const MAX_ROWS = 4000
61const SAVED_ROWS = 400
62
63const streamsA = atom({ plugin: 'streams', key: 'streams' } as const, [])
64const currentA = atom({ plugin: 'streams', key: 'current' } as const, '')
65const focusA = atom({ plugin: 'streams', key: 'focus' } as const, '')
66const viewA = atom({ plugin: 'streams', key: 'view' } as const, '')
67const showArchivedA = atom({ plugin: 'streams', key: 'showArchived' } as const, false)
68const agentsA = atom({ plugin: 'streams', key: 'agents' } as const, {})
69const foldA = atom({ plugin: 'streams', key: 'fold' } as const, {})
70const turnStartedAtA = atom({ plugin: 'streams', key: 'turnStartedAt' } as const, 0)
71const tickA = atom({ plugin: 'streams', key: 'tick' } as const, 0)
72const loopsA = atom({ plugin: 'streams', key: 'loops' } as const, {})
73const chatStyleA = atom({ plugin: 'streams', key: 'chatStyle' } as const, '')
74const tagHintA = atom({ plugin: 'streams', key: 'tagHint' } as const, null)
75const updatesA = atom({ plugin: 'streams', key: 'updates' } as const, [])
76const updatingA = atom({ plugin: 'streams', key: 'updating' } as const, false)
77const mobileOpenA = atom({ plugin: 'streams', key: 'mobileOpen' } as const, '')
78const paneCollapsedA = atom({ plugin: 'streams', key: 'paneCollapsed' } as const, false)
79const statusOpenA = atom({ plugin: 'streams', key: 'statusOpen' } as const, false)
80const statusGitA = atom({ plugin: 'streams', key: 'statusGit' } as const, [])
81const busyA = atom({ plugin: 'streams', key: 'busy' } as const, false)
82const rowsA = atom({ plugin: 'streams', key: 'rows' } as const, [])
83const agentStreamA = atom({ plugin: 'streams', key: 'agentStream' } as const, {})
84const liveA = atom({ plugin: 'streams', key: 'live' } as const, {})
85const inflightA = atom({ plugin: 'streams', key: 'inflight' } as const, {})
86const outcomeA = atom({ plugin: 'streams', key: 'outcome' } as const, {})
87const healthA = atom({ plugin: 'streams', key: 'health' } as const, {})
88const historyFiledA = atom({ plugin: 'streams', key: 'historyFiled' } as const, false)
89const importProgressA = atom({ plugin: 'streams', key: 'importProgress' } as const, { label: '', done: 0, total: 0 })
90const keyVersionA = atom({ plugin: 'streams', key: 'keyVersion' } as const, 0)
91/** 2: rows keyed by rowKey (a uuid by its first four groups). Bump when the keys change. */
92const KEY_VERSION = 2
93const transcriptA = atom({ plugin: 'streams', key: 'transcript' } as const, '')
94const foldedA = atom({ plugin: 'streams', key: 'folded' } as const, [])
95const loopStreamA = atom({ plugin: 'streams', key: 'loopStream' } as const, {})
96const ROW = { plugin: 'streams', key: 'rowStream' } as const
97const COLOR = { plugin: 'streams', key: 'streamColor' } as const
98
99type Saved = { streams: Stream[]; rows: StreamRow[]; loopStream: Record<string, string> }
100type $ = EngineInterface
101
102/**
103 * Slow work (model calls, the history import) cannot run in the hook that asks for it: the engine cuts a
104 * dispatch's calls once its hook returns. So hooks queue it here and session.start's timer drains it.
105 */
106type Job =
107 | { kind: 'route'; uuid: string; rowId: string; text: string; turnSid: string; folded: readonly Folded[] }
108 | { kind: 'import'; path: string; isCurrent: boolean }
109const jobs: Job[] = []
110let draining = false
111/** The last failure of background work, for the diagnostics file. */
112let lastError = ''
113let isWorking = false
114
115/** Starts the worker's timer once per load, from whichever hook first has work: timers outlive the hook. */
116function work($: $, job?: Job) {
117 if (job) jobs.push(job)
118 if (isWorking) return
119 isWorking = true
120 $.clock.every(1000, () => {
121 void drain($).catch(err => {
122 lastError = `${String(err)} ${(err as Error)?.stack ?? ''}`.slice(0, 600)
123 $.ui.log(`streams: background work failed: ${String(err)}`)
124 })
125 void tick($).catch(() => {})
126 })
127}
128
129async function drain($: $) {
130 if (draining) return
131 draining = true
132 try {
133 for (let job = jobs.shift(); job; job = jobs.shift()) {
134 if (job.kind === 'import') await importHistory($, job.path, job.isCurrent)
135 else await reroute($, job)
136 }
137 } finally {
138 draining = false
139 }
140}
141
142/** Moves the pane's clocks once a second, only while something runs. */
143async function tick($: $) {
144 const [busy, agents, loops] = await Promise.all([read($, busyA), read($, agentsA), read($, loopsA)])
145 if (!busy && !Object.values(agents).some(a => a.status === 'running') && Object.keys(loops).length === 0) return
146 const now = await $.clock.now()
147 if (Object.values(loops).some(l => lapsed(l, now)))
148 await update($, loopsA, m => Object.fromEntries(Object.entries(m).filter(([, l]) => !lapsed(l, now))))
149 await update($, tickA, () => now)
150}
151
152/** A reply filed under the turn's stream moves to the mid-turn prompt it answers, if it answers one. */
153async function reroute($: $, job: Job & { kind: 'route' }) {
154 const sid = await routeReply($, job.turnSid, job.folded, job.text)
155 if (sid === job.turnSid) return
156 await fileAs($, job.uuid, sid)
157 await fileAs($, textKey(job.text), sid)
158 await update($, rowsA, list => list.map(row => (row.id === job.rowId ? { ...row, streamId: sid } : row)))
159 await touch($, sid, s => ({ rows: s.rows + 1 }))
160 await touch($, job.turnSid, s => ({ rows: Math.max(0, s.rows - 1) }))
161}
162
163// Diagnostics while the matching is being proven: transcript rows drawn with no stream, written by the heartbeat.
164const unmatched = new Map<string, { component: string; requestId: string; head: string }>()
165
166function noteMatched(component: string, requestId: string) {
167 unmatched.delete(`${component}:${requestId}`)
168}
169
170function noteUnmatched(component: string, requestId: string, text: string) {
171 if (unmatched.size >= 60) return
172 unmatched.set(`${component}:${requestId}`, { component, requestId, head: text.slice(0, 80) })
173}
174
175/** Off unless the `diagnostics` setting is on: an installed copy should not write into its own folder. */
176let isDiagnosing = false
177
178async function writeDiagnostics($: $) {
179 const [rows, streams, historyFiled] = await Promise.all([read($, rowsA), read($, streamsA), read($, historyFiledA)])
180 const perStream: Record<string, number> = {}
181 for (const row of rows) perStream[row.streamId] = (perStream[row.streamId] ?? 0) + 1
182 await $.fs.write(
183 `${$.plugin.root}/debug.json`,
184 JSON.stringify(
185 {
186 at: await $.clock.now(),
187 historyFiled,
188 importing,
189 lastError,
190 lastPane,
191 queuedJobs: jobs.map(j => j.kind),
192 streams: streams.map(s => s.id),
193 rowsPerStream: perStream,
194 lastRows: rows.slice(-15).map(r => ({ id: r.id, streamId: r.streamId, kind: r.kind, head: r.text.slice(0, 60) })),
195 unmatched: [...unmatched.values()],
196 },
197 null,
198 2,
199 ),
200 )
201}
202
203/**
204 * The pane's window is the engine's: scrolled down while the pane was long, it stays put when the content
205 * shrinks, and the pane looks empty. Past the last drawn row, bring it back to the top.
206 */
207async function unstick($: $) {
208 const scroll = lastPane.scroll as { offset: number } | undefined
209 const rows = Number(lastPane.rows ?? 0)
210 lastPane = { ...lastPane, panes: await $.ui.panes() }
211 if (scroll && scroll.offset > 0 && scroll.offset >= rows) {
212 const r = await $.ui.scroll({ in: PANE, to: 'start' })
213 lastPane = { ...lastPane, unstuck: r }
214 }
215}
216
217/** The docked pane's width as last drawn: what it reopens at after being folded to the side tab. */
218let dockColumns = 0
219const PANE_KEY = 'streams:pane'
220type PaneSaved = { collapsed: boolean; columns: number }
221
222/** Folds the docked pane away to a tab at the bar's right end, keeping its width for when it comes back. */
223async function collapsePane($: $) {
224 await $.store.set(PANE_KEY, { collapsed: true, columns: dockColumns } satisfies PaneSaved)
225 await update($, paneCollapsedA, () => true)
226 await $.ui.close({ id: PANE })
227}
228
229/** Brings the pane back from the side tab at the width it had (a width the person dragged to wins anyway). */
230async function expandPane($: $, focus = false) {
231 const saved = (await $.store.get(PANE_KEY)) as PaneSaved | undefined
232 await update($, paneCollapsedA, () => false)
233 await $.store.set(PANE_KEY, { collapsed: false, columns: saved?.columns ?? 0 } satisfies PaneSaved)
234 await $.ui.open({ id: PANE, title: 'Streams', ...(saved?.columns ? { columns: saved.columns } : {}), ...(focus ? { focus: true as const } : {}) })
235}
236
237/** The pane's last draw, for the diagnostics file: when, how long, what it drew, or what it threw. */
238let lastPane: Record<string, unknown> = {}
239
240type PaneRender = Parameters<$['ui']['resolve']>[0] & {
241 props: { bodyColumns: number; placement: string; scroll?: { offset: number; bodyRows: number } }
242}
243
244async function timedPane($: $, e: PaneRender, draw: () => Promise<RenderElement>): Promise<RenderElement> {
245 const began = Date.now()
246 const view = await read($, viewA)
247 try {
248 const tree = await draw()
249 const drawn = JSON.stringify(tree)
250 // An upper bound on the rows drawn: every Text and Button is at most a row.
251 const rows = (drawn.match(/"type":"(Text|Button)"/g) ?? []).length
252 if (e.props.placement === 'dock') dockColumns = e.props.bodyColumns
253 lastPane = { at: began, ms: Date.now() - began, view, columns: e.props.bodyColumns, placement: e.props.placement, surface: e.surface, size: drawn.length, rows, scroll: e.props.scroll }
254 return tree
255 } catch (err) {
256 lastPane = { at: began, ms: Date.now() - began, view, columns: e.props.bodyColumns, error: `${String(err)} ${(err as Error)?.stack ?? ''}`.slice(0, 600) }
257 throw err
258 }
259}
260
261// Lost on reload, and that is fine: both only bridge a moment.
262let pendingSpawn = ''
263let pendingKind: StreamRowKind = 'prompt'
264
265const storeKey = (cwd: string) => `streams:v1:${cwd}`
266
267async function ensureStream($: $, name: string, summary: string): Promise<string> {
268 const id = slug(name)
269 const now = await $.clock.now()
270 await update($, streamsA, list =>
271 list.some(s => s.id === id)
272 ? list
273 : [
274 ...list,
275 {
276 id,
277 name,
278 summary,
279 createdAt: now,
280 lastAt: now,
281 rows: 0,
282 agents: 0,
283 loops: 0,
284 color: nextPastel(list.filter(s => !s.archived).map(s => s.color)),
285 },
286 ],
287 )
288 await paintStreams($)
289 return id
290}
291
292/** Files a transcript row (by uuid, tool_use_id or text key) under a stream. */
293const fileAs = ($: $, id: string, sid: string) => $.state.set({ ...ROW, id: rowKey(id) }, sid)
294
295const colorOf = (s: Stream): string => s.color ?? pastelOf(s.id)
296
297/** Mirrors each stream's colour where transcript rows read it, each its own. */
298async function paintStreams($: $) {
299 const streams = await read($, streamsA)
300 await Promise.all(streams.map(s => $.state.set({ ...COLOR, id: s.id }, colorOf(s))))
301}
302
303async function touch($: $, id: string, patch: (s: Stream) => Partial<Stream> = () => ({})) {
304 const now = await $.clock.now()
305 await update($, streamsA, list => list.map(s => (s.id === id ? { ...s, lastAt: now, ...patch(s) } : s)))
306}
307
308async function refreshStatus($: $) {
309 const [focus, current] = await Promise.all([read($, focusA), read($, currentA)])
310 $.ui.status(focus ? `◉ stream ${focus}` : current ? `stream ${current}` : undefined)
311}
312
313/** Hybrid routing: follow-ups stay put, everything else asks Haiku which stream it continues. */
314async function classify($: $, text: string): Promise<string> {
315 const [streams, current] = await Promise.all([read($, streamsA), read($, currentA)])
316 if (current && isFollowUp(text)) return current
317 // Archived streams are put away: only a #tag brings one back. The rest are offered with their state, so a
318 // finished side task is not taken for open work in the same area.
319 const recent = streams.filter(s => !s.archived).sort((a, b) => b.lastAt - a.lastAt).slice(0, 15)
320 const [now, health] = await Promise.all([$.clock.now(), healthNow($, recent)])
321 const r = await $.model.complete({ model: 'haiku', system: SYSTEM, prompt: buildPrompt(recent, current, text, now, health), maxTokens: 150 })
322 const v = r.isAnswered ? parseVerdict(r.text, recent) : undefined
323 if (!v) return current || ensureStream($, fallbackName(text), oneLine(text, 120))
324 if (v.kind === 'new') return ensureStream($, v.name, v.summary)
325 // A stream keeps the goal it was created with: a summary rewritten on every match drifts wider until it
326 // matches everything nearby.
327 return v.id
328}
329
330/** The transcript id a row was filed under: its message's uuid. */
331const uuidOf = (row: StreamRow): string => (row.id.startsWith('h:') ? row.id.split(':')[1] : row.id.split(':')[0]) ?? row.id
332
333/**
334 * `/stream move <name>`: the current stream's last prompt, and everything after it (replies, tools,
335 * subagents), filed under `<name>` instead, created when no stream has that name. Fixes a wrong guess.
336 */
337async function moveLast($: $, name: string): Promise<string> {
338 const [rows, from, streams] = await Promise.all([read($, rowsA), read($, currentA), read($, streamsA)])
339 if (!from) return 'No prompt has been filed yet.'
340 const known = streams.find(s => s.id === slug(name) || s.name.toLowerCase() === name.trim().toLowerCase())
341 const own = rows.filter(r => r.streamId === from)
342 const start = own.findLastIndex(r => r.kind === 'prompt' && !r.agentId)
343 if (start < 0) return `Nothing in ${from} to move.`
344 const moving = own.slice(start)
345 const to = known?.id ?? (await ensureStream($, name.trim(), oneLine(moving[0]?.text ?? '', 120)))
346 if (to === from) return `That prompt is already in ${to}.`
347 const ids = new Set(moving.map(r => r.id))
348 const agentIds = new Set(moving.flatMap(r => (r.agentId ? [r.agentId] : [])))
349 await update($, rowsA, list => list.map(r => (ids.has(r.id) ? { ...r, streamId: to } : r)))
350 await Promise.all(
351 moving.flatMap(r => [
352 fileAs($, uuidOf(r), to),
353 ...(r.toolId ? [fileAs($, r.toolId, to)] : []),
354 ...(r.kind === 'reply' ? [fileAs($, textKey(r.text), to)] : []),
355 ]),
356 )
357 if (agentIds.size) await update($, agentStreamA, m => Object.fromEntries(Object.entries(m).map(([a, sid]) => [a, agentIds.has(a) ? to : sid])))
358 await touch($, to, s => ({ rows: s.rows + moving.length, archived: false }))
359 await touch($, from, s => ({ rows: Math.max(0, s.rows - moving.length) }))
360 await update($, currentA, () => to)
361 await paintStreams($)
362 await refreshStatus($)
363 await save($)
364 return `Moved the last prompt and ${moving.length - 1} row${moving.length === 2 ? '' : 's'} after it from ${from} to ${to}.`
365}
366
367/** A background task's notification names its task; route it to whoever spawned that. */
368async function streamOfNotification($: $, text: string): Promise<string> {
369 const map = await read($, agentStreamA)
370 const hit = Object.keys(map).find(id => text.includes(id))
371 return (hit && map[hit]) || read($, currentA)
372}
373
374type Block = { type: string; text?: string; id?: string; name?: string; input?: unknown }
375
376type AppendedRow = {
377 message: { type: string; role?: string; isMeta?: true; content: readonly unknown[] }
378 agentId?: string
379}
380
381async function record($: $, e: AppendedRow, uuid: string) {
382 const msg = e.message
383 if (msg.type === 'attachment') {
384 // A prompt sent mid-turn is stored as this row: link it to the stream prompt.submit filed it in.
385 const prompt = (msg.content as readonly Block[]).find(b => b.type === 'text')?.text ?? ''
386 const hit = (await read($, foldedA)).findLast(f => prompt.includes(f.text) || f.text.includes(prompt.trim()))
387 // Matched to the prompt it carries; any other attachment wears the current stream's line.
388 await fileAs($, uuid, hit && prompt.trim() ? hit.streamId : await inStream($, e.agentId))
389 return
390 }
391 if (msg.isMeta) {
392 // Notices, reminders and deliveries are not activity, but they sit in the stream's part of the transcript.
393 const sid = await inStream($, e.agentId)
394 if (sid) await fileAs($, uuid, sid)
395 return
396 }
397 let sid: string
398 let folded: readonly Folded[] = []
399 if (e.agentId) {
400 const map = await read($, agentStreamA)
401 sid = map[e.agentId] ?? (pendingSpawn || (await read($, currentA)))
402 if (!map[e.agentId] && sid) await update($, agentStreamA, m => ({ ...m, [e.agentId as string]: sid }))
403 } else {
404 sid = await read($, currentA)
405 if (msg.role === 'assistant') folded = await read($, foldedA)
406 }
407 if (!sid) return
408 const marks: Promise<unknown>[] = [fileAs($, uuid, sid)]
409 const now = await $.clock.now()
410 const rows: StreamRow[] = []
411 const blocks = msg.content as readonly Block[]
412 blocks.forEach((b, i) => {
413 const id = `${uuid}:${i}`
414 const base = { id, streamId: sid, agentId: e.agentId, at: now }
415 if (b.type === 'text' && b.text) {
416 if (msg.type === 'system') rows.push({ ...base, kind: 'notice', text: b.text })
417 else if (msg.role === 'assistant') {
418 rows.push({ ...base, kind: 'reply', text: b.text })
419 marks.push(fileAs($, textKey(b.text), sid))
420 // Filed under the turn now; the worker moves it if it answers a prompt sent mid-turn.
421 if (folded.length > 0) work($, { kind: 'route', uuid, rowId: id, text: b.text, turnSid: sid, folded })
422 } else if (msg.role === 'user' && !b.text.trim().startsWith('<')) {
423 // A typed shell command and its output (<bash-input>, <bash-stdout>) are not prompts, as the import has it.
424 rows.push({ ...base, kind: e.agentId ? 'prompt' : pendingKind, text: b.text })
425 }
426 } else if (b.type === 'tool_use' && b.id) {
427 marks.push(fileAs($, b.id, sid))
428 const input = (b.input ?? {}) as { description?: string }
429 rows.push(
430 b.name === 'Agent'
431 ? { ...base, kind: 'agent', text: input.description ?? 'subagent' }
432 : { ...base, kind: 'tool', text: toolLine(b.name ?? '', b.input), toolId: b.id, ...withCode(b.name ?? '', b.input) },
433 )
434 }
435 })
436 await Promise.all(marks)
437 if (rows.length === 0) return
438 const agentId = e.agentId
439 const latest = rows.at(-1)
440 if (agentId && latest) {
441 const tools = rows.filter(r => r.kind === 'tool').length
442 await update($, agentsA, m =>
443 m[agentId] ? { ...m, [agentId]: { ...m[agentId], lastAt: now, last: oneLine(latest.text, 120), tools: m[agentId].tools + tools } } : m,
444 )
445 }
446 await update($, rowsA, list => [...list, ...rows].slice(-MAX_ROWS))
447 await touch($, sid, s => ({ rows: s.rows + rows.length }))
448}
449
450/** While prompts wait in the running turn, a reply piece goes to whichever of them, or the turn's own task, it answers. */
451async function routeReply($: $, turnSid: string, folded: readonly Folded[], text: string): Promise<string> {
452 const turn = { streamId: turnSid, text: (await read($, streamsA)).find(s => s.id === turnSid)?.summary ?? turnSid }
453 const r = await $.model.complete({ model: 'haiku', system: REPLY_SYSTEM, prompt: buildReplyPrompt(turn, folded, text), maxTokens: 20 })
454 return r.isAnswered ? pickReplyStream(r.text, turn, folded) : turnSid
455}
456
457async function save($: $) {
458 const [cwd, streams, rows, loopStream] = await Promise.all([
459 $.session.cwd(),
460 read($, streamsA),
461 read($, rowsA),
462 read($, loopStreamA),
463 ])
464 const saved: Saved = { streams, rows: rows.slice(-SAVED_ROWS), loopStream }
465 await $.store.set(storeKey(cwd), saved)
466}
467
468/** Every few seconds: recompute each stream's health, and say so when one stalls or background work finishes. */
469async function beat($: $) {
470 await unstick($).catch(() => {})
471 if (isDiagnosing) await writeDiagnostics($).catch(() => {})
472 const [streams, current, busy, live, inflight, outcome, before, now] = await Promise.all([
473 read($, streamsA),
474 read($, currentA),
475 read($, busyA),
476 read($, liveA),
477 read($, inflightA),
478 read($, outcomeA),
479 read($, healthA),
480 $.clock.now(),
481 ])
482 const after: Record<string, Health> = {}
483 for (const s of streams) {
484 after[s.id] = healthOf({
485 now,
486 lastAt: s.lastAt,
487 isTurnOn: busy && s.id === current,
488 liveAgents: Object.values(live).filter(id => id === s.id).length,
489 inflight: inflight[s.id] ?? 0,
490 outcome: outcome[s.id],
491 })
492 const was = before[s.id]
493 // A stream that wakes up opens again, whatever it was folded to.
494 if (was !== 'running' && after[s.id] === 'running') await update($, foldA, ({ [s.id]: _, ...rest }) => rest)
495 if (was !== 'stalled' && after[s.id] === 'stalled') $.ui.toast(`stream ${s.id} looks stalled: nothing for ${ago(now - s.lastAt)}`)
496 if (was === 'running' && after[s.id] === 'done' && s.id !== current) $.ui.toast(`stream ${s.id} finished`)
497 }
498 const changed = streams.some(s => before[s.id] !== after[s.id])
499 if (changed) await update($, healthA, () => after)
500}
501
502/** Under this, one read; over it (a long session's transcript), the file is streamed: a read refuses past 4 MiB. */
503const READ_LIMIT = 4_000_000
504
505/** A file's whole text, however long: long transcripts are exactly the sessions that need filing. */
506async function readWhole($: $, path: string): Promise<string> {
507 const { size } = await $.fs.stat(path)
508 if (size < READ_LIMIT) return String(await $.fs.read(path))
509 const parts: string[] = []
510 for await (const piece of $.process.spawn({ argv: ['cat', path] })) {
511 if ('text' in piece && piece.stream === 'stdout') parts.push(piece.text)
512 }
513 return parts.join('')
514}
515
516// One import at a time in this environment; a reload starts a fresh one, which is safe (see below).
517let importing = false
518
519/** Model calls in flight at once while filing history. */
520const LANES = 5
521/** Prompts per classifying call: enough for context, few enough to answer reliably. */
522const BATCH = 25
523
524const progress = ($: $, label: string, done: number, total: number) => update($, importProgressA, () => ({ label, done, total }))
525
526/**
527 * Files a transcript into streams, in parallel: prompts are classified in batches several calls at once, the
528 * names each batch proposed are merged in one pass, and turns with prompts sent mid-turn have their replies
529 * routed several at once. Then every row is filed, newest stream state saved, progress shown in the pane.
530 *
531 * This session's own transcript (isCurrent) first clears the rows this session recorded, so a run cut short
532 * by a reload is simply run again; another session's rows keep their own times and replace an earlier
533 * import of the same rows rather than doubling them.
534 */
535async function importHistory($: $, path: string, isCurrent: boolean) {
536 if (importing) return
537 importing = true
538 try {
539 await progress($, 'reading', 0, 1)
540 const [jsonl, { startedAt }, began] = await Promise.all([readWhole($, path), $.session.usage(), $.clock.now()])
541 const turns = turnsOf(readTranscript(String(jsonl)))
542 const prompts = turns.flatMap(t => [t.prompt, ...t.items.filter(i => i.kind === 'prompt')]) as (HistoryItem & { kind: 'prompt' })[]
543
544 // 1. Prompts that need no model: #tags, and follow-ups (which take the stream of the prompt before them).
545 const sids = new Map<HistoryItem, string>()
546 const texts = new Map<HistoryItem, string>()
547 for (const p of prompts) {
548 const tag = parseTag(p.text)
549 texts.set(p, tag ? tag.rest : p.text)
550 if (tag) sids.set(p, await ensureStream($, tag.name, oneLine(tag.rest, 120)))
551 }
552 const open = prompts.filter(p => !sids.has(p) && !isFollowUp(p.text))
553
554 // 2. The rest, classified in batches, several at once.
555 const streams = await read($, streamsA)
556 const batches = Array.from({ length: Math.ceil(open.length / BATCH) }, (_, i) => open.slice(i * BATCH, (i + 1) * BATCH))
557 let done = 0
558 await progress($, 'sorting prompts', 0, open.length)
559 const labels = (
560 await inParallel(batches, LANES, async batch => {
561 const said = batch.map(p => texts.get(p) ?? p.text)
562 const r = await $.model.complete({ model: 'haiku', system: BATCH_SYSTEM, prompt: buildBatchPrompt(streams, said), maxTokens: 60 + batch.length * 16 })
563 done += batch.length
564 await progress($, 'sorting prompts', done, open.length)
565 return r.isAnswered ? parseBatch(r.text, batch.length) : batch.map(() => '')
566 })
567 ).flat()
568
569 // 3. One pass merges the names the batches proposed on their own.
570 const known = new Set(streams.map(s => s.id))
571 const proposals = new Map<string, Proposal>()
572 open.forEach((p, i) => {
573 const label = labels[i] || fallbackName(texts.get(p) ?? p.text)
574 labels[i] = label
575 if (known.has(label)) return
576 const one = proposals.get(label) ?? { name: label, count: 0, samples: [] }
577 one.count += 1
578 if (one.samples.length < 2) one.samples.push(texts.get(p) ?? p.text)
579 proposals.set(label, one)
580 })
581 let merged: Record<string, string> = {}
582 if (proposals.size > 1) {
583 await progress($, 'merging streams', 0, 1)
584 const names = [...proposals.keys()]
585 const r = await $.model.complete({ model: 'haiku', system: MERGE_SYSTEM, prompt: buildMergePrompt(streams, [...proposals.values()]), maxTokens: 80 + names.length * 24 })
586 merged = r.isAnswered ? parseMerge(r.text, names) : {}
587 }
588 for (const [i, p] of open.entries()) {
589 const label = labels[i] ?? ''
590 const final = merged[label] ?? label
591 sids.set(p, known.has(final) ? final : await ensureStream($, final, oneLine(texts.get(p) ?? p.text, 120)))
592 }
593 // Follow-ups, in order: the stream of the prompt before them.
594 let previous = await read($, currentA)
595 for (const p of prompts) {
596 const sid = sids.get(p) ?? previous
597 sids.set(p, sid)
598 previous = sid
599 }
600
601 // 4. Replies in turns that had prompts sent mid-turn: which thread each answers, several turns at once.
602 const mixed = turns.filter(t => t.items.some(i => i.kind === 'prompt') && t.items.some(i => i.kind === 'reply'))
603 const replyTo = new Map<HistoryItem, string>()
604 done = 0
605 await progress($, 'routing replies', 0, mixed.length)
606 await inParallel(mixed, LANES, async (turn: HistoryTurn) => {
607 const turnSid = sids.get(turn.prompt) ?? previous
608 const first = turn.items.findIndex(i => i.kind === 'prompt')
609 const folded: Folded[] = turn.items.filter(i => i.kind === 'prompt').map(i => ({ streamId: sids.get(i) ?? turnSid, text: texts.get(i) ?? '' }))
610 const replies = turn.items.slice(first).filter(i => i.kind === 'reply') as (HistoryItem & { kind: 'reply' })[]
611 const self = { streamId: turnSid, text: texts.get(turn.prompt) ?? '' }
612 const r = await $.model.complete({ model: 'haiku', system: REPLIES_SYSTEM, prompt: buildRepliesPrompt(self, folded, replies.map(i => i.text)), maxTokens: 40 + replies.length * 16 })
613 const picks = r.isAnswered ? pickReplyStreams(r.text, self, folded, replies.length) : []
614 replies.forEach((item, j) => replyTo.set(item, picks[j] ?? turnSid))
615 done += 1
616 await progress($, 'routing replies', done, mixed.length)
617 })
618
619 // 5. File every row.
620 await progress($, 'filing rows', 0, turns.length)
621 const rows: StreamRow[] = []
622 let at = startedAt
623 for (const turn of turns) {
624 const turnSid = sids.get(turn.prompt) ?? previous
625 for (const item of [turn.prompt, ...turn.items]) {
626 try {
627 const sid = item.kind === 'prompt' ? (sids.get(item) ?? turnSid) : item.kind === 'reply' ? (replyTo.get(item) ?? turnSid) : turnSid
628 const when = isCurrent ? at++ : (item.at ?? at++)
629 await fileAs($, item.kind === 'tool' ? item.id : item.uuid, sid)
630 if (item.kind === 'reply') await fileAs($, textKey(item.text), sid)
631 const input = (item.kind === 'tool' ? item.input ?? {} : {}) as { description?: string }
632 const text =
633 item.kind === 'prompt'
634 ? (texts.get(item) ?? item.text)
635 : item.kind === 'reply'
636 ? item.text
637 : item.name === 'Agent'
638 ? (input.description ?? 'subagent')
639 : toolLine(item.name, item.input)
640 const kind: StreamRowKind = item.kind === 'tool' ? (item.name === 'Agent' ? 'agent' : 'tool') : item.kind
641 const code = kind === 'tool' && item.kind === 'tool' ? withCode(item.name, item.input) : {}
642 const toolId = item.kind === 'tool' ? { toolId: item.id } : {}
643 rows.push({ id: `h:${item.uuid}:${item.kind === 'tool' ? item.id : rows.length}`, streamId: sid, kind, text, at: when, ...code, ...toolId })
644 } catch (err) {
645 lastError = `skipped a row: ${String(err)}`.slice(0, 600)
646 }
647 }
648 }
649 const imported = new Set(rows.map(r => r.id))
650 await update($, rowsA, list =>
651 [
652 ...list.filter(row => !imported.has(row.id) && (!isCurrent || row.at < startedAt || row.at >= began)),
653 ...rows,
654 ]
655 .sort((a, b) => a.at - b.at)
656 .slice(-MAX_ROWS),
657 )
658 const counts: Record<string, number> = {}
659 for (const row of await read($, rowsA)) counts[row.streamId] = (counts[row.streamId] ?? 0) + 1
660 await update($, streamsA, list => list.map(s => ({ ...s, rows: counts[s.id] ?? 0 })))
661 if (isCurrent) {
662 await update($, currentA, () => previous)
663 await update($, historyFiledA, () => true)
664 }
665 await save($)
666 $.ui.toast(`streams: filed ${prompts.length} prompts into ${new Set(sids.values()).size} streams`)
667 } finally {
668 importing = false
669 await progress($, '', 0, 0).catch(() => {})
670 }
671}
672
673/**
674 * Each stream's health as of now, for drawing: from the same live facts the agent rows show, so a header
675 * never says idle beside a running agent. The heartbeat's stored verdict only drives its notices.
676 */
677async function healthNow($: $, streams: readonly Stream[]): Promise<Record<string, Health>> {
678 const [busy, current, agents, inflight, outcome, now] = await Promise.all([
679 read($, busyA),
680 read($, currentA),
681 read($, agentsA),
682 read($, inflightA),
683 read($, outcomeA),
684 read($, tickA).then(() => $.clock.now()),
685 ])
686 const running = Object.values(agents).filter(a => a.status === 'running')
687 return Object.fromEntries(
688 streams.map(s => [
689 s.id,
690 healthOf({
691 now,
692 lastAt: Math.max(s.lastAt, ...running.filter(a => a.streamId === s.id).map(a => a.lastAt)),
693 isTurnOn: busy && s.id === current,
694 liveAgents: running.filter(a => a.streamId === s.id).length,
695 inflight: inflight[s.id] ?? 0,
696 outcome: outcome[s.id],
697 }),
698 ]),
699 )
700}
701
702type Loops = Record<string, { kind: 'wakeup' | 'cron'; nextAt: number; label: string }>
703
704/** A wakeup that fired this long ago without re-arming ended by not scheduling another tick. */
705const LAPSE_MS = 10 * 60_000
706const lapsed = (l: { kind: string; nextAt: number }, now: number) => l.kind === 'wakeup' && now > l.nextAt + LAPSE_MS
707
708type LoopArgs = { delaySeconds?: number; stop?: boolean; cron?: string; reason?: string }
709
710/** A self-paced wakeup or a cron job arms a loop in its stream; stopping or deleting it disarms it. */
711async function noteLoop($: $, sid: string, tool: string, args: LoopArgs) {
712 if (tool === 'ScheduleWakeup') {
713 // A session has one self-paced loop: stopping it clears it whichever stream it was filed under, and
714 // re-arming it from another stream moves it there.
715 const others = (m: Loops): Loops => Object.fromEntries(Object.entries(m).filter(([, l]) => l.kind !== 'wakeup'))
716 if (args.stop) return update($, loopsA, others)
717 const nextAt = (await $.clock.now()) + (args.delaySeconds ?? 60) * 1000
718 return update($, loopsA, m => ({ ...others(m), [sid]: { kind: 'wakeup' as const, nextAt, label: args.reason ?? '' } }))
719 }
720 if (tool === 'CronCreate') return update($, loopsA, m => ({ ...m, [sid]: { kind: 'cron' as const, nextAt: 0, label: args.cron ?? '' } }))
721 if (tool === 'CronDelete') return update($, loopsA, ({ [sid]: _, ...rest }) => rest)
722}
723
724/**
725 * `/streams import` lists this project's past sessions; `/streams import <id>` (or a transcript path) files
726 * that session into streams in the background, its progress in the pane.
727 */
728async function importCommand($: $, arg: string): Promise<string> {
729 const transcript = await read($, transcriptA)
730 if (!transcript) return 'Send any prompt first: the session list comes from where this session keeps its transcript.'
731 const dir = transcript.slice(0, transcript.lastIndexOf('/'))
732 if (!arg) {
733 const now = await $.clock.now()
734 const sessions = (await $.fs.list(dir))
735 .filter(f => f.kind === 'file' && f.name.endsWith('.jsonl'))
736 .sort((a, b) => b.mtimeMs - a.mtimeMs)
737 .slice(0, 12)
738 if (sessions.length === 0) return 'No past sessions found for this project.'
739 const lines = sessions.map(f => {
740 const id = f.name.replace(/\.jsonl$/, '')
741 const mark = `${dir}/${f.name}` === transcript ? ' (this session)' : ''
742 return ` ${id} ${(f.size / 1_048_576).toFixed(1)} MB ${ago(now - f.mtimeMs)} ago${mark}`
743 })
744 return `This project's sessions, newest first:\n${lines.join('\n')}\n\nFile one into streams: /streams import <id>`
745 }
746 const path = arg.includes('/') ? arg : `${dir}/${arg.replace(/\.jsonl$/, '')}.jsonl`
747 if (!(await $.fs.exists(path))) return `No transcript at ${path}.`
748 work($, { kind: 'import', path, isCurrent: path === transcript })
749 await $.ui.open({ id: PANE, title: 'Streams' })
750 return `Filing ${arg} into streams in the background; the pane shows its progress.`
751}
752
753async function inStream($: $, agentId: string | undefined): Promise<string> {
754 if (agentId) {
755 const map = await read($, agentStreamA)
756 if (map[agentId]) return map[agentId]
757 }
758 return read($, currentA)
759}
760
761/** Opening a stream shows it in the pane and focuses the transcript on it: one act, as the bar's pills do. */
762async function openStream($: $, id: string) {
763 await update($, viewA, () => id)
764 await focusOn($, id)
765 await $.ui.scroll({ in: PANE, to: 'start' }).catch(() => {})
766}
767
768async function setArchived($: $, id: string, archived: boolean) {
769 await update($, streamsA, list => list.map(s => (s.id === id ? { ...s, archived } : s)))
770 if (archived) {
771 if ((await read($, focusA)) === id) await focusOn($, '')
772 if ((await read($, viewA)) === id) await update($, viewA, () => '')
773 }
774 await save($)
775}
776
777async function focusOn($: $, id: string) {
778 await update($, focusA, () => id)
779 await refreshStatus($)
780}
781
782/** Every ticket named in a prompt or an agent's task, with its state and latest news. */
783async function ticketStatus($: $, streamLines: readonly StatusLine[], now: number): Promise<StatusLine[]> {
784 const [rows, agents] = await Promise.all([read($, rowsA), read($, agentsA)])
785 return ticketLines({
786 rows,
787 agents: Object.values(agents),
788 streamKind: Object.fromEntries(streamLines.map(l => [l.id, l.kind ?? 'idle'])),
789 now,
790 })
791}
792
793/** A typed `status` or `status?` is a request for the card, answered here without a model turn. */
794const STATUS_ASK = /^\s*status\s*\??\s*$/i
795
796/** Shows the status card above the prompt, its git rows read now; the stream rows stay live as it shows. */
797async function openStatus($: $, band = '') {
798 await update($, statusOpenA, () => true)
799 // Opened from the bar, the band holds the keys: its ring goes to the card, so ↑↓ scroll it and Esc closes it.
800 if (band) void $.ui.focus({ requestId: band, key: 'status-close' }).catch(() => {})
801 const git = await $.process
802 .run(['git', 'status', '--porcelain=v1', '--branch'], { timeoutMs: 5000 })
803 .then(r => (r.exitCode === 0 ? gitStatus(r.stdout) : []))
804 .catch(() => [])
805 await update($, statusGitA, () => git)
806}
807
808/** Every stream's status row, in the order that needs the person first. */
809async function streamStatus($: $, streams: readonly Stream[], health: Record<string, Health>, now: number): Promise<StatusLine[]> {
810 const [agents, rows, loops] = await Promise.all([read($, agentsA), read($, rowsA), read($, loopsA)])
811 const lines = streams
812 .filter(s => !s.archived)
813 .map(s => {
814 const loop = loops[s.id]
815 return statusOf({
816 stream: s,
817 health: health[s.id] ?? 'idle',
818 ...(loop && !lapsed(loop, now) ? { loop } : {}),
819 running: Object.values(agents)
820 .filter(a => a.streamId === s.id && a.status === 'running')
821 .sort((a, b) => b.lastAt - a.lastAt),
822 lastSaid: rows.findLast(r => r.streamId === s.id && (r.kind === 'prompt' || r.kind === 'reply')),
823 now,
824 })
825 })
826 return sortStatus(lines, Object.fromEntries(streams.map(s => [s.id, s.lastAt])))
827}
828
829const UPDATE_CHECK_KEY = 'streams:updates:checkedAt'
830
831/**
832 * Which installed plugins have a newer release: each marketplace refreshed at most every six hours,
833 * then every installed plugin's version held against the one its marketplace now lists.
834 */
835async function checkUpdates($: $, force = false): Promise<Update[]> {
836 const run = (argv: string[], timeoutMs = 30_000) => $.process.run(argv, { timeoutMs })
837 const now = await $.clock.now()
838 const last = Number((await $.store.get(UPDATE_CHECK_KEY)) ?? 0)
839 if (force || now - last > CHECK_EVERY_MS) {
840 await run(['claude', 'plugin', 'marketplace', 'update'], 180_000).catch(() => undefined)
841 await $.store.set(UPDATE_CHECK_KEY, now)
842 }
843 const listed = await run(['claude', 'plugin', 'list', '--json'])
844 const installed = ((JSON.parse(listed.stdout || '{}') as { installed?: Installed[] }).installed ?? []).filter(p => p.id.includes('@'))
845 const readJson = async (path: string): Promise<unknown> => JSON.parse(String(await $.fs.read(path)))
846 const markets = new Map<string, MarketEntry[]>()
847 const latest: Record<string, string | undefined> = {}
848 for (const p of installed) {
849 const [name = '', market = ''] = p.id.split('@')
850 const dir = pluginsDirOf(p.installPath)
851 if (!dir) continue
852 const marketDir = `${dir}/marketplaces/${market}`
853 if (!markets.has(market)) {
854 const manifest = (await readJson(`${marketDir}/.claude-plugin/marketplace.json`).catch(() => ({}))) as { plugins?: MarketEntry[] }
855 markets.set(market, manifest.plugins ?? [])
856 }
857 const entry = markets.get(market)?.find(e => e.name === name)
858 if (!entry) continue
859 const path = manifestPathOf(marketDir, entry)
860 latest[p.id] = entry.version ?? (path ? ((await readJson(path).catch(() => ({}))) as { version?: string }).version : undefined)
861 }
862 const found = updatesOf(installed, latest)
863 const before = (await read($, updatesA)).map(u => u.id).join()
864 await update($, updatesA, () => found)
865 if (found.length && found.map(u => u.id).join() !== before)
866 $.ui.toast(`${found.length} plugin update${found.length === 1 ? '' : 's'} ready: press ⬆ update`)
867 return found
868}
869
870/** Installs every update found, then reloads the plugins into this session: no restart, from any surface. */
871async function applyUpdates($: $): Promise<string> {
872 const updates = await read($, updatesA)
873 if (!updates.length || (await read($, updatingA))) return updates.length ? 'An update is already running.' : 'Every plugin is up to date.'
874 await update($, updatingA, () => true)
875 const done: string[] = []
876 const failed: string[] = []
877 for (const u of updates) {
878 const r = await $.process
879 .run(['claude', 'plugin', 'update', u.id], { timeoutMs: 180_000 })
880 .catch(err => ({ exitCode: 1, stdout: '', stderr: String(err) }))
881 if (r.exitCode === 0) done.push(`${u.id} ${u.from} → ${u.to}`)
882 else failed.push(`${u.id}: ${oneLine(r.stderr || r.stdout, 160)}`)
883 }
884 await update($, updatesA, list => list.filter(u => !done.some(d => d.startsWith(`${u.id} `))))
885 await update($, updatingA, () => false)
886 const said = [done.length ? `Updated ${done.join(', ')}.` : '', failed.length ? `Failed: ${failed.join('; ')}.` : ''].filter(Boolean).join(' ')
887 if (done.length) {
888 $.ui.toast(`${said} Reloading plugins…`)
889 // The reload replaces this module, so it runs once this press or command has answered.
890 void $.clock
891 .sleep(300)
892 .then(() => $.command.run({ command: 'reload-plugins' }))
893 .catch(() => $.ui.toast('Updated: run /reload-plugins to load it'))
894 } else $.ui.toast(said)
895 return said
896}
897
898export const register: Register = (on, options) => {
899 isDiagnosing = options.diagnostics === true
900 defaultStyle = options.chatStyle === 'compact' ? 'compact' : 'full'
901
902 on('session.start', async ($, e, next) => {
903 const saved = (await $.store.get(storeKey(e.cwd))) as Saved | undefined
904 if (saved && (await read($, streamsA)).length === 0) {
905 await update($, streamsA, () => saved.streams)
906 await update($, rowsA, () => saved.rows)
907 await update($, loopStreamA, () => saved.loopStream)
908 }
909 await $.command.register({ name: 'streams', description: 'Open the streams navigator' })
910 await $.command.register({
911 name: 'stream',
912 description: 'Focus a stream (/stream <name>), or show everything again (/stream off)',
913 })
914 await paintStreams($)
915 await refreshStatus($)
916 // Rows filed under an older key scheme are not found by today's lookups: file the history again.
917 if ((await read($, keyVersionA)) !== KEY_VERSION) {
918 await update($, historyFiledA, () => false)
919 await update($, keyVersionA, () => KEY_VERSION)
920 }
921 const transcript = await read($, transcriptA)
922 if (transcript && !(await read($, historyFiledA))) work($, { kind: 'import', path: transcript, isCurrent: true })
923 const pane = (await $.store.get(PANE_KEY)) as PaneSaved | undefined
924 if (pane?.collapsed) await update($, paneCollapsedA, () => true)
925 else if (e.isInteractive) void $.ui.open({ id: PANE, title: 'Streams' })
926 $.clock.every(5000, () => void beat($).catch(() => {}))
927 if (e.isInteractive) {
928 void checkUpdates($).catch(() => {})
929 $.clock.every(CHECK_EVERY_MS, () => void checkUpdates($).catch(() => {}))
930 }
931 work($)
932 return next(e)
933 })
934
935 // A phone joining the session gets the streams accordion without asking for it.
936 on('session.attach', { surface: 'mobile' }, async ($, e, next) => {
937 const r = await next(e)
938 void $.ui.open({ id: PANE, title: 'Streams' }).catch(() => {})
939 return r
940 })
941
942 on('command.run', { command: 'streams' }, async ($, e) => {
943 const [verb, ...rest] = e.args.trim().split(/\s+/)
944 if (verb === 'import') return { text: await importCommand($, rest.join(' ')) }
945 if (verb === 'status') {
946 await openStatus($)
947 return { text: 'Status card shown above the prompt.' }
948 }
949 if (verb === 'update') {
950 await checkUpdates($, true)
951 return { text: await applyUpdates($) }
952 }
953 await update($, paneCollapsedA, () => false)
954 const opened = await $.ui.open({ id: PANE, title: 'Streams', focus: true })
955 await $.ui.scroll({ in: PANE, to: 'start' }).catch(() => {})
956 const surfaces = await $.session.surfaces().catch(() => [] as const)
957 const where = `Attached: ${surfaces.join(', ') || 'none reported'}. Pane ${opened.isPlaced ? 'drawn' : `waiting: ${opened.reason}`}.`
958 return { text: `Streams navigator opened. ${where} \`/streams import\` files a past session of this project into streams.` }
959 })
960
961 on('command.run', { command: 'stream' }, async ($, e) => {
962 const arg = e.args.trim()
963 if (!arg) {
964 const streams = await read($, streamsA)
965 return { text: streams.length ? streams.map(s => `${s.id}: ${s.summary}`).join('\n') : 'No streams yet.' }
966 }
967 if (arg === 'off') {
968 await focusOn($, '')
969 return { text: 'Showing every stream.' }
970 }
971 const move = /^move\s+(.+)$/.exec(arg)
972 if (move?.[1]) return { text: await moveLast($, move[1]) }
973 const id = await ensureStream($, arg, '')
974 await update($, currentA, () => id)
975 await update($, viewA, () => id)
976 await focusOn($, id)
977 return { text: `Focused on ${id}. Rows from other streams collapse; ctrl+o expands them.` }
978 })
979
980 on('prompt.submit', async ($, e, next) => {
981 await update($, tagHintA, () => null)
982 // From the phone the card only helps where the app draws it; otherwise the question goes to Claude as asked.
983 const phoneDraws = e.origin.kind === 'bridge' && (await $.session.surfaces().catch((): readonly string[] => [])).includes('mobile')
984 if (STATUS_ASK.test(e.text) && (e.origin.kind === 'composer' || phoneDraws)) {
985 await openStatus($)
986 if (phoneDraws) void $.ui.open({ id: PANE, title: 'Streams' }).catch(() => {})
987 return { drop: 'status shown above the prompt' }
988 }
989 await update($, statusOpenA, () => false)
990 let text = e.text
991 let id: string
992 pendingKind = 'prompt'
993 const tag = parseTag(text)
994 if (tag) {
995 id = await ensureStream($, tag.name, oneLine(tag.rest, 120))
996 text = tag.rest
997 } else if (e.origin.kind === 'scheduled-trigger') {
998 const key = loopKey(text)
999 const known = (await read($, loopStreamA))[key]
1000 id = known ?? (await classify($, text))
1001 if (!known) await update($, loopStreamA, m => ({ ...m, [key]: id }))
1002 await touch($, id, s => ({ loops: s.loops + 1 }))
1003 pendingKind = 'loop'
1004 } else if (e.origin.kind === 'task-notification') {
1005 id = await streamOfNotification($, text)
1006 pendingKind = 'notice'
1007 } else {
1008 id = await classify($, text)
1009 }
1010 if (!id) return next(text === e.text ? e : { ...e, text })
1011 if ((await read($, streamsA)).find(s => s.id === id)?.archived) {
1012 await setArchived($, id, false)
1013 $.ui.toast(`stream ${id} restored: a new prompt belongs to it`)
1014 }
1015 if (e.turnId && e.origin.kind !== 'task-notification') {
1016 // Sent mid-turn: filed in its own stream now (its row is an attachment the recorder skips),
1017 // and the running turn keeps its stream; replies are split between them as they come.
1018 const now = await $.clock.now()
1019 await update($, foldedA, list => [...list, { streamId: id, text }])
1020 await update($, rowsA, list => [...list, { id: `q:${now}`, streamId: id, kind: 'prompt' as const, text, at: now }].slice(-MAX_ROWS))
1021 await touch($, id, s => ({ rows: s.rows + 1 }))
1022 } else {
1023 await update($, currentA, () => id)
1024 await touch($, id)
1025 await refreshStatus($)
1026 }
1027 return next(text === e.text ? e : { ...e, text })
1028 }).catch(($, e, next) => next(e))
1029
1030 // The prompt hook is the one place the transcript's path is named: keep it, and file the history once.
1031 on('classic.UserPromptSubmit', async ($, e, next) => {
1032 if (e.transcript_path) await update($, transcriptA, () => e.transcript_path)
1033 if (e.transcript_path && !(await read($, historyFiledA)) && !jobs.some(j => j.kind === 'import')) {
1034 work($, { kind: 'import', path: e.transcript_path, isCurrent: true })
1035 }
1036 return next(e)
1037 })
1038
1039 on('agent.spawn', async ($, e, next) => {
1040 const map = await read($, agentStreamA)
1041 const sid = (e.parentAgentId && map[e.parentAgentId]) || (await read($, currentA))
1042 pendingSpawn = sid
1043 const r = await next(e)
1044 // The agent has started: a failure here must not fail the hook (a .catch would spawn it twice).
1045 if (r.agentId && sid) {
1046 const agentId = r.agentId
1047 const now = await $.clock.now()
1048 const run: AgentRun = { id: agentId, streamId: sid, description: e.description || e.subagentType, status: 'running', startedAt: now, lastAt: now, last: 'starting', tools: 0 }
1049 await update($, agentStreamA, m => ({ ...m, [agentId]: sid }))
1050 .then(() => update($, liveA, m => ({ ...m, [agentId]: sid })))
1051 .then(() => update($, agentsA, m => ({ ...m, [agentId]: run })))
1052 .then(() => touch($, sid, s => ({ agents: s.agents + 1 })))
1053 .catch(err => $.ui.log(`streams: could not file an agent: ${String(err)}`))
1054 }
1055 return r
1056 })
1057
1058 on('tool.call', async ($, e, next) => {
1059 const sid = await inStream($, e.agentId)
1060 if (sid) await noteLoop($, sid, String(e.tool), e as unknown as LoopArgs).catch(() => {})
1061 const bump = (d: number) => (sid ? update($, inflightA, m => ({ ...m, [sid]: Math.max(0, (m[sid] ?? 0) + d) })) : Promise.resolve())
1062 await bump(1)
1063 try {
1064 return await next(e)
1065 } finally {
1066 await bump(-1).catch(() => {})
1067 }
1068 })
1069
1070 on('session.append', async ($, e, next) => {
1071 const r = await next(e)
1072 // The row is stored by now: a failure here must not fail the hook (a .catch would append it twice).
1073 await record($, e, e.uuid).catch(err => $.ui.log(`streams: could not file a row: ${String(err)}`))
1074 return r
1075 })
1076
1077 on('turn.start', async ($, e, next) => {
1078 await update($, busyA, () => true)
1079 const startedAt = await $.clock.now()
1080 await update($, turnStartedAtA, () => startedAt)
1081 const sid = await read($, currentA)
1082 if (sid) await update($, outcomeA, ({ [sid]: _, ...rest }) => rest)
1083 return next(e)
1084 })
1085
1086 on('turn.complete', async ($, e, next) => {
1087 const r = await next(e)
1088 // A subagent's run is one turn of its loop: its end is the agent's end.
1089 const agentId = e.agentId
1090 if (agentId) {
1091 await update($, liveA, ({ [agentId]: _, ...rest }) => rest)
1092 const endedAt = await $.clock.now()
1093 const status: AgentRun['status'] = e.reason === 'answer' ? 'done' : 'error'
1094 await update($, agentsA, m => (m[agentId] ? { ...m, [agentId]: { ...m[agentId], status, endedAt } } : m))
1095 }
1096 if (!e.agentId) {
1097 await update($, foldedA, () => [])
1098 const sid = await read($, currentA)
1099 if (sid) await update($, outcomeA, m => ({ ...m, [sid]: e.reason }))
1100 await update($, busyA, () => false)
1101 await save($)
1102 }
1103 return r
1104 })
1105
1106 // The shortcut bar: every stream a colour-coded pill, its colour the heartbeat's verdict.
1107 // `#` at the start of the prompt completes stream names: the bar lists the matches, Tab takes the first,
1108 // and a tag naming a known stream is painted in that stream's colour.
1109 on('prompt.edit', async ($, e, next) => {
1110 const streams = await read($, streamsA)
1111 const typing = partialTag(e.text, e.cursor)
1112 if (e.key?.key === 'tab' && !e.key.shift && typing !== undefined) {
1113 const [first] = tagMatches(streams, typing)
1114 if (first) {
1115 await update($, tagHintA, () => null)
1116 return completeTag(e.text, e.cursor, first)
1117 }
1118 }
1119 const box = await next(e)
1120 const partial = partialTag(box.text, box.cursor)
1121 const matches = partial === undefined ? [] : tagMatches(streams, partial)
1122 await update($, tagHintA, () => (partial === undefined ? null : { partial, matches }))
1123 const named = /^\s*#([\w-]+)/.exec(box.text)
1124 const known = named?.[1] && streams.find(s => s.id === named[1])
1125 if (!named || !known) return box
1126 const start = box.text.indexOf('#')
1127 return { ...box, decorations: [...(box.decorations ?? []), { start, end: start + named[0].trim().length, color: colorOf(known), bold: true }] }
1128 })
1129
1130 on('ui.focus', { component: 'AbovePrompt' }, async ($, e, next) => {
1131 if (!e.element && (await read($, statusOpenA))) await update($, statusOpenA, () => false)
1132 return next(e)
1133 })
1134
1135 on('ui.render', { component: 'AbovePrompt' }, async ($, e, next) => {
1136 const [streams, focus, live] = await Promise.all([read($, streamsA), read($, focusA), read($, liveA)])
1137 const health = await healthNow($, streams)
1138 const hint = await read($, tagHintA)
1139 if (hint && !e.props.hasSurvey) {
1140 const { Box, Text } = $.ui.resolve(e)
1141 if (hint.matches.length === 0)
1142 return <Text dimColor>{`#${hint.partial} new stream`}</Text>
1143 return (
1144 <Box gap={1}>
1145 <Text dimColor>{`#${hint.partial} →`}</Text>
1146 {hint.matches.map((id, i) => {
1147 const s = streams.find(x => x.id === id)
1148 return (
1149 <Text key={`tag:${id}`} color={s ? colorOf(s) : undefined} bold={i === 0}>
1150 {id}
1151 </Text>
1152 )
1153 })}
1154 <Text dimColor>tab to complete</Text>
1155 </Box>
1156 )
1157 }
1158 if (!e.props.hasSurvey && e.surface !== 'mobile' && (await read($, statusOpenA))) {
1159 const { Box, Button, Text } = $.ui.resolve(e)
1160 await read($, tickA)
1161 const now = await $.clock.now()
1162 const git: StatusLine[] = await read($, statusGitA)
1163 const streamLines = await streamStatus($, streams, health, now)
1164 const tickets = (await ticketStatus($, streamLines, now)).slice(0, TICKET_ROWS)
1165 const lines = [...git, ...tickets, ...streamLines]
1166 const shown = lines.slice(0, STATUS_ROWS + tickets.length)
1167 const limits = ((await $.session.usage().catch(() => undefined))?.rateLimits ?? []).map(l => limitView(l, now))
1168 const width = e.props.bodyColumns
1169 const areaW = Math.min(24, Math.max(10, ...shown.map(l => l.area.length + 2)))
1170 // 16: room for a limit's bar and percent (`▰▰▰▰▱▱▱▱▱▱ 38%`).
1171 const stateW = Math.min(28, Math.max(limits.length ? 16 : 8, ...shown.map(l => l.state.length + 2)))
1172 const close = () => update($, statusOpenA, () => false)
1173 return (
1174 <Box flexDirection="column" borderStyle="round" borderColor="#8b949e" paddingX={1}>
1175 <Box key="st-title" gap={1}>
1176 <Text bold>Status</Text>
1177 <Text dimColor>
1178 {streams.filter(s => !s.archived).length} streams · {new Date(now).toTimeString().slice(0, 5)}
1179 </Text>
1180 <Text dimColor>· ctrl+x tab, then ↑↓ scroll · q close</Text>
1181 <Box flexGrow={1} />
1182 <Button key="status-close" plain dimColor label="✕ close" hotkey="q" onPress={close} />
1183 </Box>
1184 <Box key="st-head">
1185 <Box width={areaW} flexShrink={0}>
1186 <Text dimColor bold>Area</Text>
1187 </Box>
1188 <Box width={stateW} flexShrink={0}>
1189 <Text dimColor bold>State</Text>
1190 </Box>
1191 <Text dimColor bold>Detail</Text>
1192 </Box>
1193 {shown.flatMap((l, i) => {
1194 // With tickets on the card, tickets and streams each get a heading; without, the card reads as before.
1195 const heading =
1196 tickets.length && (l.id.startsWith('ticket:') ? i === git.length : i === git.length + tickets.length) ? (
1197 <Box key={`st-h:${l.id}`} marginTop={1}>
1198 <Text bold color="#d0bfff">{l.id.startsWith('ticket:') ? 'Tickets' : 'Streams'}</Text>
1199 </Box>
1200 ) : nullhooks/updates.ts 48 lines1/** One installed plugin with a newer release in its marketplace. */
2export type Update = { id: string; from: string; to: string }
3
4/** What `claude plugin list --json` says of one installed plugin. */
5export type Installed = { id: string; version: string; installPath: string }
6
7/** A plugin entry in a marketplace's `marketplace.json`. */
8export type MarketEntry = { name: string; version?: string; source?: unknown }
9
10/** Hours between marketplace refreshes: each one fetches every marketplace's repository. */
11export const CHECK_EVERY_MS = 6 * 3600_000
12
13const SEMVER = /^v?(\d+)\.(\d+)\.(\d+)(?:-([\w.]+))?$/
14
15/** Whether `to` is a later release than `from`; false when either is not a version (a git sha). */
16export function isNewer(to: string, from: string): boolean {
17 const a = SEMVER.exec(to.trim())
18 const b = SEMVER.exec(from.trim())
19 if (!a || !b) return false
20 for (const i of [1, 2, 3]) {
21 const d = Number(a[i]) - Number(b[i])
22 if (d) return d > 0
23 }
24 // 1.2.0 is later than 1.2.0-beta; two pre-releases compare as text.
25 if (!a[4] !== !b[4]) return !a[4]
26 return (a[4] ?? '') > (b[4] ?? '')
27}
28
29/** The plugins folder an installed plugin lives under: `<config>/plugins`, from its cache path. */
30export const pluginsDirOf = (installPath: string): string | undefined => {
31 const at = installPath.indexOf('/plugins/cache/')
32 return at < 0 ? undefined : installPath.slice(0, at + '/plugins'.length)
33}
34
35/** Where a marketplace entry's own manifest is, when its source is a folder of the marketplace. */
36export const manifestPathOf = (marketDir: string, entry: MarketEntry): string | undefined => {
37 if (typeof entry.source !== 'string') return undefined
38 const rel = entry.source.replace(/^\.\/?/, '').replace(/\/$/, '')
39 return `${marketDir}${rel ? `/${rel}` : ''}/.claude-plugin/plugin.json`
40}
41
42/** Installed plugins whose marketplace offers a later version, by the versions read for them. */
43export const updatesOf = (installed: readonly Installed[], latest: Record<string, string | undefined>): Update[] =>
44 installed.flatMap(p => {
45 const to = latest[p.id]
46 return to && isNewer(to, p.version) ? [{ id: p.id, from: p.version, to }] : []
47 })
48hooks/classify.ts 623 lines1import type { RowCode, Folded, Health, Stream } from '../types'
2
3/** A stream id is its name as a slug, so a collapsed row can show it without a lookup. */
4export const slug = (name: string): string =>
5 name
6 .toLowerCase()
7 .replace(/[^a-z0-9]+/g, '-')
8 .replace(/^-+|-+$/g, '')
9 .slice(0, 32) || 'stream'
10
11/** `#name rest` at the start of a prompt picks the stream by hand; the tag is stripped. */
12export const parseTag = (text: string): { name: string; rest: string } | undefined => {
13 const m = /^#([A-Za-z0-9][\w-]{0,31})\s+([\s\S]*)$/.exec(text.trim())
14 return m?.[1] && m[2] !== undefined ? { name: m[1], rest: m[2] } : undefined
15}
16
17/** The `#tag` being typed at the start of a draft, up to the cursor: `#bil` gives `bil`; anything else none. */
18export const partialTag = (text: string, cursor: number): string | undefined => /^\s*#([\w-]*)$/.exec(text.slice(0, cursor))?.[1]
19
20/** Streams a partial tag could complete to, live ones first, most recently active first; at most `limit`. */
21export const tagMatches = (streams: readonly { id: string; lastAt: number; archived?: boolean }[], partial: string, limit = 6): string[] =>
22 streams
23 .filter(s => s.id.startsWith(partial.toLowerCase()))
24 .sort((a, b) => Number(!!a.archived) - Number(!!b.archived) || b.lastAt - a.lastAt)
25 .slice(0, limit)
26 .map(s => s.id)
27
28/** The draft with the partial tag before the cursor replaced by `#id `, and the cursor after it. */
29export const completeTag = (text: string, cursor: number, id: string): { text: string; cursor: number } => {
30 const head = text.slice(0, cursor).replace(/#[\w-]*$/, `#${id} `)
31 const tail = text.slice(cursor).replace(/^\s+/, '')
32 return { text: head + tail, cursor: head.length }
33}
34
35/**
36 * The key a transcript row is filed under. The transcript draws a row under its uuid with the last group
37 * zeroed (stored 61ec327a-403c-4478-84b7-4b8d64663b0d, drawn 61ec327a-403c-4478-84b7-000000000000), so a
38 * uuid keys by its first four groups; anything else (a tool_use_id, a text key) keys as it is.
39 */
40export const rowKey = (id: string): string =>
41 /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i.test(id) ? id.slice(0, 23).toLowerCase() : id
42
43/** Prompts too short to carry a topic ("yes", "continue") stay on the current stream. */
44export const isFollowUp = (text: string): boolean => {
45 const t = text.trim().toLowerCase()
46 return (
47 t.length < 4 ||
48 t.startsWith('/') ||
49 /^(y|yes|yep|no|nope|ok|okay|sure|go|go ahead|continue|proceed|do it|ship it|thanks|thank you|lgtm|approved?)[.!]*$/.test(t)
50 )
51}
52
53/** A loop fires the same text each tick: key it loosely so every tick lands together. */
54export const loopKey = (text: string): string => text.trim().toLowerCase().replace(/\s+/g, ' ').slice(0, 200)
55
56/**
57 * Key for an assistant text block, for when a row's requestId is not its uuid. Letters and digits only, so
58 * the text as stored and as the transcript draws it (markdown, spacing, case) key alike.
59 */
60export const textKey = (text: string): string => {
61 const t = text.toLowerCase().replace(/[^a-z0-9]+/g, '').slice(0, 80)
62 let h = 0
63 for (let i = 0; i < t.length; i++) h = (h * 31 + t.charCodeAt(i)) | 0
64 return `t:${h.toString(36)}`
65}
66
67export const SYSTEM = `You sort a developer's messages to a coding agent into workstreams.
68A workstream is one coherent goal (a feature, a bug, an investigation, a question), not a topic area.
69Reply with JSON only, one of:
70{"stream":"<existing id>"}
71{"new":"<2-4 word name>","summary":"<one-line summary>"}
72Use an existing stream when the message continues, refines or asks about that stream's own goal.
73A stream marked finished is done: use it only when the message plainly reopens that same goal
74(the same bug, the same files, "that didn't work"). A new request in the same area is a new stream:
75after a finished "repo setup" stream, "add badges to the README" is new work, not more repo setup.
76When unsure between a finished stream and a new one, start a new one.`
77
78/** A stream idle this long, or one whose work is done, is offered to the classifier as finished. */
79export const FINISHED_MS = 30 * 60_000
80
81export const buildPrompt = (
82 streams: readonly Stream[],
83 current: string,
84 text: string,
85 now = 0,
86 health: Readonly<Record<string, Health>> = {},
87): string => {
88 const stateOf = (s: Stream): string => {
89 const verdict = health[s.id] ?? 'idle'
90 if (verdict === 'running' || verdict === 'stalled') return 'working now'
91 return now - s.lastAt > FINISHED_MS || verdict === 'done' ? `finished, last active ${ago(now - s.lastAt)} ago` : `active ${ago(now - s.lastAt)} ago`
92 }
93 const list = streams.length
94 ? streams
95 .map(s => `- id: ${s.id}${s.id === current ? ' (current)' : ''}\n name: ${s.name}\n goal: ${s.summary}\n state: ${stateOf(s)}`)
96 .join('\n')
97 : '(none yet)'
98 return `Workstreams:\n${list}\n\nNew message:\n"""\n${text.slice(0, 2000)}\n"""`
99}
100
101export type Verdict = { kind: 'existing'; id: string; summary?: string } | { kind: 'new'; name: string; summary: string }
102
103/** Reads the model's reply; anything unreadable or naming no known stream is undefined. */
104export const parseVerdict = (reply: string, streams: readonly Stream[]): Verdict | undefined => {
105 const m = /\{[\s\S]*\}/.exec(reply)
106 if (!m) return undefined
107 let v: { stream?: unknown; new?: unknown; summary?: unknown }
108 try {
109 v = JSON.parse(m[0])
110 } catch {
111 return undefined
112 }
113 const summary = typeof v.summary === 'string' ? v.summary.slice(0, 160) : undefined
114 if (typeof v.stream === 'string' && streams.some(s => s.id === v.stream)) {
115 return { kind: 'existing', id: v.stream, summary }
116 }
117 if (typeof v.new === 'string' && v.new.trim()) {
118 return { kind: 'new', name: v.new.trim().slice(0, 40), summary: summary ?? '' }
119 }
120 return undefined
121}
122
123/** Names a new stream from the prompt itself, when the model cannot be asked. */
124export const fallbackName = (text: string): string => text.trim().split(/\s+/).slice(0, 4).join(' ').slice(0, 40) || 'misc'
125
126export const ago = (ms: number): string => {
127 const s = Math.max(0, Math.round(ms / 1000))
128 if (s < 60) return `${s}s`
129 if (s < 3600) return `${Math.round(s / 60)}m`
130 if (s < 86400) return `${Math.round(s / 3600)}h`
131 return `${Math.round(s / 86400)}d`
132}
133
134/** Terminal colour codes and other control characters: a pane line holding one does not paint. */
135const CONTROL = /\x1b\[[0-9;?]*[ -/]*[@-~]|\x1b[@-_]|[\x00-\x08\x0b-\x1f\x7f]/g
136
137export const oneLine = (text: string, n: number): string => {
138 const t = text.replace(CONTROL, '').replace(/\s+/g, ' ').trim()
139 return t.length > n ? `${t.slice(0, n - 1)}…` : t
140}
141
142export const STALL_MS = 120_000
143export const STALL_IN_TOOL_MS = 600_000
144export const DONE_MS = 600_000
145
146export type Pulse = {
147 now: number
148 lastAt: number
149 isTurnOn: boolean
150 liveAgents: number
151 inflight: number
152 outcome?: 'answer' | 'aborted' | 'refusal' | 'error'
153}
154
155/** One heartbeat's verdict. A running tool (a long build) earns a longer silence before it counts as stalled. */
156export const healthOf = (p: Pulse): Health => {
157 const quiet = p.now - p.lastAt
158 if (p.isTurnOn || p.liveAgents > 0 || p.inflight > 0) {
159 return quiet > (p.inflight > 0 ? STALL_IN_TOOL_MS : STALL_MS) ? 'stalled' : 'running'
160 }
161 if (p.outcome === 'error' || p.outcome === 'refusal') return 'error'
162 return quiet < DONE_MS ? 'done' : 'idle'
163}
164
165/** Pill backgrounds: dark enough that the default text reads on them. Yellow runs, green is done, red failed. */
166export const HEALTH_COLOR: Record<Health, string> = {
167 running: '#9a6700',
168 stalled: '#bc4c00',
169 done: '#1a7f37',
170 error: '#cf222e',
171 idle: '#57606a',
172}
173
174/** The same verdicts as text on the terminal's own background. */
175export const HEALTH_TEXT: Record<Health, string> = {
176 running: '#f2cc60',
177 stalled: '#ffa657',
178 done: '#7ee787',
179 error: '#ff7b72',
180 idle: '#8b949e',
181}
182
183export const HEALTH_GLYPH: Record<Health, string> = { running: '●', stalled: '◌', done: '✓', error: '✗', idle: '○' }
184
185export const REPLY_SYSTEM = `A coding agent was working on a task when the developer sent extra messages mid-way.
186Given one piece of the agent's reply, say which of the listed threads it addresses.
187Reply with the thread id only.`
188
189export const buildReplyPrompt = (turn: Folded, folded: readonly Folded[], text: string): string =>
190 `Threads:\n${[turn, ...folded].map(f => `- ${f.streamId}: ${oneLine(f.text, 200)}`).join('\n')}\n\nReply piece:\n"""\n${text.slice(0, 1500)}\n"""`
191
192/** The thread a reply piece answers; the turn's own when the model names none of them. */
193export const pickReplyStream = (reply: string, turn: Folded, folded: readonly Folded[]): string => {
194 const ids = [turn, ...folded].map(f => f.streamId)
195 const said = reply.trim().toLowerCase()
196 return ids.find(id => said === id) ?? ids.find(id => said.includes(id)) ?? turn.streamId
197}
198
199export const REPLIES_SYSTEM = `A coding agent was working on a task when the developer sent extra messages mid-way.
200Given the numbered pieces of the agent's reply, say which listed thread each piece addresses.
201Reply with a JSON array of thread ids, one per piece, in order.`
202
203export const buildRepliesPrompt = (turn: Folded, folded: readonly Folded[], texts: readonly string[]): string =>
204 `Threads:\n${[turn, ...folded].map(f => `- ${f.streamId}: ${oneLine(f.text, 200)}`).join('\n')}\n\nReply pieces:\n${texts
205 .map((t, i) => `${i + 1}. ${oneLine(t, 400)}`)
206 .join('\n')}`
207
208/** One thread id per piece; any piece the model got wrong or left out stays with the turn. */
209export const pickReplyStreams = (reply: string, turn: Folded, folded: readonly Folded[], n: number): string[] => {
210 const ids = new Set([turn, ...folded].map(f => f.streamId))
211 let said: unknown
212 try {
213 said = JSON.parse(/\[[\s\S]*\]/.exec(reply)?.[0] ?? '[]')
214 } catch {
215 said = []
216 }
217 const list = Array.isArray(said) ? said : []
218 return Array.from({ length: n }, (_, i) => {
219 const v = list[i]
220 return typeof v === 'string' && ids.has(v) ? v : turn.streamId
221 })
222}
223
224/** A turn of the transcript: the prompt that began it, prompts sent into it, and what the agent did. */
225export type HistoryTurn = { prompt: HistoryItem & { kind: 'prompt' }; items: HistoryItem[] }
226
227/** Splits the transcript at each typed prompt; anything before the first one is dropped. */
228export const turnsOf = (items: readonly HistoryItem[]): HistoryTurn[] => {
229 const turns: HistoryTurn[] = []
230 for (const item of items) {
231 if (item.kind === 'prompt' && !item.isFolded) turns.push({ prompt: item, items: [] })
232 else turns.at(-1)?.items.push(item)
233 }
234 return turns
235}
236
237/** One line of a transcript file, as far as filing it into streams needs. */
238export type TranscriptLine = {
239 type?: string
240 uuid?: string
241 isMeta?: boolean
242 isSidechain?: boolean
243 message?: { role?: string; content?: unknown }
244 attachment?: { type?: string; prompt?: unknown; commandMode?: string }
245}
246
247export type HistoryItem =
248 | { kind: 'prompt'; uuid: string; text: string; isFolded: boolean; at?: number }
249 | { kind: 'reply'; uuid: string; text: string; at?: number }
250 | { kind: 'tool'; uuid: string; id: string; name: string; input: unknown; at?: number }
251
252/** A message's text: a string as it is, a list of blocks by its text blocks (a prompt with an image is one). */
253export const textOf = (content: unknown): string =>
254 typeof content === 'string'
255 ? content
256 : Array.isArray(content)
257 ? content
258 .filter(b => b?.type === 'text' && typeof b.text === 'string')
259 .map(b => b.text as string)
260 .join('\n')
261 : ''
262
263/** The main conversation's prompts (typed or sent mid-turn), replies and tool calls, in order. */
264export const readTranscript = (jsonl: string): HistoryItem[] => {
265 const items: HistoryItem[] = []
266 for (const line of jsonl.split('\n')) {
267 if (!line.trim()) continue
268 let d: TranscriptLine
269 try {
270 d = JSON.parse(line)
271 } catch {
272 continue
273 }
274 if (d.isSidechain || d.isMeta || !d.uuid) continue
275 // A queued task notification is the engine's, not the person's: only queued prompts count.
276 const queued = d.type === 'attachment' && d.attachment?.type === 'queued_command' && d.attachment.commandMode === 'prompt'
277 if (queued) {
278 const text = textOf(d.attachment?.prompt).trim()
279 if (text) items.push({ kind: 'prompt', uuid: d.uuid, text, isFolded: true })
280 continue
281 }
282 const content = d.message?.content
283 if (d.type === 'user') {
284 const text = textOf(content)
285 if (text.trim() && !text.trim().startsWith('<')) items.push({ kind: 'prompt', uuid: d.uuid, text: text.trim(), isFolded: false })
286 } else if (d.type === 'assistant' && Array.isArray(content)) {
287 for (const b of content) {
288 if (b?.type === 'text' && typeof b.text === 'string' && b.text.trim()) items.push({ kind: 'reply', uuid: d.uuid, text: b.text })
289 if (b?.type === 'tool_use' && typeof b.id === 'string') items.push({ kind: 'tool', uuid: d.uuid, id: b.id, name: String(b.name), input: b.input })
290 }
291 }
292 }
293 return items
294}
295
296/** Pastels that read on a dark terminal and stay apart from each other. */
297export const PASTELS = [
298 '#a5d8ff', // sky
299 '#b2f2bb', // mint
300 '#ffd8a8', // peach
301 '#d0bfff', // lavender
302 '#fcc2d7', // pink
303 '#ffec99', // butter
304 '#99e9f2', // aqua
305 '#ffc9c9', // rose
306 '#c0eb75', // lime
307 '#bac8ff', // periwinkle
308]
309
310/** A new stream takes the first pastel no live stream wears, so neighbours never share a colour. */
311export const nextPastel = (taken: readonly (string | undefined)[]): string =>
312 PASTELS.find(c => !taken.includes(c)) ?? PASTELS[taken.length % PASTELS.length] ?? '#a5d8ff'
313
314/** The colour of a stream made before colours were kept: stable, from its id. */
315export const pastelOf = (id: string): string => {
316 let h = 0
317 for (let i = 0; i < id.length; i++) h = (h * 31 + id.charCodeAt(i)) | 0
318 return PASTELS[Math.abs(h) % PASTELS.length] ?? '#a5d8ff'
319}
320
321/** How much of a stream's activity the pane shows: its share of the pane, the last 10 rows, the last one, or none. */
322export type Fold = 'all' | '10' | '1' | 'none'
323export const NEXT_FOLD: Record<Fold, Fold> = { all: '10', '10': '1', '1': 'none', none: 'all' }
324export const FOLD_LABEL: Record<Fold, string> = { all: '▾ all', '10': '▾ 10', '1': '▾ 1', none: '▸' }
325
326export const clockOf = (ms: number): string => {
327 const s = Math.max(0, Math.floor(ms / 1000))
328 return `${Math.floor(s / 60)}:${String(s % 60).padStart(2, '0')}`
329}
330
331/** Solid badges: black on bright yellow while anything runs, so movement is impossible to miss. */
332export type BadgeKind = Health | 'loop'
333export const BADGE_BG: Record<BadgeKind, string> = {
334 running: '#ffd33d',
335 loop: '#ffd33d',
336 stalled: '#ff9500',
337 done: '#2ea043',
338 error: '#da3633',
339 idle: '#30363d',
340}
341export const BADGE_FG: Record<BadgeKind, string> = {
342 running: '#000000',
343 loop: '#000000',
344 stalled: '#000000',
345 done: '#ffffff',
346 error: '#ffffff',
347 idle: '#c9d1d9',
348}
349
350/** Runs `fn` over `items` with at most `limit` at once, results in input order. Model calls wait on the network, so a few in flight file a long history several times faster. */
351export async function inParallel<T, R>(items: readonly T[], limit: number, fn: (item: T, index: number) => Promise<R>): Promise<R[]> {
352 const results: R[] = new Array(items.length)
353 let next = 0
354 const lane = async () => {
355 for (let i = next++; i < items.length; i = next++) results[i] = await fn(items[i] as T, i)
356 }
357 await Promise.all(Array.from({ length: Math.min(limit, items.length) }, lane))
358 return results
359}
360
361export const BATCH_SYSTEM = `You sort a developer's messages to a coding agent into workstreams.
362A workstream is one coherent goal (a feature, a bug, an investigation, a question).
363You get the existing workstreams and a numbered run of messages, oldest first.
364Reply with a JSON array, one string per message, in order: the id of an existing workstream,
365or a 2-4 word name for a new one. Use the same new name for every message about the same goal,
366and give a message that continues the one before it ("yes", "now also do X") that message's entry.`
367
368export const buildBatchPrompt = (streams: readonly Stream[], texts: readonly string[]): string => {
369 const list = streams.length ? streams.map(s => `- ${s.id}: ${oneLine(s.summary || s.name, 120)}`).join('\n') : '(none yet)'
370 return `Workstreams:\n${list}\n\nMessages:\n${texts.map((t, i) => `${i + 1}. ${oneLine(t, 300)}`).join('\n')}`
371}
372
373/** One label per message: an existing id or a proposed name; '' where the model gave nothing usable. */
374export const parseBatch = (reply: string, n: number): string[] => {
375 let said: unknown = []
376 try {
377 said = JSON.parse(/\[[\s\S]*\]/.exec(reply)?.[0] ?? '[]')
378 } catch {
379 said = []
380 }
381 const list = Array.isArray(said) ? said : []
382 return Array.from({ length: n }, (_, i) => (typeof list[i] === 'string' ? (list[i] as string).trim().slice(0, 40) : ''))
383}
384
385export const MERGE_SYSTEM = `Workstream names were proposed separately for different parts of one long session.
386Merge names that mean the same goal. Map a name to an existing workstream id when it is that workstream.
387Reply with a JSON object mapping every proposed name to its final name or existing id.`
388
389export type Proposal = { name: string; count: number; samples: string[] }
390
391export const buildMergePrompt = (streams: readonly Stream[], proposals: readonly Proposal[]): string =>
392 `Existing workstreams:\n${streams.length ? streams.map(s => `- ${s.id}: ${oneLine(s.summary || s.name, 100)}`).join('\n') : '(none)'}\n\nProposed names:\n${proposals
393 .map(p => `- "${p.name}" (${p.count} messages), e.g. ${p.samples.map(s => `"${oneLine(s, 80)}"`).join('; ')}`)
394 .join('\n')}`
395
396/** Every proposed name to its final one; a name the model left out keeps itself. */
397export const parseMerge = (reply: string, names: readonly string[]): Record<string, string> => {
398 let said: unknown = {}
399 try {
400 said = JSON.parse(/\{[\s\S]*\}/.exec(reply)?.[0] ?? '{}')
401 } catch {
402 said = {}
403 }
404 const map = (said && typeof said === 'object' ? said : {}) as Record<string, unknown>
405 return Object.fromEntries(names.map(n => [n, typeof map[n] === 'string' && (map[n] as string).trim() ? (map[n] as string).trim().slice(0, 40) : n]))
406}
407
408/** Longest source a row keeps for the full chat style; the pane is a glance, not the file. */
409export const CODE_LIMIT = 4000
410
411/** `text` cut to whole lines within `limit` characters, with a marker when anything was left out. */
412const clip = (text: string, limit: number): string => {
413 if (text.length <= limit) return text
414 const cut = text.slice(0, limit)
415 const kept = cut.slice(0, Math.max(0, cut.lastIndexOf('\n')))
416 return `${kept}\n… ${text.split('\n').length - kept.split('\n').length} more lines`
417}
418
419/** One unified-diff hunk replacing `before` with `after`, as the session draws an Edit. */
420const hunk = (before: string, after: string): string => {
421 const old = before.split('\n')
422 const now = after.split('\n')
423 return [`@@ -1,${old.length} +1,${now.length} @@`, ...old.map(l => `-${l}`), ...now.map(l => `+${l}`)].join('\n')
424}
425
426/** Strips control characters a Code element refuses; tab and newline stay. */
427const printable = (text: string): string => text.replace(/\r\n?/g, '\n').replace(/[\u0000-\u0008\u000b-\u001f\u007f]/g, '')
428
429/**
430 * What the full chat style draws under a tool call: the command for Bash, a diff for Edit and MultiEdit,
431 * the file for Write; none for tools whose one line says it all (Read, Grep, Glob, ...).
432 */
433export const codeOf = (tool: string, input: unknown): RowCode | undefined => {
434 const i = (input ?? {}) as Record<string, unknown>
435 const str = (k: string): string | undefined => (typeof i[k] === 'string' ? (i[k] as string) : undefined)
436 const path = str('file_path')
437 if (tool === 'Bash' && str('command')) return { source: clip(printable(str('command')!), CODE_LIMIT), language: 'bash' }
438 if (tool === 'Write' && str('content') !== undefined) return { source: clip(printable(str('content')!), CODE_LIMIT), ...(path ? { path } : {}) }
439 const edits: { old_string?: unknown; new_string?: unknown }[] =
440 tool === 'Edit' ? [i] : tool === 'MultiEdit' && Array.isArray(i.edits) ? (i.edits as { old_string?: unknown }[]) : []
441 const hunks = edits
442 .filter(e => typeof e.old_string === 'string' && typeof e.new_string === 'string')
443 .map(e => hunk(printable(e.old_string as string), printable(e.new_string as string)))
444 if (hunks.length === 0) return undefined
445 const diff = hunks.join('\n')
446 // A diff cut mid-hunk no longer parses, so an oversized one is drawn as its first hunks only.
447 if (diff.length <= CODE_LIMIT) return { source: diff, format: 'diff', ...(path ? { path } : {}) }
448 const fit: string[] = []
449 for (const h of hunks) if ([...fit, h].join('\n').length <= CODE_LIMIT) fit.push(h)
450 return fit.length ? { source: fit.join('\n'), format: 'diff', ...(path ? { path } : {}) } : { source: clip(hunks[0]!, CODE_LIMIT), ...(path ? { path } : {}) }
451}
452
453/** A tool call as the session titles it: `Edit(src/a.ts)`, `Bash(npm test)`; the raw input when no field names it. */
454export const toolLine = (tool: string, input: unknown): string => {
455 const i = (input ?? {}) as Record<string, unknown>
456 const arg = ['file_path', 'notebook_path', 'command', 'pattern', 'url', 'query', 'path', 'skill', 'prompt']
457 .map(k => i[k])
458 .find((v): v is string => typeof v === 'string' && v.trim() !== '')
459 return arg !== undefined ? `${tool}(${oneLine(arg, 100)})` : `${tool} ${oneLine(JSON.stringify(input ?? {}), 100)}`
460}
461
462/** What the status card says of one stream: a state word, its colour kind, and one line of detail. */
463export type StatusKind = BadgeKind | 'waiting'
464export type StatusLine = { id: string; area: string; kind?: StatusKind; state: string; detail: string }
465export type StatusInput = {
466 stream: Pick<Stream, 'id' | 'name' | 'summary' | 'lastAt'>
467 health: Health
468 loop?: { kind: 'wakeup' | 'cron'; nextAt: number; label: string }
469 running: { description: string; last: string; tools: number }[]
470 /** The stream's last prompt or reply: a reply ending in a question is waiting on the person. */
471 lastSaid?: { kind: string; text: string }
472 now: number
473}
474
475const STATE_WORD: Record<StatusKind, string> = {
476 running: 'RUNNING',
477 loop: 'LOOP',
478 waiting: 'WAITING FOR YOU',
479 error: 'ERROR',
480 stalled: 'STALLED',
481 done: 'DONE',
482 idle: 'IDLE',
483}
484const STATE_RANK: Record<StatusKind, number> = { running: 0, loop: 1, waiting: 2, error: 3, stalled: 4, done: 5, idle: 6 }
485
486/** The question a reply ends on, as its last sentence; undefined when it does not end on one. */
487export const questionOf = (text: string): string | undefined => {
488 const t = text.trim().replace(/[*_`\s]+$/, '')
489 if (!t.endsWith('?')) return undefined
490 const last = t.split(/(?<=[.!?:])\s+|\n+/).filter(Boolean).at(-1) ?? t
491 return oneLine(last, 300)
492}
493
494/** One stream's row on the status card: what it is doing, or what it last left the person with. */
495export function statusOf(x: StatusInput): StatusLine {
496 const { stream: s, now } = x
497 const line = (kind: StatusKind, detail: string, state = STATE_WORD[kind]): StatusLine => ({ id: s.id, area: s.name, kind, state, detail })
498 if (x.health === 'running') {
499 const top = x.running[0]
500 if (!top) return line('running', `main turn · ${s.summary}`)
501 const more = x.running.length > 1 ? `${x.running.length} agents · ` : ''
502 return line('running', `${more}${top.description}: ${top.tools ? `${top.last} (${top.tools} tools)` : 'starting up'}`)
503 }
504 if (x.loop) {
505 const when = x.loop.kind === 'cron' ? x.loop.label : `next tick in ${clockOf(Math.max(0, x.loop.nextAt - now))}`
506 return line('loop', `${when} · ${s.summary}`)
507 }
508 const question = x.lastSaid?.kind === 'reply' ? questionOf(x.lastSaid.text) : undefined
509 if (question && x.health !== 'error') return line('waiting', question)
510 return line(x.health, `${s.summary || '—'} · ${ago(now - s.lastAt)} ago`)
511}
512
513/** Status rows in the order that needs the person: running, looping, waiting, failed, then the finished. */
514export const sortStatus = (lines: StatusLine[], lastAt: Record<string, number>): StatusLine[] =>
515 [...lines].sort((a, b) => STATE_RANK[a.kind ?? 'idle'] - STATE_RANK[b.kind ?? 'idle'] || (lastAt[b.id] ?? 0) - (lastAt[a.id] ?? 0))
516
517/** The card's git rows from `git status --porcelain=v1 --branch`: branch against upstream, and what is uncommitted. */
518export function gitStatus(porcelain: string): StatusLine[] {
519 const [head = '', ...files] = porcelain.split('\n').filter(Boolean)
520 const m = /^## (?:No commits yet on )?([^.\s]+)(?:\.\.\.(\S+))?(?: \[(.*)\])?/.exec(head)
521 if (!m) return []
522 const [, branch = '', upstream, track = ''] = m
523 const ahead = Number(/ahead (\d+)/.exec(track)?.[1] ?? 0)
524 const behind = Number(/behind (\d+)/.exec(track)?.[1] ?? 0)
525 const sync = !upstream
526 ? 'no upstream; nothing pushed'
527 : ahead || behind
528 ? [ahead ? `${ahead} unpushed` : '', behind ? `${behind} behind ${upstream}` : ''].filter(Boolean).join(', ')
529 : `up to date with ${upstream}`
530 const rows: StatusLine[] = [{ id: 'git:branch', area: 'Git branch', state: branch, detail: sync }]
531 const paths = files.map(f => f.slice(3))
532 rows.push({
533 id: 'git:changes',
534 area: 'Uncommitted',
535 state: paths.length ? `${paths.length} file${paths.length === 1 ? '' : 's'}` : 'none',
536 detail: paths.length ? oneLine(paths.slice(0, 4).join(', ') + (paths.length > 4 ? ` +${paths.length - 4} more` : ''), 300) : 'clean',
537 })
538 return rows
539}
540
541/** A ticket id as trackers write them (TL-260, ENG-1042); the standards and encodings that look alike are not tickets. */
542const TICKET = /\b([A-Z][A-Z0-9]{1,9}-\d{1,6})\b/g
543const NOT_TICKETS = new Set(['UTF', 'SHA', 'ISO', 'GPT', 'RFC', 'HTTP', 'IPV', 'ES', 'MD', 'ECMA', 'WCAG', 'COVID', 'AES', 'RSA', 'P', 'H', 'X', 'TLS', 'SSL'])
544
545/** The ticket ids a text names, in order, each once. */
546export const ticketsIn = (text: string): string[] =>
547 [...new Set([...text.matchAll(TICKET)].map(m => m[1] ?? '').filter(id => id && !NOT_TICKETS.has(id.split('-')[0] ?? '')))]
548
549export type TicketInput = {
550 rows: readonly { kind: string; text: string; at: number; streamId: string }[]
551 agents: readonly { description: string; status: 'running' | 'done' | 'error'; last: string; tools: number; lastAt: number; endedAt?: number; streamId: string }[]
552 /** Each stream's own state on the card, so a ticket worked in a waiting stream is waiting too. */
553 streamKind: Record<string, StatusKind>
554 now: number
555}
556
557/** The sentence of a text that names the ticket, or its first sentence. */
558const sentenceAbout = (text: string, id: string): string => {
559 const sentences = text.replace(/\s+/g, ' ').split(/(?<=[.!?])\s+/)
560 return sentences.find(s => s.includes(id)) ?? sentences[0] ?? ''
561}
562
563/**
564 * One status row per ticket the person or an agent was set to work on: named in a prompt or an agent's task.
565 * Running agents on it say what they are doing and how long they have been quiet; otherwise its latest news.
566 */
567export function ticketLines(x: TicketInput): StatusLine[] {
568 const { now } = x
569 const named = [...new Set([...x.rows.filter(r => r.kind === 'prompt' || r.kind === 'loop').map(r => r.text), ...x.agents.map(a => a.description)].flatMap(ticketsIn))]
570 const lines = named.map(id => {
571 const mine = x.agents.filter(a => ticketsIn(a.description).includes(id))
572 const said = x.rows.filter(r => ticketsIn(r.text).includes(id)).sort((a, b) => a.at - b.at)
573 const lastAt = Math.max(0, ...mine.map(a => a.endedAt ?? a.lastAt), ...said.map(r => r.at))
574 const line = (kind: StatusKind, detail: string): StatusLine & { lastAt: number } => ({ id: `ticket:${id}`, area: id, kind, state: STATE_WORD[kind], detail, lastAt })
575 const running = mine.filter(a => a.status === 'running').sort((a, b) => b.lastAt - a.lastAt)
576 const top = running[0]
577 if (top) {
578 const quiet = now - top.lastAt > 60_000 ? `, ${ago(now - top.lastAt)} with no output` : ''
579 const more = running.length > 1 ? `${running.length} agents · ` : ''
580 return line('running', `${more}${top.description}: ${top.tools ? `${top.last} (${top.tools} tools${quiet})` : 'starting up'}`)
581 }
582 const ended = mine.filter(a => a.endedAt !== undefined).sort((a, b) => (b.endedAt ?? 0) - (a.endedAt ?? 0))[0]
583 const reply = said.filter(r => r.kind === 'reply').at(-1)
584 if (ended?.status === 'error' && (ended.endedAt ?? 0) >= (reply?.at ?? 0)) return line('error', `${ended.description} failed`)
585 const streams = [...new Set([...said.map(r => r.streamId), ...mine.map(a => a.streamId)])]
586 const kinds = streams.map(s => x.streamKind[s] ?? 'idle')
587 const kind = (['running', 'waiting', 'loop', 'error', 'stalled', 'done'] as const).find(k => kinds.includes(k)) ?? 'idle'
588 const latest = said.filter(r => r.kind !== 'tool').at(-1)
589 const news = latest ? oneLine(sentenceAbout(latest.text, id), 300) : ended ? `${ended.description} finished` : ''
590 return line(kind, news || '—')
591 })
592 return sortStatus(lines, Object.fromEntries(lines.map(l => [l.id, l.lastAt]))).map(({ id, area, kind, state, detail }) => ({ id, area, ...(kind ? { kind } : {}), state, detail }))
593}
594
595/** One plan limit as the status card's footer shows it. */
596export type LimitView = { label: string; percent: number; bar: string; resetsIn: string; resetsAt: string }
597
598const LIMIT_LABEL: Record<string, string> = { five_hour: '5h', seven_day: 'week', seven_day_opus: 'week opus', seven_day_sonnet: 'week sonnet', spend_limit: 'spend' }
599const WEEKDAY = ['Sun', 'Mon', 'Tue', 'Wed', 'Thu', 'Fri', 'Sat']
600
601/** How long until a time, in the two largest units: `3d 4h`, `2h 14m`, `9m`. */
602export const untilOf = (ms: number): string => {
603 const m = Math.max(0, Math.round(ms / 60_000))
604 const d = Math.floor(m / 1440)
605 const h = Math.floor((m % 1440) / 60)
606 return d ? `${d}d ${h}h` : h ? `${h}h ${m % 60}m` : `${m}m`
607}
608
609/** A plan limit for the card: its window, a ten-cell bar, the percent used, and when it resets with the weekday. */
610export function limitView(limit: { kind: string; percentUsed: number; resetsAt?: string }, now: number): LimitView {
611 const percent = Math.round(limit.percentUsed)
612 const filled = Math.max(0, Math.min(10, Math.round(percent / 10)))
613 const at = limit.resetsAt ? new Date(limit.resetsAt) : undefined
614 const ok = at !== undefined && !Number.isNaN(at.getTime())
615 return {
616 label: LIMIT_LABEL[limit.kind] ?? limit.kind.replace(/_/g, ' '),
617 percent,
618 bar: '▰'.repeat(filled) + '▱'.repeat(10 - filled),
619 resetsIn: ok ? untilOf(at.getTime() - now) : '',
620 resetsAt: ok ? `${WEEKDAY[at.getDay()]} ${String(at.getHours()).padStart(2, '0')}:${String(at.getMinutes()).padStart(2, '0')}` : '',
621 }
622}
623types/index.d.ts 124 lines1export type StreamRowKind = 'prompt' | 'reply' | 'tool' | 'agent' | 'loop' | 'notice'
2
3/** What the heartbeat makes of a stream: working, quiet too long while working, finished lately, failed, or asleep. */
4export type Health = 'running' | 'stalled' | 'done' | 'error' | 'idle'
5
6/** A prompt sent while a turn ran: it waits, filed in its own stream, for the reply that answers it. */
7export type Folded = { streamId: string; text: string }
8
9/** A subagent's run, as the pane shows it live. */
10export type AgentRun = {
11 id: string
12 streamId: string
13 description: string
14 status: 'running' | 'done' | 'error'
15 startedAt: number
16 endedAt?: number
17 lastAt: number
18 /** What it did last: a tool call or the head of its latest text. */
19 last: string
20 tools: number
21}
22
23export type Stream = {
24 id: string
25 name: string
26 summary: string
27 createdAt: number
28 lastAt: number
29 rows: number
30 agents: number
31 loops: number
32 /** The stream's own pastel, fixed at creation: its line down the transcript, its name in the pane. */
33 color?: string
34 /** Hidden from the bar and the list until restored; its history stays. */
35 archived?: boolean
36}
37
38/** A tool call's input as the full chat style draws it: highlighted source, or a unified diff. */
39export type RowCode = { source: string; language?: string; path?: string; format?: 'diff' }
40
41export type StreamRow = {
42 id: string
43 streamId: string
44 kind: StreamRowKind
45 text: string
46 agentId?: string
47 at: number
48 code?: RowCode
49 /** A tool row's call id: the key its transcript row is filed under. */
50 toolId?: string
51}
52
53/** How a stream's own view draws its rows: one line each, or as the session's transcript draws them. */
54export type ChatStyle = 'compact' | 'full'
55
56declare module 'claude-code' {
57 interface PluginState {
58 streams: {
59 streams: Stream[]
60 /** The stream the main loop is working on now; '' before the first prompt. */
61 current: string
62 /** The stream the transcript is focused on; '' shows everything. */
63 focus: string
64 /** The stream the pane shows in detail; '' shows the list. */
65 view: string
66 /** Whether a main-loop turn is running. */
67 busy: boolean
68 /** Subagent runs by id. */
69 agents: Record<string, AgentRun>
70 /** How much of each stream the pane shows. */
71 fold: Record<string, 'all' | '10' | '1' | 'none'>
72 /** When the main turn started, for its running clock. */
73 turnStartedAt: number
74 /** A clock the pane reads while anything runs, so elapsed times move. */
75 tick: number
76 /** Loops waiting to fire, per stream: a self-paced wakeup or a cron job. */
77 loops: Record<string, { kind: 'wakeup' | 'cron'; nextAt: number; label: string }>
78 /** The chat style chosen in the pane this session; '' follows the `chatStyle` setting. */
79 chatStyle: ChatStyle | ''
80 /** The `#tag` being typed at the start of the prompt box and the streams it could complete to. */
81 tagHint: { partial: string; matches: string[] } | null
82 /** Installed plugins with a newer release in their marketplace. */
83 updates: { id: string; from: string; to: string }[]
84 /** Whether an update is being installed now. */
85 updating: boolean
86 /** The stream open in the phone's accordion; '' with every card closed. */
87 mobileOpen: string
88 /** Whether the docked pane is folded away to the bar's side tab. */
89 paneCollapsed: boolean
90 /** Whether the status card is up above the prompt. */
91 statusOpen: boolean
92 /** The card's git rows, read as it opened: branch against upstream and what is uncommitted. */
93 statusGit: { id: string; area: string; state: string; detail: string }[]
94 /** Whether the pane lists archived streams too. */
95 showArchived: boolean
96 /** Whether this session's transcript has been filed into streams (once per session). */
97 historyFiled: boolean
98 /** The row-key scheme this session's rows were filed under; an older one means file them again. */
99 keyVersion: number
100 /** A history import under way: what it is doing and how far it has got; total 0 when none runs. */
101 importProgress: { label: string; done: number; total: number }
102 /** This session's transcript file, as the prompt hook names it. */
103 transcript: string
104 /** Prompts sent into the running turn, still waiting for the reply that answers them. */
105 folded: Folded[]
106 rows: StreamRow[]
107 agentStream: Record<string, string>
108 /** Subagents still running, by id, to the stream they work for. */
109 live: Record<string, string>
110 /** Tool calls running now, per stream. */
111 inflight: Record<string, number>
112 /** How each stream's last main-loop turn ended. */
113 outcome: Record<string, 'answer' | 'aborted' | 'refusal' | 'error'>
114 /** The heartbeat's verdict per stream. */
115 health: Record<string, Health>
116 loopStream: Record<string, string>
117 /** Which stream each transcript row (message uuid, tool_use_id, text key) belongs to. */
118 rowStream: StateFamily<string>
119 /** Each stream's colour by stream id, so a transcript row reads its own and redraws on nothing else. */
120 streamColor: StateFamily<string>
121 }
122 }
123}
124