SLOPSHOPPER

streams

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…

newpanebandrowsguardcommand
★ 2v0.3.4SEE LICENSE IN LICENSEupdated 2026-10-09macleodlabs-ai/claudeflow/plugins/streams
A shopper browsing a rack in a slop shop
Preview · a replayed session in a sandbox
claude · ~/work/app · streams
│ ┃ Streams ✕ › fix the failing auth test and add an audit log call │ ┃ 1 streams · showing all h: ⇥ hide │ ┃ collapse all expand all ⏺ Read(src/auth.ts) │ ┃ ⎿ Read 6 lines │ ┃ ✓ fix the failing auth ▾ all ✕ ⏺ Update(src/auth.ts) │ ┃ DONE · 0 rows · 0 agents · 0s ago ⎿ Added 2 lines, removed 1 line │ ┃ fix the failing auth test and add an audit … ⏺ Bash(bun test) │ ┃ ✓ DONE main turn ⎿ 3 pass, 1 fail │ │ ● Done. refresh now rejects expired claims and logs an audit event. │ │ ✻ Worked for 42s · done 4:20 PM │ │ › /streams │ ⎿ streams: Streams navigator opened. Attached: terminal. Pane draw │ │ 0: all ◉ 1: 1 ✓ fix-the-failing-auth t: status s: ≡ ────────────────────────────────────────────────────────────────────────────────────────────────────────────────────── › ? for shortcuts ⚠ streams: stream fix-the-failing-auth

Draws

Band
0: all ◉ 1: 1 ✓ fix-the-failing-auth t: status s: ≡
Pane · Streams
1 streams · showing all h: ⇥ hide collapse all expand all ✓ fix the failing auth ▾ all ✕ DONE · 0 rows · 0 agents · 0s ago fix the failing auth test and add an audit log call ✓ DONE main turn
README

claudeflow — one Claude Code session, many threads of work, untangled into live, colour-coded streams

<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.

The streams pane beside a transcript: coloured stripes per stream, live agent status, loop countdowns and the pill bar

✨ What you get

<table> <tr> <td width="33%" valign="top">

🧭 Navigator pane

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">

💊 Pill bar

One pill per stream above the prompt, coloured by health. Number keys jump between them; 0 shows everything.

</td> <td width="33%" valign="top">

🎨 Stream stripes

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">

🧠 Semantic routing

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">

🗂️ History, in parallel

Existing sessions are filed into streams in the background, with batched, parallel classification and a live progress line.

</td> <td valign="top">

📦 Archive & fold

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>

📋 Status card

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.

The status card above the prompt: git branch and uncommitted files, then each stream as running, waiting for you, done or idle, with what it is doing

  • What needs you comes first: running work, then loops, then streams waiting for you (their last reply ended on a question, shown as the detail), then failures, then finished work.
  • Git rows on top: the branch, whether anything is unpushed, and which files are uncommitted.
  • Tickets: any ticket id you name in a prompt or an agent's task (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.
  • Plan limits at the bottom: each window (5-hour, week) as a bar and percent, green, then yellow from 50%, red from 80%, with how long until it resets and the weekday and time it does.
  • Click a stream's name to open it in the pane. ✕ close (top right) or your next prompt hides the card.
  • Scroll a long card with the wheel, or press <kbd>ctrl</kbd>+<kbd>x</kbd> <kbd>tab</kbd> to move focus to it: then <kbd>↑</kbd> <kbd>↓</kbd> scroll and <kbd>q</kbd> closes it, without touching Claude's turn.
  • On your phone the card opens at the top of the streams accordion (its status button, or type status) and scrolls by touch.

📱 On your phone

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.

  • A stream waiting on you shows Claude's question with a yes button: one tap answers it, filed in that stream.
  • Tap a card to open it: its live agents, then its chat with markdown, highlighted commands and coloured diffs. Tap again to close it.
  • ⬆ update appears here too, so you can update without going back to the Mac.

🔄 Updates without a restart

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.

Status at a glance

ColourStatusMeaning
🟨RUNNINGA turn, subagent or loop is working now
🟩DONEFinished cleanly
🟦WAITING FOR YOUOn the status card: the stream's last reply asked you something
🟥ERRORA subagent or turn failed
🟧STALLEDNo activity for longer than expected
↻LOOPA /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

🚀 Install

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

Update or remove

ActionCommand
Updateclaude plugin update streams@claudeflow
Uninstallclaude plugin uninstall streams@claudeflow

🎛 Use

Commands

CommandWhat 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
/streamsOpen the navigator pane
/streams statusSame as typing status
/streams updateCheck 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 offShow 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
/streamList streams with their summaries
/streams importList this project's past sessions
/streams import <id>File a past session into streams (re-importing replaces, never duplicates)

In the pane and bar

ToDo
See the status of all workPress 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 streamClick its name in the pane, or press its pill
Show everything← all streams, or the all pill
File a prompt by handStart 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 everythingcollapse 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 viewThe 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

Keyboard

KeysAction
<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.


⚙️ Configure

Settings live in /config under streams, or in settings.json:

{
  "pluginConfigs": {
    "streams@claudeflow": {
      "options": { "chatStyle": "full", "diagnostics": false }
    }
  }
}
SettingDefaultDescription
chatStylefullHow 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.
diagnosticsfalseWrites 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.

Model use

EventHaiku requests
New prompt1
Follow-up (yes, continue), slash command, #tagged prompt0
Filing history1 per 25 prompts, plus 1 merge pass, plus 1 per turn that had a mid-turn prompt

🩺 Troubleshooting

SymptomFix
No pill barIt appears once the first prompt is sorted. Check /plugin lists streams as enabled.
Pane doesn't openThe terminal is under 144 columns or not fullscreen. Type /streams.
Dim streams: … line in the transcriptClaude Code is reporting a failed hook; the line names it. Include it in an issue.
Older rows have no stripeHistory is still filing; watch the progress line at the top of the pane.

🤝 Work with Macleod Labs

<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

Source 4 files
hooks/register.tsx 1837 lines
1import { 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              ) : null
hooks/updates.ts 48 lines
1/** 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  })
48
hooks/classify.ts 623 lines
1import 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}
623
types/index.d.ts 124 lines
1export 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