Agent teams, swarm coordination, Monitor streams, and worktree isolation — wraps 4 swarm_* + 8 agent_* MCP tools (12 total) plus 6 topologies (hierarchical /…

Agent teams, swarm coordination, Monitor streams, and worktree isolation.
/plugin marketplace add ruvnet/ruflo
/plugin install ruflo-swarm@ruflo
Monitor("npx @claude-flow/cli@latest swarm watch --stream")ruflo-core plugin (provides MCP server)@claude-flow/cli v3.6 major+minor.bash plugins/ruflo-swarm/scripts/smoke.sh is the contract.| Family | Count | Tools |
|---|---|---|
swarm_* | 4 | swarm_init, swarm_status, swarm_shutdown, swarm_health |
agent_* | 8 | agent_spawn, agent_execute, agent_terminate, agent_status, agent_list, agent_pool, agent_health, agent_update |
Sources: v3/@claude-flow/cli/src/mcp-tools/swarm-tools.ts:71, 145, 208, 270 and agent-tools.ts:182, 287, 319, 356, 395, 451, 573, 651.
This plugin pairs with Claude Code's native multi-agent tools (no MCP needed):
| Tool | Purpose |
|---|---|
Task | Spawn a sub-agent (use name: for addressability + run_in_background: true for parallel execution) |
SendMessage | Inter-agent comms (named agents only) |
TaskCreate / TaskList / TaskGet / TaskUpdate / TaskOutput / TaskStop | Shared task tracker for swarm pipelines |
Monitor | Live-stream events from a long-running process (persistent: true) — primary wake signal for /loop |
EnterWorktree / ExitWorktree | Git worktree isolation per agent |
For coding swarms, the canonical defaults that prevent agent drift:
| Setting | Value | Rationale |
|---|---|---|
topology | hierarchical | Coordinator catches divergence |
maxAgents | 6–8 | Smaller team = less drift |
strategy | specialized | Clear roles, no overlap |
consensus | raft | Leader maintains authoritative state |
memory | hybrid | SQLite + AgentDB for both fast + durable |
For 10+ agent teams, use hierarchical-mesh (queen + peer communication).
This plugin owns the swarm-state AgentDB namespace (kebab-case, follows the convention from ruflo-agentdb ADR-0001 §"Namespace convention"). Reserved namespaces (pattern, claude-memories, default) MUST NOT be shadowed.
swarm-state indexes active swarms, agent assignments, and topology snapshots. Accessed via memory_* (namespace-routed).
With Claude Code function hooks on (CLAUDE_CODE_ENABLE_FUNCTION_HOOKS=1, or the engine's own rollout), the plugin's hooks module (hooks/register.ts) draws a pane beside the transcript. Without function hooks nothing loads, and the commands, skills and agents above behave as they always have.
.claude-flow/agents/store.json, the hive queen, Claude Code's own subagents and teammates (from agent.spawn and $.agent.list()), and the main loop. The colour shows the state (idle, working, blocked, done, failed) and a word is always drawn beside it. A tile turns blue while its loop reads files and yellow while it writes. The leader is starred.$.session.usage(), and the router's last pick seen this session, with its score. A pick under routeThreshold is shown as below threshold. A figure that is not on disk, or that the engine did not give, is shown as n/a or listed under "not on disk", never as zero./ruflo swarm pane|status [json]|topology|claims|consensus, through ruflo-console's /ruflo. The old /ruflo-swarm-pane, /ruflo-swarm-status, /ruflo-swarm-topology, /ruflo-swarm-claims and /ruflo-swarm-consensus stay registered as aliases (ADR-406: no command is removed or renamed). /ruflo-swarm:watch also opens the pane./config): panel (command, the default since ruflo-console's cockpit is the pane that opens by itself | auto | off: /ruflo swarm pane refuses and /ruflo-swarm:watch does not open it), cli (npx-offline, the default, never touches the network), routeThreshold, injectAdrs (on by default: with ADRs attached to the active ruflo-console mission, appends their accepted decisions, masked and capped, to spawned subagents' prompts and names them on the member row; nothing is added without an attached ADR), injectSpawnContext (off by default: appends a one-line swarm note to spawned subagents' prompts) and audit (off by default: event names, tool names and agent ids into ruflo memory, never content).prompt.submit or tool.check hooks, and its tool.call hook only observes.modTrust: refuse-risky, the gate refuses this mod. The mod hooks tool.call and, with audit on, *, and it calls process.run. To keep it, add its provenance to modTrustAllow: ruflo-swarm@ruflo. The gate matches name@marketplace, not the plugin name. Under the default modTrust: observe the mod loads unchanged.Tests run on the engine's own kit: claude plugin test plugins/ruflo-swarm (run scripts/fetch-mod-types.sh once first if you want tsc -p plugins/ruflo-swarm to typecheck). The pure readers also run under vitest: npx vitest run --root plugins/ruflo-swarm tests/parse.spec.ts.
ruflo-swarm draws its toasts through the shared toast policy (one copy of hooks/toast-policy.ts, kept identical across the ruflo mods by scripts/sync-toast-policy.mjs): a level prefix (› info, ✓ ok, ⚠ warn, ✗ error), one clean line of at most 120 characters with secrets masked, an identical toast not repeated for a minute, and at most four a minute per source (errors are held and counted, never dropped). Next: … (when the prompt box cannot be filled) is info; a subagent whose turn ends in anything but an answer or an abort is an error; a loop that fails the same call three times running is a warn ("keeps failing …").
The person's choice is the console's Settings → Interface and updates → Toasts: all, important (warnings and errors) or off, and a mute chip per source. Every toast, drawn or not, is kept as a digest the console shows on its Events page, flagged with what became of it; without the console nothing is written and the defaults apply.
bash plugins/ruflo-swarm/scripts/smoke.sh
# Expected: "11 passed, 0 failed"
CLAUDE_CODE_ENABLE_FUNCTION_HOOKS=1 claude plugin test plugins/ruflo-swarm
CLAUDE_CODE_ENABLE_FUNCTION_HOOKS=1 claude plugin validate --strict plugins/ruflo-swarm
ruflo-agentdb — namespace convention ownerruflo-autopilot — owns the 270s cache-aware /loop heartbeat for long-running swarmsruflo-intelligence — hooks_route powers swarm agent recommendation per taskhooks/register.ts 492 lines1import type { EngineInterface, On, PluginOptions } from 'claude-code'
2
3import { paneActionsOf, type Controller } from './actions/controller'
4import { AUDIT_FLUSH_MS, auditRow, noteAudit, takeFlush } from './audit'
5import { readAdrDigest } from './adr-digest'
6import { claimsText, COMMANDS, consensusText, isSwarmSub, spawnNote, statusText, SWARM_SUBS, topologyText, type SwarmSub } from './commands'
7import type { Host, OpenResult } from './host'
8import { isAnimating, LEAD, loopLabel, newActivity, noteCall, noteDone, noteListed, noteResult, noteSpawn, stuckCall } from './model/members'
9import { parseRoute, plain } from './reader/parse'
10import { readSnapshot, type ReadCache } from './reader/snapshot'
11import { CLI_PREFIXES, newState, PANE_ID, persistedOf, restore, storeKeyOf, type State } from './state'
12import { createToastKit } from './toast-policy'
13import { paneModelOf } from './views/model'
14import { paneView } from './views/pane'
15
16const REFRESH_DEBOUNCE_MS = 400
17const TICK_MS = 2_000
18/** While idle the disk is read this often: ruflo's own CLI may change it from another terminal. */
19const IDLE_READ_MS = 10_000
20const FRAME_MS = 200
21const RUFLO_COMMAND = /(?:^|[\s/'"])(?:ruflo|claude-flow|@claude-flow\/cli)(?:@[\w.-]+)?\s/
22const ROUTE_COMMAND = /\bhooks\s+route\b/
23const ROUTE_TOOL = /(?:^|__)hooks_route$/
24
25/**
26 * Binds a Host from `$`. Declared here, each member spelled `$.noun.event(...)`, so the engine reads what the module
27 * calls off its source. The calls that answer nothing are wrapped: a refused draw or toast is not a crashed hook.
28 */
29function hostOf($: EngineInterface, cwd: string): Host {
30 const rooted = (path: string) => (path.startsWith('/') ? path : `${cwd.replace(/\/+$/, '')}/${path}`)
31 const quietly = (fn: () => void) => {
32 try {
33 fn()
34 } catch {
35 // Refused: there is nothing to do about a draw nobody may make.
36 }
37 }
38
39 const toaster = createToastKit({
40 source: 'swarm',
41 now: () => $.clock.now(),
42 show: (line, options) => $.ui.toast(line, options),
43 after: (ms, fn) => $.clock.after(ms, fn),
44 io: { read: path => $.fs.read(rooted(path)), write: (path, text) => $.fs.write(rooted(path), text), exists: path => $.fs.exists(rooted(path)) },
45 })
46
47 return {
48 fs: { read: path => $.fs.read(rooted(path)), stat: path => $.fs.stat(rooted(path)) },
49 now: () => $.clock.now(),
50 after: (ms, fn) => $.clock.after(ms, fn),
51 every: (ms, fn) => $.clock.every(ms, fn),
52 storeGet: key => $.store.get(key),
53 storeSet: (key, value) => $.store.set(key, value),
54 invalidate: () => quietly(() => $.ui.invalidate('ui.render')),
55 toast: input => quietly(() => void toaster.toast(input)),
56 log: text => quietly(() => $.ui.log(text)),
57 openPane: pane => $.ui.open(pane) as Promise<OpenResult>,
58 closePane: id => $.ui.close({ id }),
59 registerCommand: spec => $.command.register(spec),
60 fillPrompt: input => $.prompt.fill(input),
61 run: (argv, timeoutMs) => $.process.run(argv, { timeoutMs }),
62 usage: () => $.session.usage(),
63 agents: () => $.agent.list(),
64 }
65}
66
67/**
68 * The swarm mod: a live pane of the ruflo swarm on disk and of Claude Code's own agent loops, swarm commands, and,
69 * where the options ask, swarm context in spawned subagents and a content-free audit trail. `session.start` binds the
70 * engine; every hook after it works over that binding and one state object.
71 */
72export function register(on: On, raw: PluginOptions) {
73 const state: State = newState(raw, newActivity())
74 const cache: ReadCache = new Map()
75 let host: Host | null = null
76
77 const persist = () => {
78 void host?.storeSet(storeKeyOf(state.cwd), persistedOf(state)).catch(() => undefined)
79 }
80
81 async function refresh(): Promise<void> {
82 const bound = host
83
84 if (bound === null) {
85 return
86 }
87
88 if (state.isRefreshing) {
89 state.isRefreshQueued = true
90
91 return
92 }
93
94 state.isRefreshing = true
95
96 try {
97 const nowMs = Date.now()
98 const [snapshot, usage, listed] = await Promise.all([
99 readSnapshot(bound.fs, cache, nowMs),
100 bound.usage().catch(() => null),
101 bound.agents().catch(() => null),
102 ])
103 const isFirstSight = snapshot.hasSwarm && state.snapshot?.hasSwarm !== true
104
105 state.snapshot = snapshot
106 state.readError = null
107
108 if (usage !== null) {
109 state.usage = {
110 ...(usage.cost?.usd !== undefined && { costUsd: usage.cost.usd }),
111 ...(usage.context?.tokens !== undefined && { contextTokens: usage.context.tokens }),
112 ...(usage.context?.percent !== undefined && { contextPercent: usage.context.percent }),
113 ...(usage.context?.window !== undefined && { contextWindow: usage.context.window }),
114 readAtMs: nowMs,
115 }
116 }
117
118 if (listed !== null) {
119 noteListed(state.activity, listed, nowMs)
120 }
121
122 if (isFirstSight) {
123 maybeAutoOpen()
124 }
125 } catch (error) {
126 state.readError = plain(error instanceof Error ? error.message : String(error), 200)
127 } finally {
128 state.isRefreshing = false
129 bound.invalidate()
130
131 if (state.isRefreshQueued) {
132 state.isRefreshQueued = false
133 scheduleRefresh()
134 }
135 }
136 }
137
138 function scheduleRefresh(): void {
139 if (host === null || state.timers.has('refresh')) {
140 return
141 }
142
143 state.timers.set(
144 'refresh',
145 host.after(REFRESH_DEBOUNCE_MS, () => {
146 state.timers.delete('refresh')
147 void refresh()
148 }),
149 )
150 }
151
152 /** Redraws while a tile is lit, then stops: no timer runs for a pane with nothing moving. */
153 function animate(): void {
154 const bound = host
155
156 if (bound === null || state.timers.has('frames') || !state.pane.isOpen) {
157 return
158 }
159
160 state.timers.set(
161 'frames',
162 bound.every(FRAME_MS, () => {
163 bound.invalidate()
164
165 if (!state.pane.isOpen || !isAnimating(state.activity, Date.now())) {
166 state.timers.get('frames')?.cancel()
167 state.timers.delete('frames')
168 }
169 }),
170 )
171 }
172
173 async function openPane(): Promise<{ isPlaced: boolean; reason: string }> {
174 if (host === null) {
175 return { isPlaced: false, reason: 'the session has not started' }
176 }
177
178 try {
179 const result = await host.openPane({ id: PANE_ID, title: 'Swarm' })
180 const isPlaced = result === undefined || result.isPlaced !== false
181
182 state.pane.isOpen = isPlaced
183 state.pane.isClosedByPerson = false
184 persist()
185 void refresh()
186
187 return { isPlaced, reason: result?.reason ?? '' }
188 } catch (error) {
189 return { isPlaced: false, reason: plain(error instanceof Error ? error.message : String(error), 160) }
190 }
191 }
192
193 function maybeAutoOpen(viewport?: { isFullscreen?: boolean }): void {
194 state.viewport.isFullscreen = viewport?.isFullscreen ?? state.viewport.isFullscreen
195
196 const bound = host
197 const shouldOpen =
198 bound !== null &&
199 state.options.panel === 'auto' &&
200 !state.pane.isOpen &&
201 !state.pane.isClosedByPerson &&
202 !state.timers.has('auto-open') &&
203 state.snapshot?.hasSwarm === true &&
204 state.viewport.isFullscreen === true
205
206 if (shouldOpen && bound !== null) {
207 // From a timer: a render hook only draws.
208 state.timers.set(
209 'auto-open',
210 bound.after(50, () => {
211 state.timers.delete('auto-open')
212 void openPane()
213 }),
214 )
215 }
216 }
217
218 const control: Controller = { refresh, persist }
219
220 // ------------------------------------------------------------- lifecycle
221
222 on('session.start', async ($, e, next) => {
223 host = hostOf($, e.cwd)
224 state.cwd = e.cwd
225
226 for (const timer of state.timers.values()) {
227 timer.cancel()
228 }
229
230 state.timers.clear()
231 noteSpawn(state.activity, LEAD, 'lead', Date.now())
232
233 const bound = host
234
235 await Promise.all([
236 ...COMMANDS.map(spec => bound.registerCommand(spec).catch(() => undefined)),
237 bound
238 .storeGet(storeKeyOf(e.cwd))
239 .then(value => restore(state, value))
240 .catch(() => undefined),
241 ])
242
243 await refresh()
244
245 let lastReadMs = Date.now()
246
247 state.timers.set(
248 'tick',
249 bound.every(TICK_MS, () => {
250 const nowMs = Date.now()
251
252 if (state.pane.isOpen && (state.activity.isWorking || nowMs - lastReadMs >= IDLE_READ_MS)) {
253 lastReadMs = nowMs
254 scheduleRefresh()
255 }
256 }),
257 )
258
259 if (state.options.audit) {
260 state.timers.set(
261 'audit',
262 bound.every(AUDIT_FLUSH_MS, () => {
263 const flush = takeFlush(state, Date.now())
264
265 if (flush !== null) {
266 void bound.run([...CLI_PREFIXES[state.options.cli], ...flush.args], 30_000).catch(() => undefined)
267 }
268 }),
269 )
270 }
271
272 return next(e)
273 })
274
275 on('ui.render', { component: 'PromptHint' }, ($, e, next) => {
276 if (e.surface === 'terminal') {
277 maybeAutoOpen(e.viewport)
278 }
279
280 return next(e)
281 })
282
283 on('ui.render', { component: 'Pane' }, ($, e, next) => {
284 if (host === null || e.requestId !== PANE_ID) {
285 return next(e)
286 }
287
288 const table = $.ui.resolve(e)
289 const columns = Math.max(20, Math.floor(Number(e.props.bodyColumns) || 0) - 1)
290 const rows = Math.max(0, Math.floor(Number(e.props.scroll?.bodyRows) || 0))
291 const nowMs = Date.now()
292
293 state.pane.isOpen = true
294 state.pane.columns = columns
295 state.pane.rows = rows
296
297 if (isAnimating(state.activity, nowMs)) {
298 animate()
299 }
300
301 return paneView({ Box: table.Box, Text: table.Text, Button: table.Button }, paneModelOf(state, columns, rows, nowMs), paneActionsOf(state, host, control))
302 })
303
304 on('ui.close', async ($, e, next) => {
305 const result = await next(e)
306
307 if (e.id === PANE_ID && result.deny === undefined) {
308 state.pane.isOpen = false
309 state.pane.isClosedByPerson = e.origin.kind === 'person' || state.pane.isClosedByPerson
310 persist()
311 }
312
313 return result
314 })
315
316 // -------------------------------------------------------------- commands
317
318 /** The five swarm commands, each by its subcommand word: what `/ruflo swarm <sub>` and the old names both answer. */
319 async function answer(sub: SwarmSub, args: string): Promise<{ text: string }> {
320 const bound = host as Host
321
322 switch (sub) {
323 case 'pane': {
324 const arg = args.trim().toLowerCase()
325
326 if (arg === 'close' || (arg === '' && state.pane.isOpen)) {
327 state.pane.isOpen = false
328 state.pane.isClosedByPerson = true
329 persist()
330 await bound.closePane(PANE_ID).catch(() => undefined)
331
332 return { text: 'Swarm pane hidden' }
333 }
334
335 if (state.options.panel === 'off') {
336 return { text: 'The swarm pane is off (panel option); set panel to command or auto in /config to open it' }
337 }
338
339 const opened = await openPane()
340
341 return { text: opened.isPlaced ? 'Swarm pane shown' : `The swarm pane could not be shown: ${opened.reason}` }
342 }
343 case 'status':
344 await refresh()
345
346 return { text: statusText(state, Date.now(), args.trim().toLowerCase() === 'json') }
347 case 'topology':
348 await refresh()
349
350 return { text: topologyText(state, Date.now()) }
351 case 'claims':
352 await refresh()
353
354 return { text: claimsText(state) }
355 case 'consensus':
356 await refresh()
357
358 return { text: consensusText(state) }
359 }
360 }
361
362 /**
363 * `/ruflo swarm <pane|status|topology|claims|consensus>`: ruflo-console registers `/ruflo` for every ruflo mod, and
364 * this hook answers its swarm subcommands wherever it sits in the chain, passing every other word on.
365 */
366 // `/ruflo-console` is the same command as `/ruflo` (kept by ADR-406), so its swarm subcommands are answered too.
367 for (const command of ['ruflo', 'ruflo-console'] as const) {
368 on('command.run', { command }, async ($, e, next) => {
369 const [head = '', sub = '', ...rest] = e.args.trim().split(/\s+/)
370
371 if (host === null || head.toLowerCase() !== 'swarm' || !isSwarmSub(sub.toLowerCase())) {
372 return next(e)
373 }
374
375 return answer(sub.toLowerCase() as SwarmSub, rest.join(' '))
376 })
377 }
378
379 // The old names stay registered as aliases of `/ruflo swarm <sub>`: ADR-406 removes, renames or reassigns no command.
380 for (const sub of SWARM_SUBS) {
381 on('command.run', { command: `ruflo-swarm-${sub}` }, async ($, e, next) => (host === null ? next(e) : answer(sub, e.args)))
382 }
383
384 /** The plugin's markdown `/ruflo-swarm:watch` still runs as it always has; with the mod it opens the pane as well. */
385 on('command.run', { command: 'ruflo-swarm:watch' }, async ($, e, next) => {
386 if (host !== null && state.options.panel !== 'off') {
387 await openPane()
388 }
389
390 return next(e)
391 })
392
393 // ----------------------------------------------------------------- turns
394
395 on('turn.start', ($, e, next) => {
396 state.activity.isWorking = true
397 state.activity.turnId = e.turnId
398 host?.invalidate()
399
400 return next(e)
401 })
402
403 on('turn.complete', ($, e, next) => {
404 if (e.agentId !== undefined) {
405 noteDone(state.activity, e.agentId, e.reason, Date.now())
406
407 // An answer is a success and an abort is the person's own; anything else is a loop that ended badly.
408 if (e.reason !== 'answer' && e.reason !== 'aborted') host?.toast({ level: 'error', text: `swarm: ${loopLabel(state.activity, e.agentId)} failed (${plain(String(e.reason), 40)})` })
409 } else {
410 state.activity.isWorking = false
411 scheduleRefresh()
412 }
413
414 host?.invalidate()
415
416 return next(e)
417 })
418
419 // ---------------------------------------------------------------- agents
420
421 /** Claude Code's own subagents join the pane as members; with the option on, each is told which swarm it is in. */
422 on('agent.spawn', async ($, e, next) => {
423 const note = state.options.injectSpawnContext ? spawnNote(state.snapshot) : null
424 // ADR-480: the accepted decisions attached to the console's active mission, as data after the task.
425 const adr = state.options.injectAdrs && host !== null ? await readAdrDigest(host.fs, await host.now().catch(() => Date.now())) : null
426 const prompt = `${e.prompt}${note ?? ''}${adr === null ? '' : `\n\n---\n${adr.block}`}`
427 const result = await next(prompt === e.prompt ? e : { ...e, prompt })
428
429 if (result.deny === undefined && result.agentId !== undefined) {
430 noteSpawn(state.activity, result.agentId, e.subagentType, Date.now(), e.name, plain(e.description, 120), adr?.numbers ?? [])
431 host?.invalidate()
432 }
433
434 return result
435 })
436
437 // ------------------------------------------------------------ tool calls
438
439 /** Observes only: the call goes on unchanged, and what it read or wrote lights its agent's tile. */
440 on('tool.call', async ($, e, next) => {
441 const args = e as unknown as Readonly<Record<string, unknown>>
442 const command = typeof args.command === 'string' ? args.command : undefined
443 const path = typeof args.file_path === 'string' ? args.file_path : typeof args.path === 'string' ? args.path : typeof args.pattern === 'string' ? args.pattern : undefined
444 const subject = plain(command ?? path ?? '', 80)
445 const pulse = noteCall(state.activity, e.agentId, e.tool, subject, Date.now(), command)
446
447 if (pulse !== null && state.pane.isOpen) {
448 animate()
449 host?.invalidate()
450 }
451
452 const result = await next(e)
453 const isError = result.deny !== undefined || result.isError === true
454
455 noteResult(state.activity, e.agentId, e.tool, subject, isError, Date.now())
456
457 // A loop that fails the same call three times running is going round in circles (ADR-477). The policy says it once a minute at most.
458 const stuck = isError ? stuckCall(state.activity, e.agentId) : null
459
460 if (stuck !== null) host?.toast({ level: 'warn', text: `swarm: ${loopLabel(state.activity, e.agentId ?? LEAD)} keeps failing ${stuck.tool}${stuck.subject === '' ? '' : ` ${stuck.subject}`}` })
461
462 if (!isError) {
463 const text = typeof result.text === 'string' ? result.text : ''
464
465 if ((command !== undefined && ROUTE_COMMAND.test(command)) || ROUTE_TOOL.test(e.tool)) {
466 const pick = parseRoute(text, Date.now())
467
468 if (pick !== null) {
469 state.route = pick
470 persist()
471 }
472 }
473
474 if ((command !== undefined && RUFLO_COMMAND.test(command)) || e.tool.includes('claude-flow') || e.tool.includes('ruflo')) {
475 scheduleRefresh()
476 }
477 }
478
479 return result
480 })
481
482 // ----------------------------------------------------------------- audit
483
484 if (state.options.audit) {
485 on('*', ($, e, next) => {
486 noteAudit(state, auditRow(next.event, e, Date.now()))
487
488 return next(e)
489 })
490 }
491}
492hooks/actions/controller.ts 291 lines1import type { Host } from '../host'
2import { LEAD, membersOf, type Member } from '../model/members'
3import { parseRoute, plain } from '../reader/parse'
4import { CLI_PREFIXES, type State } from '../state'
5import type { PaneActions } from '../views/pane'
6import { paneModelOf } from '../views/model'
7import {
8 agentLogs,
9 argvOf,
10 claimTask,
11 handoffTask,
12 markStealable,
13 reroute,
14 setClaimStatus,
15 stealTask,
16 stopAgent,
17 vote,
18 type ActionSpec,
19} from './argv'
20
21/** A CLI call from a button gets this long; the engine's own cap is ten minutes, and the pane does not hold a person that long. */
22const RUN_MS = 60_000
23/** A confirm prompt left unanswered this long is dropped, so a stale press never runs an old action. */
24export const CONFIRM_MS = 30_000
25
26export type Controller = {
27 /** Re-reads the disk now and resolves once the snapshot is fresh. */
28 refresh: () => Promise<void>
29 persist: () => void
30}
31
32const lastLines = (text: string, count: number): string[] =>
33 text
34 .split('\n')
35 .map(line => plain(line, 400))
36 .filter(line => line !== '')
37 .slice(-count)
38
39/** Runs one action through the configured CLI and says what happened, checked against the disk where the disk can show it. */
40export async function runAction(state: State, host: Host, control: Controller, spec: ActionSpec): Promise<void> {
41 const nowMs = () => Date.now()
42
43 state.isActing = true
44 state.confirm = null
45 host.invalidate()
46
47 try {
48 const result = await host.run(argvOf(CLI_PREFIXES[state.options.cli], spec), RUN_MS)
49 const ok = result.exitCode === 0
50 let verified: 'yes' | 'no' | 'n/a' = 'n/a'
51
52 if (spec.args[0] === 'hooks' && spec.args[1] === 'route') {
53 const pick = parseRoute(result.stdout, nowMs())
54
55 if (pick !== null) {
56 state.route = pick
57 control.persist()
58 }
59
60 verified = pick !== null ? 'yes' : 'no'
61 } else if (spec.args[0] === 'agent' && spec.args[1] === 'logs') {
62 state.detail = { title: spec.label, lines: lastLines(ok ? result.stdout : result.stderr || result.stdout, 40) }
63 }
64
65 if (ok && spec.verify !== undefined) {
66 await control.refresh()
67 verified = state.snapshot !== null && spec.verify(state.snapshot) ? 'yes' : 'no'
68 }
69
70 const said = lastLines(ok ? result.stdout : result.stderr || result.stdout, 1)[0] ?? ''
71
72 state.outcome = { label: spec.label, ok, verified, detail: ok ? (verified === 'no' ? `expected ${spec.expect}` : '') : `exit ${result.exitCode}${said !== '' ? `: ${said}` : ''}`, atMs: nowMs() }
73 } catch (error) {
74 // A refused `$.process.run` (removed by an administrator, or denied above this mod) is a sentence, not a crash.
75 state.outcome = { label: spec.label, ok: false, verified: 'n/a', detail: plain(error instanceof Error ? error.message : String(error), 160), atMs: nowMs() }
76 } finally {
77 state.isActing = false
78 host.invalidate()
79 }
80}
81
82/** Asks first for a destructive action (a second press confirms), runs a safe one at once, and says so when there is nothing to act on. */
83export function request(state: State, host: Host, control: Controller, spec: ActionSpec | null, why: string): void {
84 if (state.isActing) {
85 return
86 }
87
88 if (spec === null) {
89 state.outcome = { label: why, ok: false, verified: 'n/a', detail: 'nothing to act on', atMs: Date.now() }
90 host.invalidate()
91
92 return
93 }
94
95 if (spec.isDestructive) {
96 state.confirm = { label: spec.label, spec, askedAtMs: Date.now() }
97 state.timers.get('confirm')?.cancel()
98 state.timers.set(
99 'confirm',
100 host.after(CONFIRM_MS, () => {
101 state.timers.delete('confirm')
102 state.confirm = null
103 host.invalidate()
104 }),
105 )
106 host.invalidate()
107
108 return
109 }
110
111 void runAction(state, host, control, spec)
112}
113
114/** The activity a Claude Code loop showed this session, as lines: there is no transcript API, so this is what the hooks saw. */
115function activityLines(state: State, member: Member): string[] {
116 const recent = state.activity.recent.get(member.id) ?? []
117 const loop = state.activity.loops.get(member.id)
118
119 return [
120 ...(loop?.description !== undefined ? [`task: ${plain(loop.description, 200)}`] : []),
121 ...(recent.length === 0 ? ['no tool calls seen yet'] : recent.map(call => `${call.isError ? '✗' : '·'} ${call.tool} ${call.subject}`)),
122 ]
123}
124
125export function paneActionsOf(state: State, host: Host, control: Controller): PaneActions {
126 const model = () => paneModelOf(state, state.pane.columns, state.pane.rows, Date.now())
127 const members = () => membersOf(state.snapshot, state.activity, Date.now())
128
129 const step = (by: number) => {
130 const all = members()
131
132 if (all.length === 0) {
133 return
134 }
135
136 const at = Math.max(0, all.findIndex(member => member.id === state.selected))
137
138 state.selected = all[(at + by + all.length) % all.length]?.id ?? null
139 state.detail = null
140 control.persist()
141 host.invalidate()
142 }
143
144 const stepTask = (by: number) => {
145 const rows = model().board.rows
146
147 if (rows.length === 0) {
148 return
149 }
150
151 const at = Math.max(0, rows.findIndex(row => row.id === state.selectedTask))
152
153 state.selectedTask = rows[(at + by + rows.length) % rows.length]?.id ?? null
154 control.persist()
155 host.invalidate()
156 }
157
158 const selected = () => model().selected
159 const task = () => model().board.selected
160 const ruflo = (member: Member | null) => (member !== null && member.source === 'ruflo' ? member : null)
161 const claimOf = (member: Member | null) => (member === null ? undefined : state.snapshot?.claims.find(claim => claim.claimant.id === member.id && claim.status !== 'completed'))
162 const go = (spec: ActionSpec | null, why: string) => request(state, host, control, spec, why)
163
164 return {
165 hide: () => {
166 state.pane.isOpen = false
167 state.pane.isClosedByPerson = true
168 control.persist()
169 void host.closePane('ruflo-swarm').catch(() => undefined)
170 },
171 prev: () => step(-1),
172 next: () => step(1),
173 taskPrev: () => stepTask(-1),
174 taskNext: () => stepTask(1),
175 stop: () => {
176 const member = ruflo(selected())
177
178 go(member === null ? null : stopAgent(member.id), 'stop')
179 },
180 logs: () => {
181 const member = selected()
182
183 if (member === null) {
184 return
185 }
186
187 if (member.source === 'ruflo') {
188 go(agentLogs(member.id), 'logs')
189
190 return
191 }
192
193 state.detail = { title: `activity of ${member.label}${member.id === LEAD ? '' : ` (${member.id})`}`, lines: activityLines(state, member) }
194 host.invalidate()
195 },
196 pause: () => {
197 const claim = claimOf(ruflo(selected()))
198
199 go(claim === undefined ? null : setClaimStatus(claim.issueId, 'paused'), 'pause claim')
200 },
201 resume: () => {
202 const claim = claimOf(ruflo(selected()))
203
204 go(claim === undefined ? null : setClaimStatus(claim.issueId, 'active'), 'resume claim')
205 },
206 claim: () => {
207 const member = ruflo(selected())
208 const row = task()
209
210 go(member === null || row === null ? null : claimTask(row.id, member.id, member.role), 'claim task')
211 },
212 offer: () => {
213 const row = task()
214
215 go(row === null ? null : markStealable(row.id), 'offer task')
216 },
217 steal: () => {
218 const member = ruflo(selected())
219 const row = task()
220
221 go(member === null || row === null ? null : stealTask(row.id, member.id, member.role), 'steal task')
222 },
223 handoff: () => {
224 const member = ruflo(selected())
225 const row = task()
226 const claim = row === null ? undefined : state.snapshot?.claims.find(entry => entry.issueId === row.id)
227 const from = claim?.claimant.kind === 'agent' ? { id: claim.claimant.id, type: claim.claimant.agentType ?? 'agent' } : null
228
229 go(member === null || row === null || from === null ? null : handoffTask(row.id, from, { id: member.id, type: member.role }), 'hand off')
230 },
231 reroute: () => {
232 const row = task()
233
234 go(row === null ? null : reroute(row.description || row.type), 're-route')
235 },
236 voteYes: () => {
237 const proposal = model().proposals[0]
238 const member = selected()
239
240 go(proposal === undefined || member === null ? null : vote(proposal.id, member.id, true), 'vote')
241 },
242 voteNo: () => {
243 const proposal = model().proposals[0]
244 const member = selected()
245
246 go(proposal === undefined || member === null ? null : vote(proposal.id, member.id, false), 'vote')
247 },
248 confirm: () => {
249 const asked = state.confirm
250
251 state.timers.get('confirm')?.cancel()
252 state.timers.delete('confirm')
253
254 if (asked === null || Date.now() - asked.askedAtMs > CONFIRM_MS) {
255 state.confirm = null
256 host.invalidate()
257
258 return
259 }
260
261 void runAction(state, host, control, asked.spec)
262 },
263 cancel: () => {
264 state.timers.get('confirm')?.cancel()
265 state.timers.delete('confirm')
266 state.confirm = null
267 host.invalidate()
268 },
269 fill: () => {
270 const next = model().next
271
272 if (next === null) {
273 return
274 }
275
276 void host
277 .fillPrompt({ text: next.text })
278 .then(result => {
279 if (!result.isFilled) {
280 host.toast({ level: 'info', text: `Next: ${next.text}`, timeoutMs: 8000, awayOnly: true })
281 }
282 })
283 .catch(() => host.toast({ level: 'info', text: `Next: ${next.text}`, timeoutMs: 8000, awayOnly: true }))
284 },
285 closeDetail: () => {
286 state.detail = null
287 host.invalidate()
288 },
289 }
290}
291hooks/audit.ts 60 lines1import { idOf } from './reader/parse'
2import type { State } from './state'
3
4/** At most this many rows wait for a flush; past it the oldest are dropped and counted. */
5export const AUDIT_MAX = 200
6export const AUDIT_FLUSH_MS = 30_000
7export const AUDIT_NAMESPACE = 'ruflo-swarm-audit'
8
9/** Events never recorded: the render and timer traffic would drown the trail, and the trail's own writes would feed it. */
10const SKIP = /^(ui\.|clock\.|store\.|fs\.|process\.run|session\.usage|agent\.list|env\.|config\.describe)/
11
12const NAME = /^[A-Za-z0-9_.:-]{1,64}$/
13
14/**
15 * One row of the trail: the event's name, the tool's name and the agent's id where the event has them, and the time.
16 * Never prompt text, arguments, file paths, results or anything a tool returned: the trail says what happened, not what was said.
17 */
18export function auditRow(event: string, e: unknown, nowMs: number): string | null {
19 if (SKIP.test(event) || !NAME.test(event)) {
20 return null
21 }
22
23 const value = (e ?? {}) as Record<string, unknown>
24 const tool = typeof value.tool === 'string' && NAME.test(value.tool) ? value.tool : undefined
25 const agent = idOf(value.agentId)
26
27 return JSON.stringify({ t: nowMs, event, ...(tool !== undefined && { tool }), ...(agent !== null && { agent }) })
28}
29
30/** Adds a row on the hot path: a push and, when full, a shift. Nothing waits on I/O here. */
31export function noteAudit(state: State, row: string | null): void {
32 if (row === null) {
33 return
34 }
35
36 state.audit.buffer.push(row)
37
38 if (state.audit.buffer.length > AUDIT_MAX) {
39 state.audit.buffer.shift()
40 state.audit.dropped += 1
41 }
42}
43
44/** The rows to write now and the argv after the CLI prefix that stores them, or null when there is nothing to write. */
45export function takeFlush(state: State, nowMs: number): { args: string[]; rows: number } | null {
46 if (state.audit.buffer.length === 0) {
47 return null
48 }
49
50 const rows = state.audit.buffer.splice(0, state.audit.buffer.length)
51 const dropped = state.audit.dropped
52
53 state.audit.dropped = 0
54 state.audit.flushedAtMs = nowMs
55
56 const value = JSON.stringify({ rows: rows.map(row => JSON.parse(row) as unknown), ...(dropped > 0 && { dropped }) })
57
58 return { args: ['memory', 'store', '--namespace', AUDIT_NAMESPACE, '--key', `audit-${nowMs}`, '--value', value], rows: rows.length }
59}
60hooks/adr-digest.ts 41 lines1/**
2 * The ADRs a mission carries, for the subagents it spawns (ADR-480 in ruflo-console). The console writes a small, masked digest of the
3 * project's ADRs attached to the active mission to `.claude-flow/console/adr-digest.json`; this reads it when a subagent is spawned and
4 * hands the block to the agent as data. Nothing is written here, nothing is fetched, and a file that is missing, large, stale or oddly
5 * shaped is simply no digest. The block names the decisions in force; it does not prove the agent follows them.
6 */
7import type { ReaderFs } from './reader/snapshot'
8import { plain } from './reader/parse'
9
10export const ADR_DIGEST_FILE = '.claude-flow/console/adr-digest.json'
11const MAX_BYTES = 16 * 1024
12const MAX_BLOCK = 1400
13const MAX_AGE_MS = 24 * 60 * 60 * 1000
14
15/** A terminal escape sequence: removed whole, so no remnant such as `[31m` is left behind by the control-character wash. */
16// eslint-disable-next-line no-control-regex
17const ANSI = /\u001b\[[0-9;?]*[ -/]*[@-~]|\u001b\][^\u0007\u001b]*(?:\u0007|\u001b\\)/g
18
19export type AdrDigest = { block: string; numbers: number[] }
20
21/** The digest for the project, or null. `nowMs` decides staleness (a digest older than a day is not trusted). */
22export async function readAdrDigest(fs: ReaderFs, nowMs: number): Promise<AdrDigest | null> {
23 try {
24 const stat = await fs.stat(ADR_DIGEST_FILE)
25
26 if (stat === undefined || (stat.size ?? 0) > MAX_BYTES) return null
27
28 const parsed = JSON.parse(await fs.read(ADR_DIGEST_FILE)) as Record<string, unknown>
29
30 if (typeof parsed !== 'object' || parsed === null || parsed.v !== 1 || typeof parsed.block !== 'string' || parsed.block === '') return null
31 if (typeof parsed.atMs !== 'number' || !Number.isFinite(parsed.atMs) || nowMs - parsed.atMs > MAX_AGE_MS || parsed.atMs - nowMs > MAX_AGE_MS) return null
32
33 const block = parsed.block.split('\n').slice(0, 12).map(line => plain(line.replace(ANSI, ''), 320)).filter(line => line !== '').join('\n').slice(0, MAX_BLOCK)
34 const numbers = Array.isArray(parsed.adrs) ? parsed.adrs.flatMap(each => (typeof each === 'object' && each !== null && typeof (each as { number?: unknown }).number === 'number' && (each as { status?: unknown }).status === 'accepted' ? [(each as { number: number }).number] : [])).filter(Number.isSafeInteger).slice(0, 8) : []
35
36 return block === '' ? null : { block, numbers }
37 } catch {
38 return null
39 }
40}
41hooks/commands.ts 151 lines1import type { CommandSpec } from 'claude-code'
2
3import { membersOf } from './model/members'
4import { topologyLines } from './model/topology'
5import { plain } from './reader/parse'
6import type { Snapshot } from './reader/snapshot'
7import type { State } from './state'
8import { paneModelOf } from './views/model'
9
10/**
11 * The commands this module registers: since ruflo-console's `/ruflo`, each is also reachable as `/ruflo swarm <sub>`,
12 * and each stays registered (ADR-406: no command is removed or renamed). Each name is prefixed with the plugin's so it never stands in for a built-in.
13 */
14export const COMMANDS: readonly CommandSpec[] = [
15 { name: 'ruflo-swarm-pane', description: 'Same as /ruflo swarm pane: show or hide the live swarm pane', argumentHint: '[open|close]', immediate: true },
16 { name: 'ruflo-swarm-status', description: 'Same as /ruflo swarm status: the swarm as ruflo has it on disk', argumentHint: '[json]', immediate: true },
17 { name: 'ruflo-swarm-topology', description: 'Same as /ruflo swarm topology: the topology, its leader and members', immediate: true },
18 { name: 'ruflo-swarm-claims', description: 'Same as /ruflo swarm claims: who has claimed which task', immediate: true },
19 { name: 'ruflo-swarm-consensus', description: 'Same as /ruflo swarm consensus: hive-mind proposals and votes', immediate: true },
20]
21
22/** The subcommands of `/ruflo swarm` this module answers; the old `/ruflo-swarm-<sub>` names are their aliases. */
23export const SWARM_SUBS = ['pane', 'status', 'topology', 'claims', 'consensus'] as const
24
25export type SwarmSub = (typeof SWARM_SUBS)[number]
26
27export const isSwarmSub = (word: string): word is SwarmSub => (SWARM_SUBS as readonly string[]).includes(word)
28
29const NO_SWARM = 'No ruflo swarm on disk in this folder. Start one with /ruflo-swarm:swarm init.'
30
31function missingLine(snapshot: Snapshot): string {
32 return snapshot.missing.length === 0 ? '' : `\nNot on disk: ${snapshot.missing.join(', ')}.`
33}
34
35export function statusText(state: State, nowMs: number, asJson: boolean): string {
36 const snapshot = state.snapshot
37
38 if (snapshot === null) {
39 return state.readError !== null ? `The swarm files could not be read: ${plain(state.readError, 200)}` : NO_SWARM
40 }
41
42 const model = paneModelOf(state, 100, 60, nowMs)
43
44 if (asJson) {
45 return `\`\`\`json\n${JSON.stringify(
46 {
47 swarm: snapshot.swarm,
48 members: model.members.map(member => ({ id: member.id, label: member.label, source: member.source, state: member.state, word: member.word, isLeader: member.isLeader })),
49 tasks: model.board,
50 claims: snapshot.claims,
51 hive: snapshot.hive,
52 ruos: snapshot.ruos,
53 route: state.route,
54 usage: state.usage,
55 missing: snapshot.missing,
56 },
57 null,
58 2,
59 )}\n\`\`\``
60 }
61
62 if (!snapshot.hasSwarm) {
63 return `${NO_SWARM}${missingLine(snapshot)}`
64 }
65
66 const b = model.board
67 const lines = [
68 model.title,
69 `agents: ${model.counts.map(entry => `${entry.state} ${entry.count}`).join(', ') || 'none'}`,
70 ...model.members.slice(0, 20).map(member => ` ${member.isLeader ? '★' : '·'} ${member.label} (${member.source}) ${member.word}`),
71 `tasks: ${b.pending} pending, ${b.claimed} claimed, ${b.done} done${b.failed > 0 ? `, ${b.failed} failed` : ''}`,
72 ...(model.remote.length > 0 ? ['ruOS hosts:', ...model.remote.map(line => ` ${line}`)] : []),
73 model.route,
74 model.usage,
75 ]
76
77 return `${lines.join('\n')}${missingLine(snapshot)}`
78}
79
80export function topologyText(state: State, nowMs: number): string {
81 const snapshot = state.snapshot
82
83 if (snapshot === null || !snapshot.hasSwarm) {
84 return NO_SWARM
85 }
86
87 const name = snapshot.swarm?.topology ?? snapshot.hive?.topology ?? 'unknown'
88 const members = membersOf(snapshot, state.activity, nowMs)
89 const extra = [
90 snapshot.swarm?.strategy !== undefined ? `strategy ${snapshot.swarm.strategy}` : '',
91 snapshot.swarm?.maxAgents !== undefined ? `max ${snapshot.swarm.maxAgents} agents` : '',
92 snapshot.hive?.strategy !== undefined ? `hive consensus ${snapshot.hive.strategy}` : '',
93 ].filter(part => part !== '')
94
95 return [`${name}${extra.length > 0 ? ` (${extra.join(', ')})` : ''}`, ...topologyLines(name, members, 80)].join('\n')
96}
97
98export function claimsText(state: State): string {
99 const snapshot = state.snapshot
100
101 if (snapshot === null || !snapshot.hasSwarm) {
102 return NO_SWARM
103 }
104
105 if (snapshot.missing.includes('claims')) {
106 return 'No claims on disk (.claude-flow/claims/claims.json). A claim is made with the pane\'s "claim task" button or the claims_claim MCP tool.'
107 }
108
109 if (snapshot.claims.length === 0) {
110 return 'No claims.'
111 }
112
113 return snapshot.claims
114 .slice(0, 40)
115 .map(claim => `${claim.issueId}: ${claim.status}${claim.isStealable ? ' (stealable)' : ''} by ${claim.claimant.kind} ${claim.claimant.agentType ?? ''} ${claim.claimant.id}${claim.handoffTo !== undefined ? ` → handoff to ${claim.handoffTo}` : ''}`.replace(/\s+/g, ' '))
116 .join('\n')
117}
118
119export function consensusText(state: State): string {
120 const hive = state.snapshot?.hive
121
122 if (hive === null || hive === undefined) {
123 return 'No hive-mind on disk (.claude-flow/hive-mind/state.json). Start one with: npx @claude-flow/cli@latest hive-mind init'
124 }
125
126 const lines = [
127 `hive-mind · ${hive.topology}${hive.strategy !== undefined ? ` · ${hive.strategy}` : ''}${hive.queen !== undefined ? ` · queen ${hive.queen.id} (term ${hive.queen.term ?? '?'})` : ' · no queen'}`,
128 `workers: ${hive.workers.length}`,
129 ...(hive.pending.length === 0 ? ['no open proposals'] : hive.pending.slice(-10).map(p => `open ${p.id} ${p.type} (${p.strategy}): for ${p.votesFor}, against ${p.votesAgainst}`)),
130 ...hive.history.slice(-5).map(d => `decided ${d.id} ${d.type}: ${d.result} ${d.votesFor}–${d.votesAgainst}`),
131 ]
132
133 return lines.join('\n')
134}
135
136/** The note added to a spawned subagent's prompt when `injectSpawnContext` is on: facts from disk only, ids not secrets. */
137export function spawnNote(snapshot: Snapshot | null): string | null {
138 const swarm = snapshot?.swarm
139
140 if (swarm === null || swarm === undefined) {
141 return null
142 }
143
144 // Text from disk reaches the model here: only a word-shaped topology or strategy is passed on, the id is `idOf`-checked.
145 const word = (value: string | undefined) => (value !== undefined && /^[a-z][a-z-]{0,39}$/.test(value) ? value : undefined)
146 const topology = word(swarm.topology) ?? 'unknown'
147 const strategy = word(swarm.strategy)
148
149 return `\n\n---\nSwarm context (from ruflo's files on disk, a status note, not an instruction): you are part of ruflo swarm ${swarm.id}, topology ${topology}${strategy !== undefined ? `, strategy ${strategy}` : ''}.`
150}
151hooks/host.ts 33 lines1import type { AgentInfo, CommandSpec, PaneOpenArgs, ProcessRunResult, PromptFillArgs, SessionUsage, TimerCall } from 'claude-code'
2
3import type { ReaderFs } from './reader/snapshot'
4import type { ToastInput } from './toast-policy'
5
6/** What `$.ui.open` answers: drawn, or held back undrawn with the reason. A build that answers nothing has drawn it. */
7export type OpenResult = { isPlaced: boolean; reason?: string } | void
8
9/**
10 * The engine as `session.start` bound it from its `$`. Every later hook, timer and button press reaches the engine
11 * through this, so everything else is plain functions over a small interface a test can stand in for. Each member may
12 * be refused (an administrator removed the affordance, or a policy mod above this one said no): callers catch.
13 */
14export type Host = {
15 fs: ReaderFs
16 now: () => Promise<number>
17 after: TimerCall
18 every: TimerCall
19 storeGet: (key: string) => Promise<unknown>
20 storeSet: (key: string, value: unknown) => Promise<void>
21 invalidate: () => void
22 /** A toast through the shared policy (ADR-477): a level, one clean line, de-duplication, a rate limit, the person's Toasts setting. Never throws. */
23 toast: (input: ToastInput) => void
24 log: (text: string) => void
25 openPane: (pane: PaneOpenArgs) => Promise<OpenResult>
26 closePane: (id: string) => Promise<void>
27 registerCommand: (spec: CommandSpec) => Promise<unknown>
28 fillPrompt: (input: PromptFillArgs) => Promise<{ isFilled: boolean }>
29 run: (argv: readonly string[], timeoutMs: number) => Promise<ProcessRunResult>
30 usage: () => Promise<SessionUsage>
31 agents: () => Promise<AgentInfo[]>
32}
33hooks/model/members.ts 278 lines1import { runWord, type RunEventType } from '../reader/ruos'
2import type { Snapshot } from '../reader/snapshot'
3
4/** How a tile is coloured: the five states the pane names. */
5export type MemberState = 'idle' | 'working' | 'blocked' | 'done' | 'failed'
6
7export type PulseKind = 'read' | 'write'
8
9/** One Claude Code agent loop as the engine reported it (`agent.spawn`, `$.agent.list()`), and what it has done since. */
10export type LoopRecord = {
11 id: string
12 /** The agent type it runs as, without a plugin's prefix (`ruflo-swarm:coordinator` → `coordinator`). */
13 role: string
14 name?: string
15 description?: string
16 /** The numbers of the accepted ADRs the agent was told about when it was spawned (ruflo-console, ADR-480). */
17 adrs?: number[]
18 status: string
19 calls: number
20 errors: number
21 lastTool?: string
22 lastSubject?: string
23 lastAtMs: number
24 pulse?: { kind: PulseKind; atMs: number }
25}
26
27/** What the session's loops are doing, kept by the hooks between reads of the disk. */
28export type Activity = {
29 /** Claude Code's own subagents and teammates, by agent id; `lead` is the main loop. */
30 loops: Map<string, LoopRecord>
31 turnId: string | null
32 isWorking: boolean
33 /** The last few calls of each loop, newest last, for the activity view. */
34 recent: Map<string, { tool: string; subject: string; isError: boolean; atMs: number }[]>
35}
36
37export const LEAD = 'lead'
38export const PULSE_MS = 2_000
39const RECENT = 12
40const MAX_LOOPS = 200
41
42export function newActivity(): Activity {
43 return { loops: new Map(), turnId: null, isWorking: false, recent: new Map() }
44}
45
46export type Member = {
47 id: string
48 label: string
49 role: string
50 source: 'ruflo' | 'claude' | 'lead' | 'queen'
51 state: MemberState
52 /** The word under the colour, as the source said it (`busy`, `terminated`, `running`). */
53 word: string
54 pulse?: PulseKind
55 isLeader: boolean
56 claims: number
57 calls: number
58}
59
60const WRITE_TOOLS = new Set(['Edit', 'Write', 'NotebookEdit', 'MultiEdit'])
61const READ_TOOLS = new Set(['Read', 'Glob', 'Grep', 'LS', 'NotebookRead'])
62
63/** Whether a tool call reads or writes files, or neither (a pulse is for file work only). */
64export function pulseOf(tool: string, command?: string): PulseKind | null {
65 if (WRITE_TOOLS.has(tool)) {
66 return 'write'
67 }
68
69 if (READ_TOOLS.has(tool)) {
70 return 'read'
71 }
72
73 if ((tool === 'Bash' || tool === 'PowerShell') && command !== undefined) {
74 return /(?:>|\btee\b|\bsed\s+-i|\bmv\b|\bcp\b|\brm\b|\bmkdir\b|\btouch\b|\bpatch\b|\bgit\s+(?:apply|checkout|restore|commit))/.test(command) ? 'write' : 'read'
75 }
76
77 return null
78}
79
80export const roleOf = (type: string): string => (type.includes(':') ? type.slice(type.lastIndexOf(':') + 1) : type) || 'agent'
81
82function loopFor(activity: Activity, id: string, nowMs: number): LoopRecord {
83 const held = activity.loops.get(id)
84
85 if (held !== undefined) {
86 return held
87 }
88
89 if (activity.loops.size >= MAX_LOOPS) {
90 // Bounded: the loop not heard from longest makes room.
91 const oldest = [...activity.loops.values()].filter(loop => loop.id !== LEAD).sort((a, b) => a.lastAtMs - b.lastAtMs)[0]
92
93 if (oldest !== undefined) {
94 activity.loops.delete(oldest.id)
95 activity.recent.delete(oldest.id)
96 }
97 }
98
99 const loop: LoopRecord = { id, role: id === LEAD ? 'lead' : 'agent', status: 'running', calls: 0, errors: 0, lastAtMs: nowMs }
100
101 activity.loops.set(id, loop)
102
103 return loop
104}
105
106/** A subagent the engine started (`agent.spawn`'s answer), with the type and name it was asked for. */
107export function noteSpawn(activity: Activity, id: string, type: string, nowMs: number, name?: string, description?: string, adrs: readonly number[] = []): void {
108 const loop = loopFor(activity, id, nowMs)
109
110 loop.role = roleOf(type)
111 loop.status = 'running'
112 loop.lastAtMs = nowMs
113
114 if (name !== undefined) {
115 loop.name = name
116 }
117
118 if (description !== undefined) {
119 loop.description = description
120 }
121
122 if (adrs.length > 0) {
123 loop.adrs = adrs.slice(0, 8)
124 loop.description = `${loop.description ?? ''}${loop.description === undefined ? '' : ' '}[guided by ADR ${adrs.slice(0, 8).join(', ')}]`.slice(0, 200)
125 }
126}
127
128/** A tool call made by a loop (`agentId` absent: the main loop). Returns the pulse it lit, if any. */
129export function noteCall(activity: Activity, agentId: string | undefined, tool: string, subject: string, nowMs: number, command?: string): PulseKind | null {
130 const id = agentId ?? LEAD
131 const loop = loopFor(activity, id, nowMs)
132 const pulse = pulseOf(tool, command)
133
134 loop.calls += 1
135 loop.lastTool = tool
136 loop.lastSubject = subject
137 loop.lastAtMs = nowMs
138
139 if (pulse !== null) {
140 loop.pulse = { kind: pulse, atMs: nowMs }
141 }
142
143 return pulse
144}
145
146export function noteResult(activity: Activity, agentId: string | undefined, tool: string, subject: string, isError: boolean, nowMs: number): void {
147 const id = agentId ?? LEAD
148 const loop = loopFor(activity, id, nowMs)
149
150 if (isError) {
151 loop.errors += 1
152 }
153
154 activity.recent.set(id, [...(activity.recent.get(id) ?? []), { tool, subject, isError, atMs: nowMs }].slice(-RECENT))
155}
156
157/** How many failures in a row of one call make a loop stuck (ADR-477). */
158export const STUCK_AFTER = 3
159
160/** What a loop is called to the person, as its tile says it (a loop known only from its tool calls is said to be a subagent). */
161export function loopLabel(activity: Activity, id: string): string {
162 const loop = activity.loops.get(id)
163
164 if (id === LEAD) return 'claude (main)'
165
166 return loop === undefined ? `subagent …${id.slice(-4)}` : (loop.name ?? (loop.role === 'agent' ? `subagent …${id.slice(-4)}` : loop.role))
167}
168
169/** The call a loop has just failed STUCK_AFTER times running with the same subject (nothing in between), else null: the sign of a loop going round in circles. */
170export function stuckCall(activity: Activity, agentId: string | undefined): { tool: string; subject: string } | null {
171 const last = (activity.recent.get(agentId ?? LEAD) ?? []).slice(-STUCK_AFTER)
172 const [first] = last
173
174 if (first === undefined || last.length < STUCK_AFTER) return null
175
176 return last.every(call => call.isError && call.tool === first.tool && call.subject === first.subject) ? { tool: first.tool, subject: first.subject } : null
177}
178
179/** A loop's own `turn.complete`: a subagent's answer is in. */
180export function noteDone(activity: Activity, agentId: string, reason: string, nowMs: number): void {
181 const loop = loopFor(activity, agentId, nowMs)
182
183 loop.status = reason === 'answer' ? 'completed' : reason === 'aborted' ? 'killed' : 'failed'
184 loop.lastAtMs = nowMs
185}
186
187/** `$.agent.list()` as the engine answered it: types and statuses for loops the hooks saw start before this module did. */
188export function noteListed(activity: Activity, listed: readonly { id: string; type: string; status: string; description?: string; name?: string }[], nowMs: number): void {
189 for (const agent of listed.slice(0, MAX_LOOPS)) {
190 const loop = loopFor(activity, agent.id, nowMs)
191
192 loop.role = roleOf(agent.type)
193 loop.status = agent.status
194
195 // A named Agent call runs as an in-process teammate, listed as type `teammate`: its name is what tells it apart.
196 if (agent.name !== undefined && agent.name !== '') {
197 loop.name = agent.name.slice(0, 40)
198 }
199
200 if (agent.description !== undefined && loop.description === undefined) {
201 loop.description = agent.description
202 }
203 }
204}
205
206const RUFLO_STATES: Record<string, MemberState> = { idle: 'idle', busy: 'working', terminated: 'done' }
207const RUN_STATES: Partial<Record<RunEventType, MemberState>> = { 'run.started': 'working', 'run.output': 'working', 'run.completed': 'done', 'run.failed': 'failed', 'run.stopped': 'done' }
208const LOOP_STATES: Record<string, MemberState> = { running: 'working', pending: 'working', completed: 'done', failed: 'failed', killed: 'failed' }
209
210/**
211 * The tiles, in a stable order: the leader first, then ruflo's agents as the store lists them, then Claude Code's loops.
212 * A ruflo agent with a blocked claim is blocked; a loop that has not called a tool for a while is idle, not working.
213 */
214export function membersOf(snapshot: Snapshot | null, activity: Activity, nowMs: number): Member[] {
215 const claims = snapshot?.claims ?? []
216 const claimsOf = (id: string) => claims.filter(claim => claim.claimant.id === id && claim.status !== 'completed')
217 const queen = snapshot?.hive?.queen
218 const members: Member[] = []
219
220 if (queen !== undefined) {
221 members.push({ id: queen.id, label: 'queen', role: 'queen', source: 'queen', state: 'working', word: `term ${queen.term ?? '?'}`, isLeader: true, claims: 0, calls: 0 })
222 }
223
224 const coordinator = (snapshot?.agents ?? []).find(agent => /coordinator|queen/.test(agent.type) && agent.status !== 'terminated')
225
226 const events = snapshot?.ruos.events ?? []
227
228 for (const agent of snapshot?.agents ?? []) {
229 const held = claimsOf(agent.id)
230 const isBlocked = held.some(claim => claim.status === 'blocked')
231 const remote = agent.remote
232 // On a ruOS desktop the newest lifecycle event for its run (or the agent) says where it stands.
233 const run = remote === undefined ? undefined : events.filter(event => event.type !== 'desktop.state' && ((remote.runId !== undefined && event.runId === remote.runId) || event.agentId === agent.id)).sort((a, b) => b.ts - a.ts)[0]
234 const runState = run === undefined ? undefined : RUN_STATES[run.type]
235
236 members.push({
237 id: agent.id,
238 label: remote === undefined ? agent.type : `${agent.type} @${remote.desktopName}`,
239 role: agent.type,
240 source: 'ruflo',
241 state: isBlocked ? 'blocked' : (runState ?? RUFLO_STATES[agent.status] ?? 'idle'),
242 word: isBlocked ? 'blocked' : run !== undefined ? runWord(run) : agent.status,
243 isLeader: queen === undefined && coordinator?.id === agent.id,
244 claims: held.length,
245 calls: 0,
246 })
247 }
248
249 const loops = [...activity.loops.values()].sort((a, b) => (a.id === LEAD ? -1 : b.id === LEAD ? 1 : 0))
250
251 for (const loop of loops) {
252 const pulse = loop.pulse !== undefined && nowMs - loop.pulse.atMs < PULSE_MS ? loop.pulse.kind : undefined
253 const isLead = loop.id === LEAD
254 const isQuiet = nowMs - loop.lastAtMs > 30_000
255 const base = isLead ? (activity.isWorking ? 'working' : 'idle') : (LOOP_STATES[loop.status] ?? 'idle')
256
257 members.push({
258 id: loop.id,
259 // A loop known only from its tool calls (no spawn or listing reached the hooks) is said to be one, not given a role.
260 label: isLead ? 'claude (main)' : (loop.name ?? (loop.role === 'agent' ? `subagent …${loop.id.slice(-4)}` : loop.role)),
261 role: loop.role,
262 source: isLead ? 'lead' : 'claude',
263 state: base === 'working' && isQuiet && !isLead ? 'idle' : base,
264 word: isLead ? (activity.isWorking ? 'turn running' : 'waiting') : loop.status,
265 ...(pulse !== undefined && { pulse }),
266 isLeader: isLead && queen === undefined && coordinator === undefined,
267 claims: 0,
268 calls: loop.calls,
269 })
270 }
271
272 return members
273}
274
275/** True while some tile is still lit: the frame timer keeps drawing until every pulse has faded. */
276export const isAnimating = (activity: Activity, nowMs: number): boolean =>
277 [...activity.loops.values()].some(loop => loop.pulse !== undefined && nowMs - loop.pulse.atMs < PULSE_MS)
278hooks/reader/parse.ts 365 lines1/**
2 * Parsers for the files ruflo writes under `.claude-flow/` and `.swarm/`, and for the JSON `hooks route` prints.
3 * Every one takes text from disk that another process wrote, so each tolerates any shape: what it cannot read
4 * is left out, never guessed. Nothing here reads, keeps or returns the hive's `hiveToken`.
5 */
6import { remoteHostOf, type RemoteHost } from './ruos'
7
8/** Text longer than this is not read: a store that size is not one the CLI wrote, and parsing it would stall a hook. */
9export const MAX_TEXT = 4_000_000
10/** At most this many records of one kind are kept; the rest are counted, not drawn. */
11export const MAX_RECORDS = 1_000
12
13export type SwarmInfo = {
14 id: string
15 topology: string
16 status: string
17 maxAgents?: number
18 strategy?: string
19 consensus?: string
20 agentIds: string[]
21 updatedAt?: string
22}
23
24export type AgentRecord = {
25 id: string
26 type: string
27 status: string
28 health?: number
29 taskCount?: number
30 model?: string
31 createdAt?: string
32 /** Where the agent runs, when ruflo-ruos placed it on a ruOS desktop (`config.host`, see ./ruos). */
33 remote?: RemoteHost
34}
35
36export type TaskRecord = {
37 id: string
38 type: string
39 description: string
40 priority: string
41 status: string
42 progress?: number
43 assignedTo: string[]
44}
45
46export type ClaimRecord = {
47 issueId: string
48 status: string
49 claimant: { kind: 'agent' | 'human'; id: string; agentType?: string }
50 progress?: number
51 handoffTo?: string
52 isStealable: boolean
53}
54
55export type Proposal = { id: string; type: string; status: string; strategy: string; votesFor: number; votesAgainst: number; proposedBy: string }
56export type Decision = { id: string; type: string; result: string; votesFor: number; votesAgainst: number; decidedAt: string }
57
58export type HiveInfo = {
59 topology: string
60 strategy?: string
61 queen?: { id: string; term?: number }
62 workers: string[]
63 pending: Proposal[]
64 history: Decision[]
65}
66
67export type RoutePick = {
68 task: string
69 agent: string
70 confidence: number
71 matched: boolean
72 pattern?: string
73 alternatives: { agent: string; confidence: number }[]
74 method?: string
75 atMs: number
76}
77
78const ID = /^[A-Za-z0-9][A-Za-z0-9._:@-]{0,127}$/
79
80/** Plain printable text of at most `max` characters: no control characters reach the terminal. */
81export function plain(value: unknown, max = 200): string {
82 if (typeof value !== 'string') {
83 return ''
84 }
85
86 const cleaned = value.replace(/[\u0000-\u001f\u007f-\u009f---]/g, ' ').replace(/\s+/g, ' ').trim()
87
88 return cleaned.length <= max ? cleaned : `${cleaned.slice(0, Math.max(0, max - 1))}…`
89}
90
91/** An id as ruflo mints them (`agent-…`, `task-…`, `proposal-…`), or null: only such a string ever reaches an argv. */
92export function idOf(value: unknown): string | null {
93 return typeof value === 'string' && ID.test(value) ? value : null
94}
95
96const numberOf = (value: unknown): number | undefined => (typeof value === 'number' && Number.isFinite(value) ? value : undefined)
97const stringOf = (value: unknown, max = 80): string | undefined => (typeof value === 'string' && value !== '' ? plain(value, max) : undefined)
98const recordOf = (value: unknown): Record<string, unknown> | null =>
99 value !== null && typeof value === 'object' && !Array.isArray(value) ? (value as Record<string, unknown>) : null
100
101/** JSON text to a plain object, or null for anything else (too long, malformed, an array, a scalar). */
102export function jsonObject(text: string | null): Record<string, unknown> | null {
103 if (text === null || text.length > MAX_TEXT) {
104 return null
105 }
106
107 try {
108 return recordOf(JSON.parse(text))
109 } catch {
110 return null
111 }
112}
113
114const valuesOf = (value: unknown): unknown[] => {
115 const record = recordOf(value)
116
117 return record === null ? [] : Object.values(record).slice(0, MAX_RECORDS)
118}
119
120/** `.claude-flow/swarm/swarm-state.json`: the running swarm, else the one updated last. */
121export function parseSwarmStore(text: string | null): SwarmInfo | null {
122 const swarms = valuesOf(jsonObject(text)?.swarms)
123 .map(recordOf)
124 .flatMap(swarm => {
125 const id = idOf(swarm?.swarmId)
126
127 if (swarm === null || id === null) {
128 return []
129 }
130
131 const config = recordOf(swarm.config)
132
133 return [
134 {
135 id,
136 topology: stringOf(swarm.topology, 40) ?? 'unknown',
137 status: stringOf(swarm.status, 40) ?? 'unknown',
138 ...(numberOf(swarm.maxAgents) !== undefined && { maxAgents: numberOf(swarm.maxAgents) }),
139 ...(stringOf(config?.strategy, 40) !== undefined && { strategy: stringOf(config?.strategy, 40) }),
140 ...(stringOf(config?.consensusMechanism, 40) !== undefined && { consensus: stringOf(config?.consensusMechanism, 40) }),
141 agentIds: (Array.isArray(swarm.agents) ? swarm.agents : []).slice(0, MAX_RECORDS).flatMap(agent => (idOf(agent) !== null ? [agent as string] : [])),
142 ...(stringOf(swarm.updatedAt, 40) !== undefined && { updatedAt: stringOf(swarm.updatedAt, 40) }),
143 } satisfies SwarmInfo,
144 ]
145 })
146
147 const byRecency = (a: SwarmInfo, b: SwarmInfo) => (b.updatedAt ?? '').localeCompare(a.updatedAt ?? '')
148
149 return [...swarms.filter(swarm => swarm.status === 'running')].sort(byRecency)[0] ?? [...swarms].sort(byRecency)[0] ?? null
150}
151
152/** `.swarm/state.json`: the pointer `swarm init` leaves, which names the swarm and its strategy. */
153export function parseSwarmPointer(text: string | null): { id: string; topology?: string; strategy?: string; status?: string } | null {
154 const value = jsonObject(text)
155 const id = idOf(value?.id)
156
157 if (value === null || id === null) {
158 return null
159 }
160
161 return {
162 id,
163 ...(stringOf(value.topology, 40) !== undefined && { topology: stringOf(value.topology, 40) }),
164 ...(stringOf(value.strategy, 40) !== undefined && { strategy: stringOf(value.strategy, 40) }),
165 ...(stringOf(value.status, 40) !== undefined && { status: stringOf(value.status, 40) }),
166 }
167}
168
169/** `.claude-flow/agents/store.json`. */
170export function parseAgents(text: string | null): AgentRecord[] {
171 return valuesOf(jsonObject(text)?.agents).flatMap(entry => {
172 const agent = recordOf(entry)
173 const id = idOf(agent?.agentId)
174
175 if (agent === null || id === null) {
176 return []
177 }
178
179 return [
180 {
181 id,
182 type: stringOf(agent.agentType, 40) ?? 'agent',
183 status: stringOf(agent.status, 20) ?? 'unknown',
184 ...(numberOf(agent.health) !== undefined && { health: numberOf(agent.health) }),
185 ...(numberOf(agent.taskCount) !== undefined && { taskCount: numberOf(agent.taskCount) }),
186 ...(stringOf(agent.model, 40) !== undefined && { model: stringOf(agent.model, 40) }),
187 ...(stringOf(agent.createdAt, 40) !== undefined && { createdAt: stringOf(agent.createdAt, 40) }),
188 ...(remoteHostOf(agent.config) !== null && { remote: remoteHostOf(agent.config) as RemoteHost }),
189 },
190 ]
191 })
192}
193
194/** `.claude-flow/tasks/store.json`. */
195export function parseTasks(text: string | null): TaskRecord[] {
196 return valuesOf(jsonObject(text)?.tasks).flatMap(entry => {
197 const task = recordOf(entry)
198 const id = idOf(task?.taskId)
199
200 if (task === null || id === null) {
201 return []
202 }
203
204 return [
205 {
206 id,
207 type: stringOf(task.type, 40) ?? 'task',
208 description: plain(task.description, 300),
209 priority: stringOf(task.priority, 20) ?? 'normal',
210 status: stringOf(task.status, 20) ?? 'unknown',
211 ...(numberOf(task.progress) !== undefined && { progress: numberOf(task.progress) }),
212 assignedTo: (Array.isArray(task.assignedTo) ? task.assignedTo : []).slice(0, 50).flatMap(agent => (idOf(agent) !== null ? [agent as string] : [])),
213 },
214 ]
215 })
216}
217
218/** `.claude-flow/claims/claims.json`: issue claims (who works on what), not the authorization file `.claude-flow/claims.json`. */
219export function parseClaims(text: string | null): ClaimRecord[] {
220 const store = jsonObject(text)
221 const stealable = recordOf(store?.stealable) ?? {}
222
223 return valuesOf(store?.claims).flatMap(entry => {
224 const claim = recordOf(entry)
225 const issueId = idOf(claim?.issueId)
226 const claimant = recordOf(claim?.claimant)
227
228 if (claim === null || issueId === null || claimant === null) {
229 return []
230 }
231
232 const isAgent = claimant.type === 'agent'
233 const who = idOf(isAgent ? claimant.agentId : claimant.userId)
234
235 if (who === null) {
236 return []
237 }
238
239 const handoff = recordOf(claim.handoffTo)
240 const handoffTo = idOf(handoff?.agentId ?? handoff?.userId)
241
242 return [
243 {
244 issueId,
245 status: stringOf(claim.status, 30) ?? 'unknown',
246 claimant: { kind: isAgent ? 'agent' : 'human', id: who, ...(stringOf(claimant.agentType, 40) !== undefined && { agentType: stringOf(claimant.agentType, 40) }) },
247 ...(numberOf(claim.progress) !== undefined && { progress: numberOf(claim.progress) }),
248 ...(handoffTo !== null && { handoffTo }),
249 isStealable: claim.status === 'stealable' || Object.hasOwn(stealable, issueId),
250 },
251 ]
252 })
253}
254
255const votesOf = (value: unknown): { votesFor: number; votesAgainst: number } => {
256 const votes = valuesOf(value)
257
258 return { votesFor: votes.filter(vote => vote === true).length, votesAgainst: votes.filter(vote => vote === false).length }
259}
260
261/** `.claude-flow/hive-mind/state.json`, without its capability token. */
262export function parseHive(text: string | null): HiveInfo | null {
263 const hive = jsonObject(text)
264
265 if (hive === null || hive.initialized !== true) {
266 return null
267 }
268
269 const queen = recordOf(hive.queen)
270 const queenId = idOf(queen?.agentId)
271 const consensus = recordOf(hive.consensus)
272
273 const pending = (Array.isArray(consensus?.pending) ? consensus.pending : []).slice(-50).flatMap(entry => {
274 const proposal = recordOf(entry)
275 const id = idOf(proposal?.proposalId)
276
277 return proposal === null || id === null
278 ? []
279 : [
280 {
281 id,
282 type: stringOf(proposal.type, 40) ?? 'proposal',
283 status: stringOf(proposal.status, 20) ?? 'pending',
284 strategy: stringOf(proposal.strategy, 20) ?? 'unknown',
285 proposedBy: stringOf(proposal.proposedBy, 60) ?? 'unknown',
286 ...votesOf(proposal.votes),
287 },
288 ]
289 })
290
291 const history = (Array.isArray(consensus?.history) ? consensus.history : []).slice(-50).flatMap(entry => {
292 const decision = recordOf(entry)
293 const id = idOf(decision?.proposalId)
294 const votes = recordOf(decision?.votes)
295
296 return decision === null || id === null
297 ? []
298 : [
299 {
300 id,
301 type: stringOf(decision.type, 40) ?? 'proposal',
302 result: stringOf(decision.result, 20) ?? 'unknown',
303 votesFor: numberOf(votes?.for) ?? 0,
304 votesAgainst: numberOf(votes?.against) ?? 0,
305 decidedAt: stringOf(decision.decidedAt, 40) ?? '',
306 },
307 ]
308 })
309
310 return {
311 topology: stringOf(hive.topology, 40) ?? 'unknown',
312 ...(stringOf(hive.consensusStrategy, 30) !== undefined && { strategy: stringOf(hive.consensusStrategy, 30) }),
313 ...(queenId !== null && { queen: { id: queenId, ...(numberOf(queen?.term) !== undefined && { term: numberOf(queen?.term) }) } }),
314 workers: (Array.isArray(hive.workers) ? hive.workers : []).slice(0, MAX_RECORDS).flatMap(worker => (idOf(worker) !== null ? [worker as string] : [])),
315 pending,
316 history,
317 }
318}
319
320/**
321 * What `hooks route --format json` printed (log lines may come first), or the `hooks_route` MCP tool answered, as one pick.
322 * The router's own `matched` flag is kept as it said it; whether the confidence clears a threshold is the view's call.
323 */
324export function parseRoute(text: string, atMs: number): RoutePick | null {
325 const start = text.indexOf('{')
326 const value = start < 0 ? null : jsonObject(text.slice(start, Math.min(text.length, start + 200_000)))
327 const primary = recordOf(value?.primaryAgent)
328 const agent = stringOf(primary?.type, 40)
329 const confidence = numberOf(primary?.confidence)
330
331 if (value === null || agent === undefined || confidence === undefined) {
332 return null
333 }
334
335 const routing = recordOf(value.routing)
336
337 return {
338 task: plain(value.task, 200),
339 agent,
340 confidence: Math.max(0, Math.min(1, confidence)),
341 matched: value.matched === true,
342 ...(stringOf(value.matchedPattern, 40) !== undefined && { pattern: stringOf(value.matchedPattern, 40) }),
343 alternatives: (Array.isArray(value.alternativeAgents) ? value.alternativeAgents : []).slice(0, 3).flatMap(entry => {
344 const alt = recordOf(entry)
345 const type = stringOf(alt?.type, 40)
346 const score = numberOf(alt?.confidence)
347
348 return type !== undefined && score !== undefined ? [{ agent: type, confidence: Math.max(0, Math.min(1, score)) }] : []
349 }),
350 ...(stringOf(routing?.method, 40) !== undefined && { method: stringOf(routing?.method, 40) }),
351 atMs,
352 }
353}
354
355/** A pick held in `$.store` read back: every field checked, as anything else from storage is. */
356export function routeFromStore(value: unknown): RoutePick | null {
357 const pick = recordOf(value)
358
359 if (pick === null || typeof pick.task !== 'string' || typeof pick.agent !== 'string' || numberOf(pick.confidence) === undefined || numberOf(pick.atMs) === undefined) {
360 return null
361 }
362
363 return parseRoute(JSON.stringify({ task: pick.task, matched: pick.matched, matchedPattern: pick.pattern, primaryAgent: { type: pick.agent, confidence: pick.confidence }, alternativeAgents: [], routing: { method: pick.method } }), numberOf(pick.atMs) ?? 0)
364}
365hooks/reader/snapshot.ts 138 lines1import {
2 parseAgents,
3 parseClaims,
4 parseHive,
5 parseSwarmPointer,
6 parseSwarmStore,
7 parseTasks,
8 type AgentRecord,
9 type ClaimRecord,
10 type HiveInfo,
11 type SwarmInfo,
12 type TaskRecord,
13} from './parse'
14import { EVENTS_MAX, parseEvents, parseHosts, RUOS_EVENTS, RUOS_HOSTS, type RunEvent, type RuosHost } from './ruos'
15
16/** The file calls a reader makes, as `$.fs` answers them; every one may be refused. */
17export type ReaderFs = {
18 read: (path: string) => Promise<string>
19 stat: (path: string) => Promise<{ mtimeMs?: number; size?: number } | undefined>
20}
21
22/** Where ruflo keeps each fact, relative to the session's working directory. */
23export const FILES = {
24 swarm: '.claude-flow/swarm/swarm-state.json',
25 pointer: '.swarm/state.json',
26 agents: '.claude-flow/agents/store.json',
27 tasks: '.claude-flow/tasks/store.json',
28 claims: '.claude-flow/claims/claims.json',
29 hive: '.claude-flow/hive-mind/state.json',
30} as const
31
32export type FileKey = keyof typeof FILES
33
34export type Snapshot = {
35 swarm: SwarmInfo | null
36 agents: AgentRecord[]
37 tasks: TaskRecord[]
38 claims: ClaimRecord[]
39 hive: HiveInfo | null
40 /** ruOS desktops that host swarm agents (ruflo-ruos): null hosts when the snapshot file is absent. */
41 ruos: { hosts: RuosHost[] | null; events: RunEvent[]; note?: string }
42 /** The facts that are not on disk (no file, unreadable, or not in a shape ruflo writes), by key. */
43 missing: FileKey[]
44 /** True when any of ruflo's swarm files is there at all: the pane has something of its own to show. */
45 hasSwarm: boolean
46 readAtMs: number
47}
48
49/** The text of each file as last read, by path, with the mtime it was read at: an unchanged file is not read again. */
50export type ReadCache = Map<string, { mtimeMs: number; size: number; text: string }>
51
52async function textOf(fs: ReaderFs, cache: ReadCache, path: string): Promise<string | null> {
53 try {
54 const stat = await fs.stat(path)
55 const mtimeMs = stat?.mtimeMs ?? -1
56 const size = stat?.size ?? -1
57 const held = cache.get(path)
58
59 if (held !== undefined && mtimeMs >= 0 && held.mtimeMs === mtimeMs && held.size === size) {
60 return held.text
61 }
62
63 const text = await fs.read(path)
64
65 cache.set(path, { mtimeMs, size, text })
66
67 return text
68 } catch {
69 cache.delete(path)
70
71 return null
72 }
73}
74
75/** Reads every swarm fact ruflo keeps on disk, in parallel; a refused or absent file is a missing fact, never an error. */
76export async function readSnapshot(fs: ReaderFs, cache: ReadCache, nowMs: number): Promise<Snapshot> {
77 const keys = Object.keys(FILES) as FileKey[]
78 const texts = await Promise.all(keys.map(key => textOf(fs, cache, FILES[key])))
79 const text = (key: FileKey) => texts[keys.indexOf(key)] ?? null
80
81 const pointer = parseSwarmPointer(text('pointer'))
82 const stored = parseSwarmStore(text('swarm'))
83 // The pointer names the swarm `swarm init` made last; the store has its members. Either alone still names a swarm.
84 const swarm: SwarmInfo | null =
85 stored ??
86 (pointer !== null
87 ? { id: pointer.id, topology: pointer.topology ?? 'unknown', status: pointer.status ?? 'unknown', agentIds: [], ...(pointer.strategy !== undefined && { strategy: pointer.strategy }) }
88 : null)
89
90 const ruos = await readRuos(fs, cache)
91 const parsed = {
92 swarm,
93 agents: parseAgents(text('agents')),
94 tasks: parseTasks(text('tasks')),
95 claims: parseClaims(text('claims')),
96 hive: parseHive(text('hive')),
97 ruos,
98 }
99
100 const missing = keys.filter(key => {
101 switch (key) {
102 case 'swarm':
103 return stored === null
104 case 'pointer':
105 return pointer === null
106 case 'hive':
107 return parsed.hive === null
108 default:
109 return text(key) === null
110 }
111 })
112
113 return {
114 ...parsed,
115 missing,
116 hasSwarm: swarm !== null || parsed.agents.length > 0 || parsed.hive !== null,
117 readAtMs: nowMs,
118 }
119}
120
121/** The ruOS host snapshot and the tail of its event log. The log is read whole (there is no ranged read), so past a size it is not read. */
122async function readRuos(fs: ReaderFs, cache: ReadCache): Promise<Snapshot['ruos']> {
123 const hosts = parseHosts(await textOf(fs, cache, RUOS_HOSTS))
124 let size = -1
125
126 try {
127 size = (await fs.stat(RUOS_EVENTS))?.size ?? -1
128 } catch {
129 return { hosts, events: [] }
130 }
131
132 if (size > EVENTS_MAX) {
133 return { hosts, events: [], note: `ruOS event log is ${Math.round(size / 1_000_000)} MB: not read` }
134 }
135
136 return { hosts, events: parseEvents(await textOf(fs, cache, RUOS_EVENTS)) }
137}
138hooks/state.ts 135 lines1import type { PluginOptions, Timer } from 'claude-code'
2
3import type { ActionSpec } from './actions/argv'
4import type { Activity } from './model/members'
5import { routeFromStore, type RoutePick } from './reader/parse'
6import type { Snapshot } from './reader/snapshot'
7
8export const PANE_ID = 'ruflo-swarm'
9
10/**
11 * The `$.store` key the few facts that must outlive a hot reload sit under. `$.store` is the plugin's, not the folder's,
12 * so the key carries the working directory: a router pick or a selection never follows the person into another project.
13 */
14export const storeKeyOf = (cwd: string): string => `ruflo-swarm/ui:${cwd}`
15
16/** How the pane reaches the ruflo CLI. Each is a fixed argv prefix; only `npx` may touch the network. */
17export const CLI_PREFIXES = {
18 'npx-offline': ['npx', '--offline', '-y', '@claude-flow/cli@latest'],
19 npx: ['npx', '-y', '@claude-flow/cli@latest'],
20 ruflo: ['ruflo'],
21 'claude-flow': ['claude-flow'],
22} as const satisfies Record<string, readonly string[]>
23
24export type CliChoice = keyof typeof CLI_PREFIXES
25
26export type Options = {
27 /** `auto` opens the pane once a swarm is on disk and the layout docks panes; `command` only on request; `off` never. */
28 panel: 'auto' | 'command' | 'off'
29 cli: CliChoice
30 /** A router pick under this confidence is drawn as "below threshold", with its score. */
31 routeThreshold: number
32 /** Appends a short swarm note to the prompt of each subagent Claude Code spawns. */
33 injectSpawnContext: boolean
34 /** Records a bounded, content-free trail of engine events into ruflo memory. */
35 audit: boolean
36 /** Appends the ADRs attached to the active ruflo-console mission to the prompt of each subagent (default on; nothing is added without an attached ADR). */
37 injectAdrs: boolean
38}
39
40const PANELS = new Set(['auto', 'command', 'off'])
41
42/** The options as the settings hold them, each one checked: a value the plugin does not know is its default. */
43export function optionsOf(raw: PluginOptions): Options {
44 const value = (raw ?? {}) as Record<string, unknown>
45 // `command` by default since ruflo-console: its cockpit is the pane that opens by itself, and two would crowd the dock.
46 const panel = typeof value.panel === 'string' && PANELS.has(value.panel) ? (value.panel as Options['panel']) : 'command'
47 const cli = typeof value.cli === 'string' && value.cli in CLI_PREFIXES ? (value.cli as CliChoice) : 'npx-offline'
48 const threshold = typeof value.routeThreshold === 'number' ? value.routeThreshold : Number(value.routeThreshold)
49
50 return {
51 panel,
52 cli,
53 routeThreshold: Number.isFinite(threshold) && threshold >= 0 && threshold <= 1 ? threshold : 0.5,
54 injectSpawnContext: value.injectSpawnContext === true,
55 audit: value.audit === true,
56 injectAdrs: value.injectAdrs !== false,
57 }
58}
59
60/** A destructive or mutating action waiting for the person's second press. */
61export type PendingConfirm = { label: string; spec: ActionSpec; askedAtMs: number }
62
63/** What an action did, as the pane says it: what ran, how it exited, and whether the disk shows the change. */
64export type ActionOutcome = { label: string; ok: boolean; verified: 'yes' | 'no' | 'n/a'; detail: string; atMs: number }
65
66export type Usage = { costUsd?: number; contextTokens?: number; contextPercent?: number; contextWindow?: number; readAtMs: number }
67
68export type State = {
69 options: Options
70 cwd: string
71 snapshot: Snapshot | null
72 readError: string | null
73 isRefreshing: boolean
74 isRefreshQueued: boolean
75 activity: Activity
76 usage: Usage | null
77 route: RoutePick | null
78 pane: { isOpen: boolean; isClosedByPerson: boolean; columns: number; rows: number }
79 viewport: { isFullscreen?: boolean }
80 /** The tile the buttons act on, by member id: survives a re-read that reorders the tiles. */
81 selected: string | null
82 selectedTask: string | null
83 confirm: PendingConfirm | null
84 outcome: ActionOutcome | null
85 isActing: boolean
86 /** Lines an action asked to show (an agent's logs), with whose they are. */
87 detail: { title: string; lines: string[] } | null
88 timers: Map<string, Timer>
89 audit: { buffer: string[]; dropped: number; flushedAtMs: number }
90}
91
92export function newState(raw: PluginOptions, activity: Activity): State {
93 return {
94 options: optionsOf(raw),
95 cwd: '',
96 snapshot: null,
97 readError: null,
98 isRefreshing: false,
99 isRefreshQueued: false,
100 activity,
101 usage: null,
102 route: null,
103 pane: { isOpen: false, isClosedByPerson: false, columns: 0, rows: 0 },
104 viewport: {},
105 selected: null,
106 selectedTask: null,
107 confirm: null,
108 outcome: null,
109 isActing: false,
110 detail: null,
111 timers: new Map(),
112 audit: { buffer: [], dropped: 0, flushedAtMs: 0 },
113 }
114}
115
116/** The part of the state written to `$.store`, so a hot reload keeps what the person chose. */
117export type Persisted = { selected: string | null; selectedTask: string | null; isClosedByPerson: boolean; route: RoutePick | null }
118
119export function persistedOf(state: State): Persisted {
120 return { selected: state.selected, selectedTask: state.selectedTask, isClosedByPerson: state.pane.isClosedByPerson, route: state.route }
121}
122
123export function restore(state: State, value: unknown): void {
124 if (value === null || typeof value !== 'object') {
125 return
126 }
127
128 const held = value as Partial<Persisted>
129
130 state.selected = typeof held.selected === 'string' ? held.selected : null
131 state.selectedTask = typeof held.selectedTask === 'string' ? held.selectedTask : null
132 state.pane.isClosedByPerson = held.isClosedByPerson === true
133 state.route = routeFromStore(held.route)
134}
135hooks/toast-policy.ts 462 lines1/**
2 * Toast policy (ADR-477): levels, one-line washing, de-duplication, a per-source rate limit, the person's setting, and digests for
3 * the console's Events page. One CANONICAL source, plugins/ruflo-mods/hooks/toast/policy.ts; every other plugin carries a byte-identical
4 * copy at hooks/toast-policy.ts, written and checked by scripts/sync-toast-policy.mjs (a plugin ships alone through the marketplace and
5 * cannot import a sibling at run time). Edit the canonical file, run `node scripts/sync-toast-policy.mjs`, never a copy.
6 *
7 * Dependency-free and engine-free: no `$`, no import. The engine is reached only through the functions a plugin hands in, so every
8 * rule here is a plain function a test can drive with a fake clock. Nothing in this file throws to its caller: a refused toast, an
9 * unreadable setting or a failed write is a result, never a crash.
10 */
11
12export type ToastLevel = 'info' | 'ok' | 'warn' | 'error'
13/** all: every level; important: warn and error, plus a toast marked `always`; off: none (all are still recorded). */
14export type ToastMode = 'all' | 'important' | 'off'
15export type ToastPrefs = { mode: ToastMode; muted: readonly string[] }
16/** Why a toast was or was not drawn. `shown` is the only one that drew. `coalesced` is an error held for one later line. */
17export type Why = 'shown' | 'deduped' | 'rate-limited' | 'coalesced' | 'muted' | 'off' | 'filtered' | 'away' | 'refused'
18
19export const TOAST_SOURCES = ['console', 'swarm', 'protector', 'mods'] as const
20export const TOAST_MODES: readonly ToastMode[] = ['all', 'important', 'off']
21export const LEVELS: readonly ToastLevel[] = ['info', 'ok', 'warn', 'error']
22export const DEFAULT_PREFS: ToastPrefs = { mode: 'all', muted: [] }
23export const PREFIX: Readonly<Record<ToastLevel, string>> = { info: '›', ok: '✓', warn: '⚠', error: '✗' }
24
25/** A whole toast line, prefix included, is never longer than this. */
26export const LINE_MAX = 120
27/** An identical (source, text) is not drawn again inside this window. */
28export const DEDUPE_MS = 60_000
29/** A source draws at most RATE_MAX toasts in RATE_WINDOW_MS; an error past that is held and said once, with a count. */
30export const RATE_MAX = 4
31export const RATE_WINDOW_MS = 60_000
32/** Digests kept per source (a ring), and how long the setting is believed before the file is read again. */
33export const RING_MAX = 60
34export const PREFS_TTL_MS = 4_000
35export const CONSOLE_TTL_MS = 30_000
36/** The setting (written by the console's Settings) and the folder of per-source digest files (read by the console). */
37export const CONSOLE_DIR = '.claude-flow/console'
38export const PREFS_FILE = `${CONSOLE_DIR}/toast-prefs.json`
39export const TOAST_DIR = `${CONSOLE_DIR}/toasts`
40export const MASK = '‹masked›'
41
42// ------------------------------------------------------------------------------------------------------------ washing
43
44const ESCAPES = new RegExp('\\u001b\\][^\\u0007\\u001b]*(?:\\u0007|\\u001b\\\\)|\\u009d[^\\u0007\\u009c]*[\\u0007\\u009c]|(?:\\u001b\\[|\\u009b)[0-9;?]*[ -/]*[@-~]', 'g')
45/** White space of every kind (and the line and paragraph separators) becomes one space. */
46const SPACES = new RegExp('[\\s\\u0085\\u2028\\u2029]+', 'g')
47/** A control, zero-width, bidi or tag character is deleted, not spaced: `sk-ant-AAAA<NUL>BBBB` is one credential. */
48const DROPPED = new RegExp('[\\u0000-\\u0008\\u000e-\\u001f\\u007f-\\u0084\\u0086-\\u009f\\u00ad\\u034f\\u061c\\u115f\\u1160\\u17b4\\u17b5\\u180b-\\u180f\\u200b-\\u200f\\u202a-\\u202e\\u2060-\\u206f\\u3164\\ufe00-\\ufe0f\\ufeff\\uffa0\\ufff9-\\ufffb]|[\\u{e0000}-\\u{e0fff}]', 'gu')
49const SECRETISH = new RegExp(
50 [
51 String.raw`\b(?:sk|pk|ghp|gho|ghs|github_pat|xox[abprs]|xapp|AKIA|ASIA|AIza)[-_A-Za-z0-9]{12,}`,
52 String.raw`\b(?:glpat|npm|hf|dop_v1|shpat|whsec|rk_live|sk_live|ya29)[-_.][-_.A-Za-z0-9]{12,}`,
53 String.raw`\bBearer\s+\S{8,}`,
54 String.raw`\beyJ[A-Za-z0-9_-]{6,}\.[A-Za-z0-9_-]{6,}\.?[A-Za-z0-9_-]*`,
55 String.raw`\b[A-Za-z0-9+_-]{32,}={0,2}`,
56 String.raw`-----BEGIN [A-Z ]*PRIVATE KEY-----[\s\S]*?(?:-----END [A-Z ]*PRIVATE KEY-----|$)`,
57 String.raw`\b[a-z][a-z0-9+.-]*://[^\s/:@]+:[^\s/@]+@`,
58 String.raw`(?:key|token|secret|passw(?:or)?d|pwd|passphrase|credential|authorization|cookie)["']?\s*[=:]\s*(?:(?:Bearer|Basic|Token)\s+)?(?:"[^"]*"|'[^']*'|\S+)`,
59 String.raw`(?<![A-Za-z0-9])pass["']?\s*=\s*(?:"[^"]*"|'[^']*'|\S+)`,
60 String.raw`(?:^|\s)--?(?:token|password|passwd|pwd|secret|api-?key|auth(?:orization)?|access-?key|client-?secret)(?:=|\s+)\S+`,
61 ].join('|'),
62 'gi',
63)
64const HOME_PATH = /\/(?:home|Users)\/[^/\s'"]+/g
65const EMAIL = /\b[\w.+-]{1,64}@[A-Za-z0-9-]{1,63}(?:\.[A-Za-z0-9-]{1,63})+\b/g
66
67/** One line of at most `max` characters: escapes and control characters gone, white space collapsed, credentials, e-mail addresses and home paths masked, an ellipsis where it was cut. Never throws; a non-string is ''. */
68export function tidy(value: unknown, max: number = LINE_MAX): string {
69 if (typeof value !== 'string') return ''
70
71 const washed = value.slice(0, 4096).replace(ESCAPES, '').replace(SPACES, ' ').replace(DROPPED, '').trim()
72 const masked = washed.replace(SECRETISH, MASK).replace(EMAIL, MASK).replace(HOME_PATH, '~')
73
74 return masked.length > max ? `${masked.slice(0, Math.max(0, max - 1))}…` : masked
75}
76
77/** The line a toast draws: its level's prefix, a space, the washed text; the whole is at most LINE_MAX. '' when there is no text. */
78export function lineOf(level: ToastLevel, text: unknown): string {
79 const body = tidy(text, LINE_MAX - 2)
80
81 return body === '' ? '' : `${PREFIX[level]} ${body}`
82}
83
84// ------------------------------------------------------------------------------------------------------------ the setting
85
86const SOURCE_NAME = /^[a-z][a-z0-9-]{0,23}$/
87
88/** The setting from its file's text: anything unreadable, or an unknown mode, is the default (all, nothing muted). At most 8 names are muted. */
89export function parsePrefs(text: string | null | undefined): ToastPrefs {
90 if (typeof text !== 'string' || text.length > 4096) return DEFAULT_PREFS
91
92 try {
93 const o: unknown = JSON.parse(text)
94
95 if (typeof o !== 'object' || o === null || Array.isArray(o)) return DEFAULT_PREFS
96
97 const r = o as { mode?: unknown; muted?: unknown }
98 const mode = TOAST_MODES.find(candidate => candidate === r.mode) ?? 'all'
99 const muted = Array.isArray(r.muted) ? [...new Set(r.muted.filter((name): name is string => typeof name === 'string' && SOURCE_NAME.test(name)))].slice(0, 8) : []
100
101 return { mode, muted }
102 } catch {
103 return DEFAULT_PREFS
104 }
105}
106
107export const encodePrefs = (prefs: ToastPrefs): string => `${JSON.stringify({ v: 1, mode: prefs.mode, muted: [...new Set(prefs.muted)].filter(name => SOURCE_NAME.test(name)).slice(0, 8) })}\n`
108
109// ------------------------------------------------------------------------------------------------------------ digests
110
111/** What is kept of every toast, drawn or not: masked, short, and flagged with what became of it. */
112export type Digest = { t: number; source: string; level: ToastLevel; text: string; shown: boolean; why: Why; /** How many identical, consecutive ones this stands for (absent: one). */ n?: number }
113
114const WHYS: readonly Why[] = ['shown', 'deduped', 'rate-limited', 'coalesced', 'muted', 'off', 'filtered', 'away', 'refused']
115
116/** Adds a digest to a ring (newest last). An identical neighbour (source, text, outcome) is counted, not repeated. Mutates `ring`. */
117export function pushDigest(ring: Digest[], d: Digest, cap: number = RING_MAX): void {
118 const last = ring[ring.length - 1]
119
120 if (last !== undefined && last.source === d.source && last.text === d.text && last.why === d.why && last.level === d.level) {
121 last.n = (last.n ?? 1) + 1
122 last.t = d.t
123 } else ring.push({ ...d })
124
125 if (ring.length > cap) ring.splice(0, ring.length - cap)
126}
127
128export const encodeRing = (ring: readonly Digest[]): string =>
129 ring.map(d => `${JSON.stringify({ v: 1, t: Math.round(d.t), s: d.source, l: d.level, x: d.text, w: d.why, ...(d.n !== undefined && d.n > 1 && { n: d.n }) })}\n`).join('')
130
131/** The digests in a ring file's text; a line that is not ours, or is cut off, is skipped. Text is washed again: a file is not trusted. */
132export function decodeRing(text: string | null | undefined): Digest[] {
133 if (typeof text !== 'string' || text === '') return []
134
135 const out: Digest[] = []
136
137 for (const line of text.slice(-300_000).split('\n')) {
138 if (line.length < 2 || line.length > 600 || line[0] !== '{') continue
139
140 try {
141 const o = JSON.parse(line) as Record<string, unknown>
142 const level = LEVELS.find(candidate => candidate === o.l)
143 const why = WHYS.find(candidate => candidate === o.w)
144 const source = typeof o.s === 'string' && SOURCE_NAME.test(o.s) ? o.s : undefined
145 const body = tidy(o.x, LINE_MAX)
146
147 if (o.v !== 1 || typeof o.t !== 'number' || !Number.isFinite(o.t) || level === undefined || why === undefined || source === undefined || body === '') continue
148 out.push({ t: o.t, source, level, text: body, shown: why === 'shown', why, ...(typeof o.n === 'number' && o.n > 1 && o.n < 1e6 && { n: Math.floor(o.n) }) })
149 } catch {
150 /* a half-written line */
151 }
152 }
153
154 return out
155}
156
157// ------------------------------------------------------------------------------------------------------------ the toaster
158
159type Maybe<T> = T | Promise<T>
160
161const isThenable = (value: unknown): value is Promise<unknown> => typeof value === 'object' && value !== null && typeof (value as { then?: unknown }).then === 'function'
162
163/** `f` over a value that may be a promise: synchronous when the value is, and a failure becomes `fallback()`, never a throw. */
164function chain<T, U>(value: Maybe<T>, f: (v: T) => Maybe<U>, fallback: () => Maybe<U>): Maybe<U> {
165 try {
166 return isThenable(value) ? ((value as Promise<T>).then(f, fallback) as Maybe<U>) : f(value as T)
167 } catch {
168 return fallback()
169 }
170}
171
172export type ToastInput = {
173 level?: ToastLevel
174 text: string
175 timeoutMs?: number
176 /** Passes the `important` filter whatever its level (still muted by `off` and by a per-source mute). */
177 always?: boolean
178 /** Drawn only when the person is away: held back (and recorded) when `away()` answers false; drawn as usual when the host cannot say. */
179 awayOnly?: boolean
180}
181
182export type ToasterDeps = {
183 source: string
184 /** Milliseconds; may be a promise (`$.clock.now()`). The whole toast is synchronous when this and `prefs` are. */
185 now: () => Maybe<number>
186 /** Draws the line (`$.ui.toast`). May throw: that is `refused`. */
187 show: (line: string, options: { timeoutMs?: number }) => void
188 prefs?: () => Maybe<ToastPrefs>
189 /** Gets every digest, drawn or not. Fire and forget. */
190 persist?: (digest: Digest) => unknown
191 away?: () => boolean | undefined
192 /** Schedules the release of held errors; without it they are released by the next toast or `release()`. */
193 after?: (ms: number, fn: () => void) => unknown
194 dedupeMs?: number
195 rateMax?: number
196 windowMs?: number
197}
198
199export type Toaster = {
200 /** Decides, draws or holds, records. Resolves to what became of it; never rejects. */
201 toast: (input: ToastInput) => Maybe<Why | 'empty'>
202 /** Draws held errors once the window has room: `✗ <first> … and N more`. */
203 release: () => Maybe<void>
204}
205
206export function createToaster(deps: ToasterDeps): Toaster {
207 const dedupeMs = deps.dedupeMs ?? DEDUPE_MS
208 const rateMax = deps.rateMax ?? RATE_MAX
209 const windowMs = deps.windowMs ?? RATE_WINDOW_MS
210 const seen = new Map<string, number>()
211 let shownAt: number[] = []
212 let held: { text: string; n: number; timeoutMs?: number } | null = null
213 let isScheduled = false
214
215 const record = (d: Digest): void => {
216 try {
217 const r = deps.persist?.(d)
218
219 if (isThenable(r)) r.catch(() => undefined)
220 } catch {
221 /* a failed record never changes what was shown */
222 }
223 }
224
225 const draw = (now: number, line: string, timeoutMs: number | undefined): boolean => {
226 try {
227 deps.show(line, timeoutMs === undefined ? {} : { timeoutMs })
228 shownAt.push(now)
229
230 return true
231 } catch {
232 return false
233 }
234 }
235
236 const room = (now: number): boolean => {
237 shownAt = shownAt.filter(at => now - at < windowMs)
238
239 return shownAt.length < rateMax
240 }
241
242 const freeHeld = (now: number): void => {
243 if (held === null || !room(now)) return
244
245 const { text, n, timeoutMs } = held
246
247 held = null
248 draw(now, lineOf('error', n > 1 ? `${text} … and ${n - 1} more` : text), timeoutMs)
249 }
250
251 const arm = (now: number): void => {
252 if (isScheduled || deps.after === undefined || held === null) return
253
254 isScheduled = true
255
256 const wait = Math.max(50, windowMs - (now - (shownAt[0] ?? now)) + 50)
257
258 const lost = (): void => {
259 isScheduled = false
260 }
261
262 try {
263 const timer = deps.after(wait, () => {
264 isScheduled = false
265 void releaseNow()
266 })
267
268 if (isThenable(timer)) timer.catch(lost)
269 } catch {
270 lost()
271 }
272 }
273
274 /** A held error is said only while the setting still lets this source draw: switching toasts off, or muting the source, drops what waits (it is already recorded). */
275 const releaseAt = (now: number, prefs: ToastPrefs): void => {
276 if (prefs.mode === 'off' || prefs.muted.includes(deps.source)) held = null
277 freeHeld(now)
278 arm(now)
279 }
280
281 /** The work of a call at a moment: the setting is read, then `then` runs with it; a setting that cannot be read is the default. */
282 const withPrefs = <T>(then: (prefs: ToastPrefs) => T): Maybe<T> => {
283 let asked: Maybe<ToastPrefs>
284
285 try {
286 asked = deps.prefs?.() ?? DEFAULT_PREFS
287 } catch {
288 return then(DEFAULT_PREFS)
289 }
290
291 return chain(asked, then, () => then(DEFAULT_PREFS))
292 }
293
294 /** A clock that fails (a refused `clock.now`) is the wall clock: the toast is still decided. */
295 const atNow = <T>(then: (now: number) => Maybe<T>, fallback: () => T): Maybe<T> => {
296 let clock: Maybe<number>
297
298 try {
299 clock = deps.now()
300 } catch {
301 clock = Date.now()
302 }
303
304 return chain(clock, then, () => chain(Date.now(), then, fallback))
305 }
306
307 const releaseNow = (): Maybe<void> => atNow(now => withPrefs(prefs => releaseAt(now, prefs)), () => undefined)
308
309 const run = (now: number, prefsIn: ToastPrefs, input: ToastInput): Why | 'empty' => {
310 const level = LEVELS.find(candidate => candidate === input.level) ?? 'info'
311 const text = tidy(input.text, LINE_MAX - 2)
312
313 if (text === '') return 'empty'
314
315 const prefs = parsePrefs(JSON.stringify(prefsIn))
316 let why: Why
317
318 releaseAt(now, prefs)
319
320 if (prefs.mode === 'off') why = 'off'
321 else if (prefs.muted.includes(deps.source)) why = 'muted'
322 else if (prefs.mode === 'important' && (level === 'info' || level === 'ok') && input.always !== true) why = 'filtered'
323 else if (input.awayOnly === true && deps.away?.() === false) why = 'away'
324 else {
325 const last = seen.get(text)
326
327 if (last !== undefined && now - last < dedupeMs) why = 'deduped'
328 else if (room(now)) {
329 why = draw(now, lineOf(level, text), input.timeoutMs) ? 'shown' : 'refused'
330 if (why === 'shown') seen.set(text, now)
331 } else if (level === 'error') {
332 // An error is never dropped: it waits, and the next free slot says it once with a count.
333 held = held === null ? { text, n: 1, ...(input.timeoutMs !== undefined && { timeoutMs: input.timeoutMs }) } : { ...held, n: held.n + 1 }
334 seen.set(text, now)
335 why = 'coalesced'
336 arm(now)
337 } else why = 'rate-limited'
338 }
339
340 if (seen.size > 200) for (const [key, at] of seen) if (now - at >= dedupeMs) seen.delete(key)
341
342 record({ t: now, source: deps.source, level, text, shown: why === 'shown', why })
343
344 return why
345 }
346
347 return {
348 toast: input => atNow(now => withPrefs(prefs => run(now, prefs, input)), () => 'refused' as const),
349 release: releaseNow,
350 }
351}
352
353// ------------------------------------------------------------------------------------------------------------ the kit: the setting and the digests through files
354
355/** The three file calls a plugin has (`$.fs.read`, `$.fs.write`, `$.fs.exists`), relative to the project. Each may be refused. */
356export type ToastIo = { read: (path: string) => Promise<string>; write: (path: string, text: string) => Promise<void>; exists: (path: string) => Promise<boolean> }
357
358export type KitDeps = Omit<ToasterDeps, 'prefs' | 'persist'> & {
359 /** Without files the toaster draws by the defaults and records nothing on disk. */
360 io?: ToastIo
361 prefsTtlMs?: number
362}
363
364/**
365 * A toaster whose setting is the console's file (read at most every few seconds; none means all, nothing muted) and whose digests go
366 * to this source's own ring file under the console's folder, rewritten whole (a plugin has no append), one write in flight at a time.
367 * Nothing is written unless that folder exists: it is the console's, and its own `.gitignore` covers it. Without the console a plugin
368 * runs on the defaults and leaves no file.
369 */
370export function createToastKit(deps: KitDeps): Toaster {
371 const { io } = deps
372 const ttl = deps.prefsTtlMs ?? PREFS_TTL_MS
373 let prefs: ToastPrefs = DEFAULT_PREFS
374 let prefsAt = -Infinity
375 const ring: Digest[] = []
376 let isLoaded = false
377 let isDirty = false
378 let isWriting = false
379 let isConsole = false
380 let consoleAt = -Infinity
381
382 /**
383 * One write in flight at a time, the newest ring in it: digests that arrive meanwhile are written by the next pass, so a burst costs
384 * a few writes, not one each, and the last digest is always on disk. No timer: a toast is never lost to a debounce that did not fire.
385 */
386 const flush = async (): Promise<void> => {
387 if (io === undefined || isWriting) return
388
389 isWriting = true
390
391 try {
392 while (isDirty) {
393 isDirty = false
394
395 if (!isLoaded) {
396 isLoaded = true
397 const old = decodeRing(await io.read(`${TOAST_DIR}/${deps.source}.jsonl`).catch(() => ''))
398
399 ring.unshift(...old.filter(d => d.source === deps.source).slice(-RING_MAX))
400 if (ring.length > RING_MAX) ring.splice(0, ring.length - RING_MAX)
401 }
402
403 await io.write(`${TOAST_DIR}/${deps.source}.jsonl`, encodeRing(ring))
404 }
405 } catch {
406 /* the digest is a courtesy: a refused write is dropped */
407 } finally {
408 isWriting = false
409 }
410 }
411
412 /** Whether the console's folder exists, believed for a while: no console, no digest. */
413 const hasConsole = (now: number): Promise<boolean> => {
414 if (now - consoleAt < CONSOLE_TTL_MS) return Promise.resolve(isConsole)
415
416 consoleAt = now
417
418 return (io as ToastIo).exists(CONSOLE_DIR).then(
419 yes => (isConsole = yes),
420 () => (isConsole = false),
421 )
422 }
423
424 const readPrefs = (): Maybe<ToastPrefs> => {
425 if (io === undefined) return DEFAULT_PREFS
426
427 const at = (now: number): Maybe<ToastPrefs> => {
428 if (now - prefsAt < ttl) return prefs
429
430 prefsAt = now
431
432 return io.read(PREFS_FILE).then(
433 text => (prefs = parsePrefs(text)),
434 () => (prefs = DEFAULT_PREFS),
435 )
436 }
437
438 return chain(
439 deps.now(),
440 at,
441 () => at(Date.now()),
442 )
443 }
444
445 return createToaster({
446 ...deps,
447 prefs: readPrefs,
448 persist: d => {
449 if (io === undefined) return
450
451 // The ring is seeded from the file on the first write; digests made before it are kept in order.
452 return hasConsole(d.t).then(yes => {
453 if (!yes) return
454
455 pushDigest(ring, d)
456 isDirty = true
457 void flush()
458 })
459 },
460 })
461}
462hooks/views/model.ts 213 lines1import { membersOf, type Member, type MemberState } from '../model/members'
2import { topologyLines } from '../model/topology'
3import type { Decision, Proposal, TaskRecord } from '../reader/parse'
4import type { FileKey } from '../reader/snapshot'
5import type { ActionOutcome, PendingConfirm, State } from '../state'
6
7/** Below this many columns the pane draws its narrow form: a list, not a grid, and the board as counts. */
8export const NARROW = 44
9export const TILE = 18
10
11export type TaskRow = { id: string; status: string; type: string; description: string; owner: string; isSelected: boolean }
12
13/** Everything the pane draws, worked out from the state with no element in sight, so it is tested as data. */
14export type PaneModel = {
15 columns: number
16 rows: number
17 isNarrow: boolean
18 hasSwarm: boolean
19 title: string
20 members: Member[]
21 selected: Member | null
22 counts: { state: MemberState; count: number }[]
23 topology: string[]
24 board: { pending: number; claimed: number; done: number; failed: number; rows: TaskRow[]; selected: TaskRow | null; total: number }
25 proposals: Proposal[]
26 decisions: Decision[]
27 usage: string
28 route: string
29 /** ruOS desktops hosting swarm agents, one line each; empty when no agent runs remotely and no host file is there. */
30 remote: string[]
31 missing: string[]
32 confirm: PendingConfirm | null
33 outcome: ActionOutcome | null
34 detail: State['detail']
35 next: { text: string; why: string } | null
36 isActing: boolean
37}
38
39const MISSING_WORDS: Record<FileKey, string> = {
40 swarm: 'swarm store',
41 pointer: 'swarm pointer',
42 agents: 'agent store',
43 tasks: 'task store',
44 claims: 'claims',
45 hive: 'hive-mind',
46}
47
48const CLAIMED = new Set(['in_progress', 'running', 'assigned'])
49const DONE = new Set(['completed', 'done'])
50const FAILED = new Set(['failed', 'cancelled'])
51
52const ageOf = (ms: number): string => {
53 const s = Math.max(0, Math.round(ms / 1000))
54
55 return s < 60 ? `${s}s` : s < 3600 ? `${Math.floor(s / 60)}m` : `${Math.floor(s / 3600)}h`
56}
57
58/** The ruOS hosts section: each desktop with its state, heartbeat age and agents; "not on disk" when agents name a host and no snapshot exists. */
59export function remoteLines(snapshot: State['snapshot'], nowMs: number): string[] {
60 if (snapshot === null) {
61 return []
62 }
63
64 const hosts = snapshot.ruos.hosts
65 const remoteAgents = snapshot.agents.filter(agent => agent.remote !== undefined).length
66 const note = snapshot.ruos.note !== undefined ? [snapshot.ruos.note] : []
67
68 if (hosts === null) {
69 return remoteAgents > 0 ? [`ruOS hosts: not on disk (${remoteAgents} agent${remoteAgents === 1 ? '' : 's'} name one)`, ...note] : note
70 }
71
72 return [
73 ...hosts.slice(0, 20).map(host => {
74 const beat = host.heartbeatAt !== undefined ? ` · heartbeat ${ageOf(nowMs - Date.parse(host.heartbeatAt))} ago` : ' · no heartbeat'
75
76 return `${host.name} · ${host.state}${beat} · ${host.agents.length} agent${host.agents.length === 1 ? '' : 's'}`
77 }),
78 ...(hosts.length > 20 ? [`+${hosts.length - 20} more hosts`] : []),
79 ...(hosts.length === 0 ? ['no ruOS hosts'] : []),
80 ...note,
81 ]
82}
83
84const short = (id: string) => (id.length > 14 ? `…${id.slice(-8)}` : id)
85
86/** A cost or a token count as the engine gave it; a figure it left out is `n/a`, never zero. */
87export function usageLine(usage: State['usage']): string {
88 if (usage === null) {
89 return 'usage: not read yet'
90 }
91
92 const cost = usage.costUsd !== undefined ? `$${usage.costUsd.toFixed(usage.costUsd < 1 ? 3 : 2)}` : 'n/a'
93 const tokens = usage.contextTokens !== undefined ? `${Math.round(usage.contextTokens / 1000)}k` : 'n/a'
94 const window = usage.contextWindow !== undefined ? `/${Math.round(usage.contextWindow / 1000)}k` : ''
95 const percent = usage.contextPercent !== undefined ? ` (${Math.round(usage.contextPercent)}%)` : ''
96
97 return `cost ${cost} · context ${tokens}${window}${percent}`
98}
99
100/**
101 * The router's last pick the pane has seen (kept in `$.store`, so it may come from an earlier session: its age says),
102 * with its score; under the threshold it says so instead of naming a winner.
103 */
104export function routeLine(route: State['route'], threshold: number, nowMs = Date.now()): string {
105 if (route === null) {
106 return 'router: no pick seen yet'
107 }
108
109 const score = `${Math.round(route.confidence * 100)}%`
110 const alt = route.alternatives[0] !== undefined ? ` · alt ${route.alternatives[0].agent} ${Math.round(route.alternatives[0].confidence * 100)}%` : ''
111
112 const age = route.atMs > 0 ? ` · ${ageOf(nowMs - route.atMs)} ago` : ''
113
114 if (!route.matched) {
115 return `router: no match (best ${route.agent} ${score})${age}`
116 }
117
118 return route.confidence < threshold
119 ? `router: below threshold ${Math.round(threshold * 100)}% (best ${route.agent} ${score})${alt}${age}`
120 : `router: ${route.agent} ${score}${route.pattern !== undefined ? ` (${route.pattern})` : ''}${alt}${age}`
121}
122
123function nextOf(state: State, members: readonly Member[], board: PaneModel['board']): PaneModel['next'] {
124 const snapshot = state.snapshot
125
126 if (snapshot === null || !snapshot.hasSwarm) {
127 return { text: '/ruflo-swarm:swarm init', why: 'no swarm on disk here: start one' }
128 }
129
130 if (!members.some(member => member.source === 'ruflo')) {
131 return { text: 'npx @claude-flow/cli@latest agent spawn -t coder --name coder', why: 'the swarm has no ruflo agents yet' }
132 }
133
134 if (board.total === 0) {
135 return { text: 'npx @claude-flow/cli@latest task create --type implementation --description "…"', why: 'no tasks on the board' }
136 }
137
138 if (board.pending > 0 && board.claimed === 0) {
139 return { text: '/ruflo-swarm:swarm status', why: `${board.pending} task${board.pending === 1 ? '' : 's'} waiting and none claimed` }
140 }
141
142 return null
143}
144
145export function paneModelOf(state: State, columns: number, rows: number, nowMs: number): PaneModel {
146 const snapshot = state.snapshot
147 const members = membersOf(snapshot, state.activity, nowMs)
148 const selected = members.find(member => member.id === state.selected) ?? members[0] ?? null
149 const counts = (['working', 'blocked', 'idle', 'done', 'failed'] as const)
150 .map(memberState => ({ state: memberState, count: members.filter(member => member.state === memberState).length }))
151 .filter(entry => entry.count > 0)
152
153 const claimsByIssue = new Map((snapshot?.claims ?? []).map(claim => [claim.issueId, claim]))
154 const typeOf = new Map(members.map(member => [member.id, member.label]))
155
156 const taskRow = (task: TaskRecord): TaskRow => {
157 const claim = claimsByIssue.get(task.id)
158 const ownerId = claim?.claimant.id ?? task.assignedTo[0]
159
160 return {
161 id: task.id,
162 status: claim !== undefined && claim.status !== 'active' ? `${task.status}/${claim.status}` : task.status,
163 type: task.type,
164 description: task.description,
165 owner: ownerId === undefined ? 'unassigned' : (typeOf.get(ownerId) ?? short(ownerId)),
166 isSelected: task.id === state.selectedTask,
167 }
168 }
169
170 const tasks = snapshot?.tasks ?? []
171 // Open work first, newest store order kept within each group.
172 const rank = (task: TaskRecord) => (CLAIMED.has(task.status) ? 0 : task.status === 'pending' ? 1 : FAILED.has(task.status) ? 2 : 3)
173 const ordered = [...tasks].sort((a, b) => rank(a) - rank(b)).map(taskRow)
174 const selectedTask = ordered.find(row => row.id === state.selectedTask) ?? ordered[0] ?? null
175
176 const board = {
177 pending: tasks.filter(task => task.status === 'pending').length,
178 claimed: tasks.filter(task => CLAIMED.has(task.status) || claimsByIssue.get(task.id)?.status === 'active').length,
179 done: tasks.filter(task => DONE.has(task.status)).length,
180 failed: tasks.filter(task => FAILED.has(task.status)).length,
181 rows: ordered.map(row => ({ ...row, isSelected: row.id === selectedTask?.id })),
182 selected: selectedTask,
183 total: tasks.length,
184 }
185
186 const swarm = snapshot?.swarm
187 const topologyName = swarm?.topology ?? snapshot?.hive?.topology ?? 'unknown'
188
189 return {
190 columns,
191 rows,
192 isNarrow: columns < NARROW,
193 hasSwarm: snapshot?.hasSwarm === true,
194 title: swarm !== undefined && swarm !== null ? `swarm ${short(swarm.id)} · ${topologyName} · ${swarm.status}` : snapshot?.hive !== null && snapshot?.hive !== undefined ? `hive-mind · ${topologyName}` : 'ruflo swarm',
195 members,
196 selected,
197 counts,
198 topology: topologyLines(topologyName, members, columns),
199 board,
200 proposals: snapshot?.hive?.pending.filter(proposal => proposal.status === 'pending') ?? [],
201 decisions: (snapshot?.hive?.history ?? []).slice(-2).reverse(),
202 usage: usageLine(state.usage),
203 route: routeLine(state.route, state.options.routeThreshold, nowMs),
204 remote: remoteLines(snapshot, nowMs),
205 missing: (snapshot?.missing ?? []).map(key => MISSING_WORDS[key]),
206 confirm: state.confirm,
207 outcome: state.outcome,
208 detail: state.detail,
209 next: nextOf(state, members, board),
210 isActing: state.isActing,
211 }
212}
213