SLOPSHOPPER

ruflo-swarm

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

newpanespinnerguardcommandtoast
★ 74,184v0.3.5MITupdated 2026-10-07ruvnet/ruflo/plugins/ruflo-swarm
A shopper browsing a rack in a slop shop
Preview · a replayed session in a sandbox
claude · ~/work/app · ruflo-swarm
│ ┃ Swarm ✕ › fix the failing auth test and add an audit log call │ ┃ ruflo swarm [ hide ] │ ┃ ⏺ Read(src/auth.ts) │ ┃ No ruflo swarm on disk in this folder. ⎿ Read 6 lines │ ┃ ⏺ Update(src/auth.ts) │ ┃ Next /ruflo-swarm:swarm init ⎿ Added 2 lines, removed 1 line │ ┃ no swarm on disk here: start one [ put in pr ⏺ Bash(bun test) │ ┃ ⎿ 3 pass, 1 fail │ ┃ cost $0.420 · context 97k/200k (49%) │ ● Done. refresh now rejects expired claims and logs an audit event. │ │ ✻ Worked for 42s · done 4:20 PM │ │ › /ruflo-swarm-pane │ ⎿ ruflo-swarm: Swarm pane shown │ │ ────────────────────────────────────────────────────────────────────────────────────────────────────────────────────── › ? for shortcuts

Draws

Pane · Swarm
ruflo swarm [ hide ] No ruflo swarm on disk in this folder. Next /ruflo-swarm:swarm init no swarm on disk here: start one [ put in prompt ] cost $0.420 · context 97k/200k (49%)
README

ruflo-swarm

Agent teams, swarm coordination, Monitor streams, and worktree isolation.

Install

/plugin marketplace add ruvnet/ruflo
/plugin install ruflo-swarm@ruflo

What's Included

  • Agent Teams: TeamCreate, SendMessage, and Task tool integration for multi-agent coordination
  • Topologies: hierarchical, mesh, hierarchical-mesh, ring, star, adaptive
  • Monitor Streams: Real-time swarm status via Monitor("npx @claude-flow/cli@latest swarm watch --stream")
  • Worktree Isolation: Each agent works in its own git worktree to avoid conflicts
  • Hive-Mind Consensus: Byzantine, Raft, Gossip, CRDT, and Quorum strategies
  • Anti-Drift: hierarchical topology with specialized strategy for tight coordination

Requires

  • ruflo-core plugin (provides MCP server)

Compatibility

  • CLI: pinned to @claude-flow/cli v3.6 major+minor.
  • Verification: bash plugins/ruflo-swarm/scripts/smoke.sh is the contract.

MCP surface (12 tools)

FamilyCountTools
swarm_*4swarm_init, swarm_status, swarm_shutdown, swarm_health
agent_*8agent_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.

Built-in Claude Code coordination tools

This plugin pairs with Claude Code's native multi-agent tools (no MCP needed):

ToolPurpose
TaskSpawn a sub-agent (use name: for addressability + run_in_background: true for parallel execution)
SendMessageInter-agent comms (named agents only)
TaskCreate / TaskList / TaskGet / TaskUpdate / TaskOutput / TaskStopShared task tracker for swarm pipelines
MonitorLive-stream events from a long-running process (persistent: true) — primary wake signal for /loop
EnterWorktree / ExitWorktreeGit worktree isolation per agent

Anti-drift defaults (per CLAUDE.md)

For coding swarms, the canonical defaults that prevent agent drift:

SettingValueRationale
topologyhierarchicalCoordinator catches divergence
maxAgents6–8Smaller team = less drift
strategyspecializedClear roles, no overlap
consensusraftLeader maintains authoritative state
memoryhybridSQLite + AgentDB for both fast + durable

For 10+ agent teams, use hierarchical-mesh (queen + peer communication).

Namespace coordination

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

Live swarm pane (Claude Code mod, early access)

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.

  • Tiles: one per agent: ruflo's agents from .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.
  • Topology (hierarchical, mesh, ring, star) as ruflo wrote it, plus tasks with their claims, hive-mind proposals and their votes, cost and context from $.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.
  • Buttons act through the ruflo CLI with fixed argv and ids checked against ruflo's own format: logs or activity, claim, steal, hand off, pause or resume a claim, stop, offer for stealing, re-route, vote, and put the next command in the prompt (a person still presses Enter). Stop, steal and hand off ask for a second press. After each action the pane re-reads the disk and says whether the change shows there. A CLI that exits 0 while nothing changes on disk is reported as unverified.
  • Commands: /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.
  • Options (/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).
  • The mod does not register prompt.submit or tool.check hooks, and its tool.call hook only observes.
  • With the ruflo-mods trust gate: if you set 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.

Toasts (ADR-477)

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.

Verification

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

Architecture Decisions

Related Plugins

  • ruflo-agentdb — namespace convention owner
  • ruflo-autopilot — owns the 270s cache-aware /loop heartbeat for long-running swarms
  • ruflo-intelligence — hooks_route powers swarm agent recommendation per task
Source 17 files
hooks/register.ts 492 lines
1import 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}
492
hooks/actions/controller.ts 291 lines
1import 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}
291
hooks/audit.ts 60 lines
1import { 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}
60
hooks/adr-digest.ts 41 lines
1/**
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}
41
hooks/commands.ts 151 lines
1import 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}
151
hooks/host.ts 33 lines
1import 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}
33
hooks/model/members.ts 278 lines
1import { 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)
278
hooks/reader/parse.ts 365 lines
1/**
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}
365
hooks/reader/snapshot.ts 138 lines
1import {
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}
138
hooks/state.ts 135 lines
1import 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}
135
hooks/toast-policy.ts 462 lines
1/**
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}
462
hooks/views/model.ts 213 lines
1import { 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