Opt-in Claude Code mod for MCP Task Orchestrator: live work graph, gate band, and in-process orchestration hooks (early-access function-hooks API)

An opt-in Claude Code mod for MCP Task Orchestrator: a live work-graph pane, an in-flight band and status line, and in-process orchestration hooks. It ships as a separate plugin in the same marketplace as the task-orchestrator plugin and does not replace it.
A mod is a plugin of in-process TypeScript function hooks (hooks/hooks.json points at ./register.ts). It is additive: the server stays the gate authority, and the task-orchestrator plugin's command hooks remain the baseline for clients that do not load the mod.
The function-hooks API this mod uses is early access. A Claude Code update can break the mod or switch it off. It was built and tested against Claude Code 2.1.286, and the per-build types file shipped with your Claude Code is the authority on the API.
If the mod misbehaves, disable it. The task-orchestrator command hooks resume by themselves, because the dedupe flags (see Dedupe with the command hooks) are only set when the mod handled the event.
Verification status:
/to-graph pane and the band.mcp-task-orchestrator. The name is fixed in the mod and matches the registration shown in the root README..taskorchestrator/config.yaml with project.rootId, read relative to the session working directory.curl on PATH. See Live updates.Marketplace install, after the existing marketplace add from the root README:
/plugin install task-orchestrator-mod@task-orchestrator-marketplace
To disable it:
/plugin disable task-orchestrator-mod@task-orchestrator-marketplace
From a checkout on disk:
claude --plugin-dir <repo>/claude-plugins/task-orchestrator-mod
On the desktop app, where no flag can be passed, set CLAUDE_CODE_PLUGIN_DIRS in the process environment or in the env block of ~/.claude/settings.json. Do not put it in project settings. Both forms are watched and hot-reload on save.
Headless claude -p and claude plugin test need CLAUDE_CODE_ENABLE_FUNCTION_HOOKS=1.
Do not load the marketplace copy and a disk copy at the same time.
The marketplace install path has not yet been confirmed end to end. If your build needs an extra enable step, it is not described here.
The graph pane draws boxes and edges differently per surface. The terminal falls back to a raster; every other surface uses an SVG edge layer.
| Surface | Rendering | Status |
|---|---|---|
| Terminal | Boxes with a box-drawing edge raster | Kit-tested, not live-verified |
| Desktop Code tab | Boxes with an SVG edge layer; cell size comes from userConfig | Live-verified |
| VS Code and mobile | Take the SVG path | Unverified |
On all surfaces:
clip, pbcopy, wl-copy or xclip). On mobile that is the host machine's clipboard, which is unverified./to-graphroot or an item id./to-graph, or the band opening the pane, starts in auto: while the pane is open, a TO call or a spawned agent that names an item outside the graph moves the pane to that item's owning scope. An id argument, root, This feature, Whole project, a breadcrumb, Open graph or a click on a box pins the scope, and a pinned scope never moves; the Follow button (shown only while pinned) or a bare /to-graph returns to auto. Besides this session's own activity, an item transition made by another session in this project moves an auto pane after a 3 second pause, when the server sends the writer's actor and this session has been quiet for 30 seconds; without an actor only the echo window applies and nothing is followed.This feature and Whole project. The whole-project view is an overview of the root's children with role roll-ups, from a single call limited to 100 children; past that it shows "More children than shown."properties.planLabel (1 to 12 characters).Show done steps / Hide done steps toggle.Loading ... shows while the scope switches.stalled marks an item whose last move was 30 minutes ago or more. gate: marks this session's gate-blocked advance and is kept for 10 minutes.seat · model. A star marks a recent change.Reconnect shows only when live updates are degraded.item-id · done/total done · ready: ... · waiting: .... Otherwise it shows the in-flight item and its gate./to-band hides or shows it.TO, the item id prefix and a finished/remaining count.curl stream, or a 15 second poll) runs for the whole session.Opt-in through graphToasts. A toast appears when an item completes, becomes ready, is sent back from review to work, or one of this session's advances is gate-blocked. At most three toasts are shown per snapshot.
A change made by another session, seen over the live stream outside a short window after your own writes, is marked in the graph. Poll mode never marks changes.
The mod offers the session retrospective itself, by spawning the retrospective agent, when the session is not headless, a config is readable (at AGENT_CONFIG_DIR or the working directory) and the orchestration mode is not off. Test-kit verified, not live-verified.
An in-process version of the phase guard that blocks a sub-agent from stopping with its phase notes unfilled. Test-kit verified, not live-verified.
Test-kit verified, not live-verified. The mod repairs two common call shapes before the server sees them, and logs one transcript line starting call-shape: for each rewrite.
ancestorId=<project rootId> into list-mode query_items search, get_next_item, get_blocked_items and non-item get_context calls. It applies only when all of these hold: ancestorId is absent, there is no parentId, depth is not 0, there is no tags or type filter, and a root id resolves.actor into any manage_notes upsert element that lacks one, and keeps the top-level actor.To opt out of the ancestorId injection for a call, pass ancestorId: "". The server treats a blank value as absent.
The mod and the plugin's command hooks would otherwise act twice on the same event. The mod forwards a flag on the classic event, and the matching command hook exits early when it sees it.
| Area | Flag the mod sets | Command hooks that step aside |
|---|---|---|
| Retrospective | to_mod_retro: true, on classic.PostToolUse (advance_item, complete_tree) and classic.Stop | retro-trigger and retro-backstop |
| Phase guard | to_mod_active: true, after a clean run, on classic.PostToolUse and classic.SubagentStop | phase-guard-record and phase-guard |
If the phase guard hits an internal error it sends no flag, and the command hook acts instead.
Two command hooks are deliberately not ported: skill enforcement and actor-attribution enforcement. A flag on PreToolUse would land inside the tool input, so those command hooks stay authoritative.
Call-shape rewrites need no dedupe: the mod's tool-call hook runs before the classic PreToolUse, so the command hooks see the repaired input.
Version: the early exits exist only in a task-orchestrator plugin newer than 3.9.0. With 3.9.0 or older plus this mod, the retrospective is offered twice (a spawned agent plus the command-hook directive), and both phase guards block.
Live updates need the server's REST API (API_ENABLED=true) and curl on PATH.
The API URL is resolved in this order:
TASK_ORCHESTRATOR_API_URL environment variable..taskorchestrator/client.json in the working directory (a bare loopback origin only).client.json under TASK_ORCHESTRATOR_HOME, HOME or USERPROFILE.A bearer token from TASK_ORCHESTRATOR_API_TOKEN is passed to curl on stdin, never on the command line. The mod streams GET /api/v1/events for the project root. An event that carries the writer's actor is this session's own when that actor (or its parent) is one this session wrote with; only an event without an actor falls back to a 5 second echo window after this session's last write.
With no URL, or after three failed connects, a 15 second poll takes over. Reconnect in the pane restarts the stream.
Set these from the plugin's config menu rows, or in settings under pluginConfigs["task-orchestrator-mod"].options.
| Option | Default | Meaning |
|---|---|---|
cellWidthPx | 8.4 | Width of one character cell in the desktop pane; clamped to 4-32 |
cellHeightPx | 18 | Height of one character cell in the desktop pane; clamped to 8-64 |
graphToasts | false | Show graph toasts |
client.json are resolved from the session working directory only. There is no walk-up, no main-checkout lookup from a linked worktree, and no user-level config, unlike the task-orchestrator hooks. A session in a linked worktree falls back to polling unless the API URL is set through the environment or a user-level client.json.mcp-task-orchestrator./to-graph and /to-band are registered lazily, on the first band draw or snapshot./to-graph <prefix> of the root id walks the subtree instead of showing the overview; Copy UUID on Linux may report unavailable; a critical path can drop the SVG layer; and a teardown error can replace a tool result.ancestorId: "" to opt out):work-summary, Scoped Mode get_context(): a process-global scope such as Retrospective Trends, Observations or Proposals loses its active, blocked and stalled items.dependency-manager, Path B get_blocked_items(includeDetails=true): the broad view is narrowed to the project, so blocked process-global items vanish.ralph, troubleshooting query_items(operation="search", claimStatus="claimed"): claims outside the root are hidden.The skills are not changed by the mod; the skill-side fixes are tracked separately.
Authoring and testing rules for this plugin are in CLAUDE.md.
hooks/register.ts 21 lines1// Entry module for the task-orchestrator-mod plugin. Thin by design: each feature lives under
2// src/<feature>/ and owns its own hooks; this file only wires them. Add a feature by importing its
3// register function here — do not put hook logic in this file.
4import type { Register } from 'claude-code'
5
6import { registerBand } from '../src/band/index.ts'
7import { registerCallShape } from '../src/call-shape/index.ts'
8import { registerGraphData } from '../src/graph-data/index.ts'
9import { registerGraphPane } from '../src/graph-pane/index.ts'
10import { registerPhaseGuard } from '../src/phase-guard/index.ts'
11import { registerRetro } from '../src/retro/index.ts'
12
13export const register: Register = (on, options) => {
14 registerGraphData(on, options)
15 registerGraphPane(on, options)
16 registerBand(on)
17 registerRetro(on)
18 registerPhaseGuard(on)
19 registerCallShape(on)
20}
21src/band/index.ts 119 lines1// Status line + AbovePrompt band: in-flight item and gate progress (T4).
2// Owned by work item ec1e2f91-1d16-43dc-a9c9-f86aea28a050. Read-only: it draws from T2's snapshot
3// atom and makes no TO call and no write beyond its own state.
4//
5// The validator follows `$` only into functions of the same file, so the atoms below are this
6// file's own same-literal declarations of graph-data's keys (an atom is a typed reference to a
7// host-held value); only the pure steps and types come from graph-data.
8import { atom, read, update } from 'claude-code'
9import type { EngineInterface, On } from 'claude-code'
10
11import type { GraphSnapshot } from '../../types'
12import { addSubscriber, requestLiveSync } from '../graph-data/index.ts'
13import { bandLine, bandModel } from './model.ts'
14
15/** Below this many columns the band hides; the status line still carries the item. */
16export const MIN_BAND_COLUMNS = 40
17
18/** The /to-graph pane's id (graph-pane PANE_ID); the band's text opens it. */
19const PANE_ID = 'to-graph'
20
21const graphSnapshot = atom({ plugin: 'task-orchestrator-mod', key: 'graphSnapshot' } as const, null as GraphSnapshot | null)
22const graphSubscribers = atom({ plugin: 'task-orchestrator-mod', key: 'graphSubscribers' } as const, 0)
23const graphPaneOpen = atom({ plugin: 'task-orchestrator-mod', key: 'graphPaneOpen' } as const, false)
24const graphScopeMode = atom({ plugin: 'task-orchestrator-mod', key: 'graphScopeMode' } as const, 'auto' as 'auto' | 'pinned')
25const bandHidden = atom({ plugin: 'task-orchestrator-mod', key: 'bandHidden' } as const, false)
26const bandSubscribed = atom({ plugin: 'task-orchestrator-mod', key: 'bandSubscribed' } as const, false)
27
28const toggle = (hidden: boolean): boolean => !hidden
29const hide = (): boolean => true
30const yes = (): boolean => true
31const auto = (): 'auto' | 'pinned' => 'auto'
32
33/** Set once per module load, so a hot reload (which drops the registered command) registers it again. */
34let commandRegistered = false
35
36/**
37 * Brings the status line to the current snapshot, subscribes the band once per session (bandSubscribed
38 * lives in $.state, so a hot reload does not count twice) and registers /to-band, the way back after
39 * Hide. Idempotent; safe to call from any hook.
40 */
41async function sync($: EngineInterface): Promise<void> {
42 $.ui.status(bandModel(await read($, graphSnapshot), true).status)
43 if (!(await read($, bandSubscribed))) {
44 await update($, bandSubscribed, yes)
45 await update($, graphSubscribers, addSubscriber)
46 requestLiveSync(await read($, graphSubscribers))
47 }
48 if (!commandRegistered) {
49 commandRegistered = true
50 await $.command.register({ name: 'to-band', description: 'Show or hide the Task Orchestrator in-flight band' })
51 }
52}
53
54/**
55 * The band's text pressed: open the work graph pane on the current scope, and count it once as a
56 * subscriber, the same steps /to-graph runs (graph-pane index.ts). Run directly, not through a hook on
57 * our own state write (a plugin's own `$.state.set` does not reliably reach its own hooks).
58 */
59async function openPane($: EngineInterface): Promise<void> {
60 await $.ui.open({ id: PANE_ID, title: 'TO graph' })
61 if (!(await read($, graphPaneOpen))) {
62 // The band opens the pane on the current scope without choosing one: follow this session's activity.
63 await update($, graphScopeMode, auto)
64 await update($, graphPaneOpen, yes)
65 await update($, graphSubscribers, addSubscriber)
66 requestLiveSync(await read($, graphSubscribers))
67 }
68}
69
70export function registerBand(on: On): void {
71 // Graph-data writes the snapshot atom; this observes the write after it lands. There is no
72 // session.start hook here (graph-data owns the one unmatched session.start and the validator refuses
73 // a second). A write raised from a timer callback may not reach this hook in every engine, so the
74 // prompt and turn hooks below run the same sync from events that always dispatch the whole chain;
75 // the status line is then at worst one prompt or turn behind.
76 on('state.set', { plugin: 'task-orchestrator-mod', key: 'graphSnapshot' }, async ($, e, next) => {
77 const wrote = await next(e)
78 const result = wrote as { isSet?: boolean; value?: { isSet?: boolean } }
79 if (result.isSet === true || result.value?.isSet === true) await sync($)
80
81 return wrote
82 })
83
84 on('prompt.submit', async ($, e, next) => {
85 await sync($)
86
87 return next(e)
88 })
89
90 on('turn.complete', async ($, e, next) => {
91 await sync($)
92
93 return next(e)
94 })
95
96 on('command.run', { command: 'to-band' }, async $ => {
97 const hidden = await update($, bandHidden, toggle)
98
99 return { text: hidden ? 'In-flight band hidden. Run /to-band to show it again.' : 'In-flight band shown.' }
100 })
101
102 on('ui.render', { component: 'AbovePrompt' }, async ($, e, next) => {
103 if (e.props.hasSurvey || e.props.bodyColumns < MIN_BAND_COLUMNS) return next(e)
104 if (await read($, bandHidden)) return next(e)
105 // The smart line while the graph is scoped to an item, else today's in-flight line (model.ts).
106 const text = bandLine(await read($, graphSnapshot), e.props.bodyColumns)
107 if (text === undefined) return next(e)
108
109 const { Box, Button } = $.ui.resolve(e)
110
111 return h(
112 Box,
113 null,
114 h(Button, { key: 'band-open', label: text, plain: true, onPress: () => openPane($) }),
115 h(Button, { key: 'hide', label: 'Hide', onPress: () => update($, bandHidden, hide) }),
116 )
117 })
118}
119src/call-shape/index.ts 36 lines1// Call-shape repair rewrites for TO tool calls (T7).
2// Owned by work item cedcfc11-6205-432d-8968-371bdf1cb917.
3//
4// `tool.call` sits above `classic.PreToolUse`, so the TO command hooks (actor attribution, skill
5// enforcement) see the repaired input. Only the call's shape is repaired; server validation and
6// write semantics are untouched.
7import type { On } from 'claude-code'
8
9import { PLUGIN, toToolName } from '../shared/constants.ts'
10import { CONFIG_PATH, parseProjectRootId } from '../shared/config.ts'
11import { isOwnCall, rewriteCall } from './rewrite.ts'
12
13/** The TO tools whose input this feature can repair; other TO calls never reach the hook. */
14const REPAIRED_TOOLS = /^mcp__.*task-orchestrator.*__(?:query_items|get_next_item|get_blocked_items|get_context|manage_notes)$/
15
16export function registerCallShape(on: On): void {
17 on('tool.call', { tool: REPAIRED_TOOLS }, async ($, e, next) => {
18 const tool = toToolName(e.tool)
19 if (tool === null || isOwnCall(next.origin?.plugin, PLUGIN)) return next(e)
20
21 let rootId: string | null = null
22 try {
23 rootId = parseProjectRootId(await $.fs.read(CONFIG_PATH))
24 } catch {
25 rootId = null
26 }
27
28 const { input, changes } = rewriteCall(tool, e as unknown as Record<string, unknown>, rootId)
29 if (changes.length === 0) return next(e)
30
31 for (const change of changes) $.ui.log(`call-shape: ${change}`)
32
33 return next(input as unknown as typeof e)
34 })
35}
36src/graph-data/index.ts 253 lines1// Graph snapshot data layer + live invalidation (T2).
2// Owned by work item 32492bfa-93d5-4e72-93cc-e56ecf8b68c8. Read-only: the data layer renders
3// nothing and never calls a TO write tool. T3 (pane) and T4 (band) import from this file.
4//
5// This is the only file here that touches `$`: the plugin validator follows `$` only into functions
6// declared in the same file, so `graphIo` turns `$` into plain functions (GraphIo) and every other
7// module takes that. Consumers control the layer through the atoms (see the atoms below), not by calling
8// helpers with their own `$`.
9import { atom, read, update } from 'claude-code'
10import type { EngineInterface, On, PluginOptions } from 'claude-code'
11
12import { PLUGIN, TO_SERVER, toToolName } from '../shared/constants.ts'
13import type { GraphSnapshot, GraphStatus } from '../../types'
14import { touchedIds } from '../graph-pane/activity.ts'
15import { snapshotEvents, toastLines } from './events.ts'
16import type { GraphIo } from './io.ts'
17import { invalidateLabels } from './labels.ts'
18import { restartLive, syncLive } from './live.ts'
19import { refresh, refreshNow, sameSnapshot, sameStatus } from './refresh.ts'
20import { actorIdsOf, clearRemote, rememberLocalActors, resultItemIds, withLocalWrite } from './remote.ts'
21
22export type { GateInfo, GraphEdge, GraphNode, GraphSnapshot, GraphStatus } from '../../types'
23// ── Atoms: the control and data surface for T3 (pane) and T4 (band) ──
24// Consumers cannot call helpers that take `$` (the validator follows `$` only into functions of the
25// same file), so they use the library's `update($, atom, step)` / `read($, atom)` on these:
26// - show a scope: update($, graphScope, scopeTo(id)) // null = the project root
27// - start/stop live: update($, graphSubscribers, addSubscriber) // on pane/band open
28// update($, graphSubscribers, removeSubscriber) // on close
29// - ask for a refresh: update($, graphRefreshRequest, bump)
30// The write alone does not reach this plugin's own state.set hooks reliably, so each write is followed by
31// the matching `$`-free request function below (requestRefresh, requestLiveSync, requestReconnect).
32// The state.set hooks further down remain a backstop. Plugin and key are literals so
33// `claude plugin validate` can list them.
34
35/** The item id in view; null means the project root. Writing it re-snapshots at once. */
36export const graphScope = atom({ plugin: 'task-orchestrator-mod', key: 'graphScope' } as const, null as string | null)
37
38/** The latest snapshot, or null before the first one lands. Written by graph-data only. */
39export const graphSnapshot = atom({ plugin: 'task-orchestrator-mod', key: 'graphSnapshot' } as const, null as GraphSnapshot | null)
40
41/** Whether a refresh is running, the last error, and which live source keeps the snapshot current. */
42export const graphStatus = atom(
43 { plugin: 'task-orchestrator-mod', key: 'graphStatus' } as const,
44 { refreshing: false, liveSource: 'none' } as GraphStatus,
45)
46
47/** Consumers showing the graph. The SSE stream or the poll runs only while this is above 0. */
48export const graphSubscribers = atom({ plugin: 'task-orchestrator-mod', key: 'graphSubscribers' } as const, 0)
49
50/** Bump it to request a debounced re-snapshot. */
51export const graphRefreshRequest = atom({ plugin: 'task-orchestrator-mod', key: 'graphRefreshRequest' } as const, 0)
52
53/** Bump it to restart the live source (SSE/poll) and re-snapshot. */
54export const graphReconnectRequest = atom({ plugin: 'task-orchestrator-mod', key: 'graphReconnectRequest' } as const, 0)
55
56/** Items whose last change came from another session (SSE, not our echo): id -> when (ms). Written through graphIo.updateRemote only. */
57export const graphRemote = atom({ plugin: 'task-orchestrator-mod', key: 'graphRemote' } as const, {} as Record<string, number>)
58
59/** Steps for `update($, atom, step)`. */
60export const addSubscriber = (n: number): number => n + 1
61export const removeSubscriber = (n: number): number => Math.max(0, n - 1)
62export const bump = (n: number): number => n + 1
63export const scopeTo =
64 (id: string | null) =>
65 (): string | null =>
66 id
67
68export type { GraphIo } from './io.ts'
69export { snapshot } from './snapshot.ts'
70export { invalidateLabels } from './labels.ts'
71export { remoteAdvanceListener, setRemoteAdvanceListener } from './live.ts'
72
73/** TO tools whose calls can change the graph; a successful call re-snapshots. */
74export const WRITE_TOOLS: ReadonlySet<string> = new Set([
75 'advance_item',
76 'manage_notes',
77 'create_work_tree',
78 'manage_items',
79 'manage_dependencies',
80 'complete_tree',
81 'claim_item',
82])
83
84/** The engine's spelling of those tools (`mcp__<server>__<tool>`), for the hook matcher. */
85export const WRITE_TOOL_NAME = /^mcp__.*task-orchestrator.*__(advance_item|manage_notes|create_work_tree|manage_items|manage_dependencies|complete_tree|claim_item)$/
86
87/**
88 * Whether a finished tool call should refresh the graph: a TO write tool, not raised by this plugin's
89 * own `$` call (its reads re-enter `tool.call`), and neither denied nor errored.
90 */
91export function shouldRefresh(tool: string, originPlugin: string, result: { deny?: unknown; isError?: unknown }): boolean {
92 const name = toToolName(tool)
93
94 return name !== null && WRITE_TOOLS.has(name) && originPlugin !== PLUGIN && result.deny === undefined && result.isError !== true
95}
96
97/** The parsed JSON text of a TO tool result; throws on an MCP error result or bad JSON. */
98export function parseToolResult(tool: string, result: { content: readonly unknown[]; isError?: boolean }): unknown {
99 let text = ''
100 for (const block of result.content) {
101 const b = block as { type?: string; text?: string }
102 if (b.type === 'text' && typeof b.text === 'string') {
103 text = b.text
104 break
105 }
106 }
107 if (result.isError) throw new Error(`${tool}: ${text || 'error result'}`)
108
109 return JSON.parse(text)
110}
111
112/**
113 * The IO bound when the session started. A plugin's own `$.state.set` does not reliably reach its own
114 * `state.set` hooks (verified live 2026-10-05: scope changes never refreshed), so consumers in this
115 * plugin call these `$`-free functions right after writing a control atom.
116 */
117let sessionIo: GraphIo | null = null
118
119/** The `graphToasts` option (default off): set by registerGraphData, read when a snapshot lands. */
120let toastsOn = false
121
122/** Re-snapshot for the current scope now (a scope change) or after the debounce window. */
123export function requestRefresh(now = true): void {
124 if (sessionIo !== null) void (now ? refreshNow(sessionIo) : refresh(sessionIo))
125}
126
127/** Restart the live source and re-snapshot (the Reconnect button). */
128export function requestReconnect(subscribers: number): void {
129 if (sessionIo === null) return
130 invalidateLabels()
131 restartLive(sessionIo, subscribers)
132 void refreshNow(sessionIo)
133}
134
135/** Start or stop the live source for a new subscriber count (pane or band opened or closed). */
136export function requestLiveSync(subscribers: number): void {
137 if (sessionIo !== null) syncLive(sessionIo, subscribers)
138}
139
140function graphIo($: EngineInterface): GraphIo {
141 return {
142 callTool: async (tool, args) => parseToolResult(tool, await $.mcp.call(TO_SERVER, tool, args)),
143 readFile: path => $.fs.read(path),
144 now: () => $.clock.now(),
145 after: (ms, fn) => $.clock.after(ms, fn),
146 every: (ms, fn) => $.clock.every(ms, fn),
147 sleep: (ms, signal) => $.clock.sleep(ms, { signal }),
148 spawn: request => $.process.spawn(request),
149 readScope: () => read($, graphScope),
150 // Skip writes that change nothing: every write redraws the pane and band (and flashes the Svg frame).
151 // A real change may raise toasts (graphToasts on): only against the previous snapshot of the same
152 // scope, so the first snapshot after a load or a scope switch is quiet (events.ts).
153 setSnapshot: async value => {
154 const prev = await read($, graphSnapshot)
155 if (sameSnapshot(prev, value)) return
156 await update($, graphSnapshot, () => value)
157 if (toastsOn) for (const line of toastLines(snapshotEvents(prev, value))) $.ui.toast(line)
158 },
159 updateRemote: async step => {
160 const current = await read($, graphRemote)
161 if (step(current) === current) return
162 await update($, graphRemote, step)
163 },
164 updateStatus: async change => {
165 // Pre-check only: the write applies `change` to the value current at write time, so a
166 // concurrent setLiveSource is never overwritten by a stale precomputed status.
167 const current = await read($, graphStatus)
168 if (sameStatus(current, change(current))) return
169 await update($, graphStatus, change)
170 },
171 envApiUrl: () => $.env.get('TASK_ORCHESTRATOR_API_URL'),
172 envApiToken: () => $.env.get('TASK_ORCHESTRATOR_API_TOKEN'),
173 homeDir: async () => (await $.env.get('TASK_ORCHESTRATOR_HOME')) || (await $.env.get('HOME')) || (await $.env.get('USERPROFILE')),
174 }
175}
176
177/** A successful local write clears the remote marks of every item it named or cascaded to. Advisory. */
178async function clearLocalMarks($: EngineInterface, e: unknown, text: unknown): Promise<void> {
179 const ids = [...touchedIds(e).map(t => t.id), ...resultItemIds(text)]
180 if (ids.length === 0) return
181 try {
182 await graphIo($).updateRemote?.(map => clearRemote(map, ids))
183 } catch {
184 // marks are advisory
185 }
186}
187
188export function registerGraphData(on: On, options: PluginOptions = {}): void {
189 toastsOn = options.graphToasts === true
190 // Matched (the validator refuses a second unmatched `tool.call`): same names as WRITE_TOOLS.
191 // Matcher spelled as a literal: the engine resolves matchers from source, and an imported or exported constant stays unresolved (never matches).
192 on('tool.call', { tool: /^mcp__.*task-orchestrator.*__(advance_item|manage_notes|create_work_tree|manage_items|manage_dependencies|complete_tree|claim_item)$/ }, async ($, e, next) => {
193 if (next.origin.plugin === PLUGIN) return next(e)
194 // This session's actor ids, recorded before the call: an echo can arrive before it returns (remote.ts).
195 rememberLocalActors(actorIdsOf(e))
196 // The local-write window: SSE events for this write (its echo) are not marked as remote (remote.ts).
197 const ran = await withLocalWrite(() => next(e), () => $.clock.now())
198 // Debounced and never awaited: the model's call is not slowed by the snapshot.
199 if (shouldRefresh(e.tool, next.origin.plugin, ran)) {
200 if (toToolName(e.tool) === 'manage_items') invalidateLabels()
201 void refresh(graphIo($))
202 void clearLocalMarks($, e, (ran as { text?: unknown }).text)
203 }
204
205 return ran
206 })
207
208 on('session.start', async ($, e, next) => {
209 const started = await next(e)
210 sessionIo = graphIo($)
211 // A hot reload resets module state (the live source) but $.state survives, so resume live from the
212 // surviving subscriber count. A fresh session has 0, which makes this a no-op.
213 syncLive(sessionIo, await read($, graphSubscribers))
214 void refresh(sessionIo)
215
216 return started
217 })
218
219 // Consumer controls: they write these atoms, and the layer reacts once the write has landed.
220 // state.set matchers use the literal plugin name for the same reason (an imported PLUGIN never matched).
221 on('state.set', { plugin: 'task-orchestrator-mod', key: 'graphScope' }, async ($, e, next) => {
222 const wrote = await next(e)
223 if (wrote.isSet) void refreshNow(graphIo($))
224
225 return wrote
226 })
227
228 on('state.set', { plugin: 'task-orchestrator-mod', key: 'graphRefreshRequest' }, async ($, e, next) => {
229 const wrote = await next(e)
230 if (wrote.isSet) void refresh(graphIo($))
231
232 return wrote
233 })
234
235 on('state.set', { plugin: 'task-orchestrator-mod', key: 'graphReconnectRequest' }, async ($, e, next) => {
236 const wrote = await next(e)
237 if (wrote.isSet) {
238 invalidateLabels()
239 restartLive(graphIo($), await read($, graphSubscribers))
240 void refreshNow(graphIo($))
241 }
242
243 return wrote
244 })
245
246 on('state.set', { plugin: 'task-orchestrator-mod', key: 'graphSubscribers' }, async ($, e, next) => {
247 const wrote = await next(e)
248 if (wrote.isSet) syncLive(graphIo($), typeof e.value === 'number' ? e.value : 0)
249
250 return wrote
251 })
252}
253src/graph-pane/index.ts 723 lines1// /to-graph pane: a top-down dependency graph of keyed, state-filled boxes (T3 a47dcd6f, rewritten).
2// Read-only: the pane makes only query_items, get_context and query_notes calls and no TO write tool, ever.
3//
4// This is the only file here that touches `$` (the plugin validator follows `$` only into functions
5// of the same file). Control is atom-driven, so the atoms below are this file's own same-literal
6// declarations of graph-data's keys; only the pure steps come from the other modules.
7import { atom, read, update } from 'claude-code'
8import type { EngineInterface, On, PluginOptions } from 'claude-code'
9
10import type { GraphActivity, GraphAgentInfo, GraphDetail, GraphSnapshot, GraphStatus } from '../../types'
11import { CONFIG_PATH, parseProjectRootId } from '../shared/config.ts'
12import { PLUGIN, TO_SERVER } from '../shared/constants.ts'
13import { parseToResult } from '../shared/to-client.ts'
14import { addSubscriber, bump, removeSubscriber, requestLiveSync, requestReconnect, requestRefresh, scopeTo, setRemoteAdvanceListener, shouldRefresh } from '../graph-data/index.ts'
15import { GATE_BLOCK_MS, WORKING_TTL_MS, RECENT_MS, gateBlocks, pruneActivity, recordActivity, recordGateBlocks, rememberAgent, resolveId, seatOf, touchedIds } from './activity.ts'
16import type { GateBlockRow } from './activity.ts'
17import { edgePlan, extrasChars, marksFit } from './budget.ts'
18import { cellSize } from './cell.ts'
19import { layoutTD } from './layout.ts'
20import { cut, titleLines } from './wrap.ts'
21import type { TopDown } from './layout.ts'
22import { cardsOf, criticalPath, stepsOf } from './model.ts'
23import type { Card, CriticalPath, Model } from './model.ts'
24import { recordFocusOrder, redirectFocus } from './focus.ts'
25import { gateToastLines } from '../graph-data/events.ts'
26import { activeScope, detailLines, firstSpawnUuid, isDegraded, isSwitching, loadingTitle, owningScope, parseScopeArg, scopeHeader, summaryLine } from './pane-model.ts'
27import type { ToolCall } from './pane-model.ts'
28import { CRIT_COLOR, raster, runs } from './raster.ts'
29import { routes } from './route.ts'
30import { KIND, READY, id8, legendItems } from './shared.ts'
31import type { GraphView } from './shared.ts'
32import { edgeSvg, px } from './svg.ts'
33
34export const PANE_ID = 'to-graph'
35
36/** Every TO tool this pane calls; none is a write tool. */
37export const READ_TOOLS: readonly string[] = ['query_items', 'get_context', 'query_notes']
38
39const graphSnapshot = atom({ plugin: 'task-orchestrator-mod', key: 'graphSnapshot' } as const, null as GraphSnapshot | null)
40const graphStatus = atom({ plugin: 'task-orchestrator-mod', key: 'graphStatus' } as const, { refreshing: false, liveSource: 'none' } as GraphStatus)
41const graphScope = atom({ plugin: 'task-orchestrator-mod', key: 'graphScope' } as const, null as string | null)
42const graphScopeMode = atom({ plugin: 'task-orchestrator-mod', key: 'graphScopeMode' } as const, 'auto' as 'auto' | 'pinned')
43const graphSubscribers = atom({ plugin: 'task-orchestrator-mod', key: 'graphSubscribers' } as const, 0)
44const graphReconnectRequest = atom({ plugin: 'task-orchestrator-mod', key: 'graphReconnectRequest' } as const, 0)
45const graphPaneOpen = atom({ plugin: 'task-orchestrator-mod', key: 'graphPaneOpen' } as const, false)
46const graphDetail = atom({ plugin: 'task-orchestrator-mod', key: 'graphDetail' } as const, null as GraphDetail | null)
47const graphActivity = atom({ plugin: 'task-orchestrator-mod', key: 'graphActivity' } as const, { working: {}, changed: {} } as GraphActivity)
48const graphAgents = atom({ plugin: 'task-orchestrator-mod', key: 'graphAgents' } as const, {} as Record<string, GraphAgentInfo>)
49const graphShowDone = atom({ plugin: 'task-orchestrator-mod', key: 'graphShowDone' } as const, false)
50const graphRemote = atom({ plugin: 'task-orchestrator-mod', key: 'graphRemote' } as const, {} as Record<string, number>)
51
52/** The engine's spelling of every TO tool (`mcp__<server>__<tool>`); reads count as "working on it". */
53const TO_TOOL = /^mcp__.*task-orchestrator.*__[a-z_]+$/
54
55const yes = (): boolean => true
56const no = (): boolean => false
57const pinned = (): 'auto' | 'pinned' => 'pinned'
58const auto = (): 'auto' | 'pinned' => 'auto'
59
60/** Items already resolved to an owning scope (null: none). Capped; a repeat touch of an id reads nothing. */
61const SCOPE_MEMO_MAX = 200
62const scopeMemo = new Map<string, string | null>()
63
64/** Set once per module load, so a hot reload (which drops the registered command) registers it again. */
65let commandRegistered = false
66
67/**
68 * Copies `text` with the host's clipboard tool (`clip` on Windows, `pbcopy` on macOS, `wl-copy` or
69 * `xclip` on Linux). Used when `$.ui.copy` cannot reach the surface, as on the desktop app today.
70 */
71async function hostCopy($: EngineInterface, text: string): Promise<boolean> {
72 const windows = (await $.env.get('OS')) === 'Windows_NT'
73 const tools: string[][] = windows ? [['clip']] : [['pbcopy'], ['wl-copy'], ['xclip', '-selection', 'clipboard']]
74 for (const argv of tools) {
75 try {
76 const ran = await $.process.run(argv, { stdin: text, timeoutMs: 3000 })
77 if (ran.exitCode === 0) return true
78 } catch {
79 // that tool is not installed here; try the next
80 }
81 }
82
83 return false
84}
85
86/**
87 * The scope `/to-graph` opens with no argument: the owning scope of the latest transition under the
88 * project root (pane-model `activeScope`), else the most recently modified feature-implementation in
89 * work, else null. A failed or empty scope read falls back to the legacy query; no root id makes no call.
90 */
91async function activeFeature($: EngineInterface): Promise<string | null> {
92 try {
93 const rootId = parseProjectRootId(await $.fs.read(CONFIG_PATH))
94 if (rootId === null) return null
95 const call = async (tool: string, args: Record<string, unknown>): Promise<unknown> => parseToResult<unknown>(tool, await $.mcp.call(TO_SERVER, tool, args))
96 try {
97 const scope = await activeScope(call, rootId, await $.clock.now())
98 if (scope !== null) return scope
99 } catch {
100 // fall through to the legacy rule
101 }
102 const found = (await call('query_items', {
103 operation: 'search',
104 ancestorId: rootId,
105 type: 'feature-implementation',
106 role: 'work',
107 sortBy: 'modifiedAt',
108 sortOrder: 'desc',
109 limit: 1,
110 })) as { items?: { id?: unknown }[] }
111 const id = found.items?.[0]?.id
112
113 return typeof id === 'string' ? id : null
114 } catch {
115 return null
116 }
117}
118
119const message = (err: unknown): string => (err instanceof Error ? err.message : String(err))
120
121/** Reads the three read-only calls for one node into the detail atom. Each failure is tolerated. In-pane, read-only. */
122async function openDetail($: EngineInterface, itemId: string): Promise<void> {
123 // Clicking a node is a choice of focus: stop following this session's activity.
124 await update($, graphScopeMode, pinned)
125 const call = async (tool: string, args: Record<string, unknown>): Promise<unknown> => parseToResult<unknown>(tool, await $.mcp.call(TO_SERVER, tool, args))
126 const [ctx, item, notes] = await Promise.allSettled([
127 call('get_context', { itemId }),
128 call('query_items', { operation: 'get', itemId, includeTimestamps: true }),
129 call('query_notes', { operation: 'list', itemId, includeBody: false }),
130 ])
131 const now = await $.clock.now()
132 const lines =
133 ctx.status === 'rejected'
134 ? [`Could not read ${id8(itemId)}: ${message(ctx.reason)}`]
135 : detailLines({
136 itemId,
137 ctx: ctx.value,
138 now,
139 ...(item.status === 'fulfilled' ? { item: item.value } : { itemFailed: true }),
140 ...(notes.status === 'fulfilled' ? { notes: notes.value } : { notesFailed: true }),
141 })
142 await update($, graphDetail, () => ({ itemId, lines }))
143}
144
145/**
146 * Auto mode only, pane open: re-scopes the pane to the owning scope of the first candidate id the
147 * snapshot cannot place (one attempt per call). Never throws; the caller's result is untouched.
148 */
149async function followActivity($: EngineInterface, pick: (rootId: string) => string[]): Promise<void> {
150 try {
151 if ((await read($, graphScopeMode)) !== 'auto' || !(await read($, graphPaneOpen))) return
152 const snap = await read($, graphSnapshot)
153 if (snap === null) return
154 const rootId = parseProjectRootId(await $.fs.read(CONFIG_PATH))
155 if (rootId === null) return
156 const known = snap.nodes.map(n => n.id)
157 const target = pick(rootId).find(c => resolveId(c, known) === null)
158 if (target === undefined) return
159 const key = `${rootId}:${target}`
160 let scope: string | null
161 if (scopeMemo.has(key)) scope = scopeMemo.get(key) ?? null
162 else {
163 const call: ToolCall = async (tool, args) => parseToResult<unknown>(tool, await $.mcp.call(TO_SERVER, tool, args))
164 try {
165 scope = await owningScope(call, target, rootId)
166 } catch {
167 scope = null
168 }
169 if (scopeMemo.size >= SCOPE_MEMO_MAX) scopeMemo.delete(scopeMemo.keys().next().value as string)
170 scopeMemo.set(key, scope)
171 }
172 if (scope === null || scope === (await read($, graphScope))) return
173 await update($, graphDetail, () => null)
174 await update($, graphScope, scopeTo(scope))
175 requestRefresh()
176 } catch {
177 // following is best effort: the model's call is never affected
178 }
179}
180
181/** Seat and expiry bookkeeping of one TO tool call: who is working on which snapshot item, what changed. */
182async function noteActivity(
183 $: EngineInterface,
184 e: { tool: string; agentId?: string },
185 changed: boolean,
186): Promise<void> {
187 const snap = await read($, graphSnapshot)
188 if (snap === null) return
189 const known = snap.nodes.map(n => n.id)
190 const agents = await read($, graphAgents)
191 const hits: { id: string; seat: string; model?: string }[] = []
192 for (const t of touchedIds(e)) {
193 const id = resolveId(t.id, known)
194 if (id === null) continue
195 const model = e.agentId !== undefined ? agents[e.agentId]?.model : undefined
196 hits.push({ id, seat: seatOf(e.agentId, agents, t.actor), ...(model !== undefined ? { model } : {}) })
197 }
198 if (hits.length === 0) return
199 const now = await $.clock.now()
200 const fold = (cur: GraphActivity): GraphActivity => hits.reduce((acc, h) => recordActivity(acc, { ...h, agentId: e.agentId ?? 'main', changed }, now), cur)
201 const before = await read($, graphActivity)
202 if (fold(before) === before) return
203 await update($, graphActivity, fold)
204 // Entries expire even when nothing else redraws the pane (and a hot reload drops these timers: every
205 // record call prunes first, so stale entries clear on the next write too).
206 const prune = async (): Promise<void> => {
207 const t = await $.clock.now()
208 const cur = await read($, graphActivity)
209 if (pruneActivity(cur, t) !== cur) await update($, graphActivity, c => pruneActivity(c, t))
210 }
211 if (changed) $.clock.after(RECENT_MS + 500, () => void prune())
212 $.clock.after(WORKING_TTL_MS + 500, () => void prune())
213}
214
215/** Marks this session's gate-blocked advance_item rows on their snapshot items; they expire after GATE_BLOCK_MS. */
216async function noteGateBlocks($: EngineInterface, rows: readonly GateBlockRow[]): Promise<void> {
217 const snap = await read($, graphSnapshot)
218 if (snap === null) return
219 const known = snap.nodes.map(n => n.id)
220 const hits: { id: string; missing: string[]; target?: string }[] = []
221 for (const r of rows) {
222 const id = resolveId(r.itemId, known)
223 if (id !== null) hits.push({ id, missing: r.missing, ...(r.targetRole !== undefined ? { target: r.targetRole } : {}) })
224 }
225 if (hits.length === 0) return
226 const now = await $.clock.now()
227 const before = await read($, graphActivity)
228 if (recordGateBlocks(before, hits, now) === before) return
229 await update($, graphActivity, cur => recordGateBlocks(cur, hits, now))
230 $.clock.after(GATE_BLOCK_MS + 500, async () => {
231 const t = await $.clock.now()
232 const cur = await read($, graphActivity)
233 if (pruneActivity(cur, t) !== cur) await update($, graphActivity, c => pruneActivity(c, t))
234 })
235}
236
237/** With the `graphToasts` option on: one toast per gate-blocked row of this session (capped; events.ts). */
238async function toastGateBlocks($: EngineInterface, rows: readonly GateBlockRow[]): Promise<void> {
239 for (const line of gateToastLines(rows, await read($, graphSnapshot))) $.ui.toast(line)
240}
241
242/** What the tool.call hook needs of the outside world: record one call's activity and its gate blocks. */
243export interface ActivityIo {
244 note: (e: { tool: string; agentId?: string }, changed: boolean) => Promise<void>
245 gate: (rows: readonly GateBlockRow[]) => Promise<void>
246}
247
248/** The toast goes first: a block on an item outside the snapshot is still worth one, and marking it may throw. */
249const activityIo = ($: EngineInterface, toasts: boolean): ActivityIo => ({
250 note: async (e, changed) => {
251 await noteActivity($, e, changed)
252 await followActivity($, () => touchedIds(e).map(t => t.id))
253 },
254 gate: async rows => {
255 if (toasts) await toastGateBlocks($, rows)
256 await noteGateBlocks($, rows)
257 },
258})
259
260/**
261 * The tool.call hook body: skips the pane's own calls, always returns the call's result untouched.
262 * The gate blocks are recorded after the call's activity, so the call's own write does not clear them.
263 */
264export async function trackToolCall<E extends { tool: string; agentId?: string }, R extends { deny?: unknown; isError?: unknown; text?: unknown }>(
265 io: ActivityIo,
266 e: E,
267 next: ((e: E) => Promise<R>) & { origin: { plugin: string } },
268): Promise<R> {
269 if (next.origin.plugin === PLUGIN) return next(e)
270 const ran = await next(e)
271 try {
272 await io.note(e, shouldRefresh(e.tool, next.origin.plugin, ran))
273 const blocks = ran.deny === undefined ? gateBlocks(e.tool, ran.text) : []
274 if (blocks.length > 0) await io.gate(blocks)
275 } catch {
276 // bookkeeping only: the model's call is never affected
277 }
278
279 return ran
280}
281
282/** Registers /to-graph once per module load. Idempotent. */
283async function ensureCommand($: EngineInterface): Promise<void> {
284 if (commandRegistered) return
285 commandRegistered = true
286 try {
287 await $.command.register({ name: 'to-graph', description: 'Open the Task Orchestrator work graph (active feature, an item id, or root)' })
288 } catch {
289 commandRegistered = false
290 }
291}
292
293type Ui = ReturnType<EngineInterface['ui']['resolve']>
294
295interface BoxSpec {
296 key: string
297 id: string
298 rect: { left: number; top: number; width: number; height: number }
299 fill: string
300 line1: string
301 line3: string
302 /** The state line is a warning: drawn undimmed. */
303 warn: boolean
304 recent: boolean
305 /** Last changed by another session: `⇄ ` after the `✱ ` on line 1. */
306 remote: boolean
307}
308
309/**
310 * Below this many boxes every line of a box is clickable; from it on only the state line. 90 whole-click
311 * cards (COST.wholeCardChars) leave room for step chips in the tree budget; canvasOf also falls back to
312 * one Button per box whenever whole-click boxes alone would not fit.
313 */
314export const WHOLE_CLICK_MAX = 90
315
316/** One state-filled box: a pure function of its spec (an id-derived key, nothing view-wide). */
317function boxOf(ui: Ui, s: BoxSpec, open: (id: string) => void, wholeClick = true): unknown {
318 const { Box, Text, Button } = ui
319 const n = s.rect.width - 2
320 // A recently changed item is marked in its title (a Button label takes no colour); a change made by
321 // another session adds `⇄ ` after it.
322 const [first, second] = titleLines(`${s.recent ? '✱ ' : ''}${s.remote ? '⇄ ' : ''}${s.line1}`, n)
323 const press = () => open(s.id)
324 const dim = s.warn ? {} : { dimColor: true }
325
326 const frame = { key: s.key, position: 'absolute', top: s.rect.top, left: s.rect.left, width: s.rect.width, height: s.rect.height, backgroundColor: s.fill, flexDirection: 'column', paddingX: 1 }
327 if (!wholeClick) {
328 return h(
329 Box,
330 frame,
331 h(Text, { color: '#ffffff', bold: true, wrap: 'truncate-end' }, first),
332 h(Text, { color: '#ffffff', wrap: 'truncate-end' }, second),
333 h(Button, { key: `open:${s.id}`, label: cut(s.line3, n) || ' ', plain: true, ...dim, onPress: press }),
334 )
335 }
336
337 // Box takes no onPress and a Button label is a single string, so every line of the box is a plain
338 // Button opening the same detail: a click anywhere on the box opens it.
339 return h(
340 Box,
341 frame,
342 h(Button, { key: `open:${s.id}:1`, label: first || ' ', plain: true, onPress: press }),
343 h(Button, { key: `open:${s.id}:2`, label: second || ' ', plain: true, onPress: press }),
344 h(Button, { key: `open:${s.id}`, label: cut(s.line3, n) || ' ', plain: true, ...dim, onPress: press }),
345 )
346}
347
348interface CanvasInput {
349 model: Model
350 lay: TopDown
351 desktop: boolean
352 cell: { w: number; h: number }
353 steps: number
354 open: (id: string) => void
355 count: number
356 /** The project-root overview: the root plus one row of children, no steps. */
357 overview: boolean
358 /** What the detail panel and breadcrumb add (charged to the edge drawing only). */
359 extraChars: number
360 /** The critical path (empty on the overview). */
361 crit: CriticalPath
362}
363
364/** The graph canvas (or a one-line notice when the tree would not fit). */
365function canvasOf(ui: Ui, i: CanvasInput): unknown[] {
366 const { Box, Text, Svg } = ui
367 const { model, lay } = i
368 const rs = routes(lay, model.cards, i.desktop ? 'px' : 'cell', i.crit.edges)
369 const critIds = new Set(i.crit.ids.filter(id => lay.cards.has(id)))
370 const marks = critIds.size
371 let edges: unknown = null
372 let plan: ReturnType<typeof edgePlan>
373 const cardCount = lay.cards.size + (model.root !== undefined ? 1 : 0)
374 // Whole-click costs more per card: fall back to one Button per card rather than refuse the boxes.
375 const wholeClick = model.cards.length < WHOLE_CLICK_MAX && edgePlan({ desktop: i.desktop, wholeClick: true, cards: cardCount, chips: lay.chips.length, cells: 0, runs: 0, svgChars: 0, svgWidth: 0, svgHeight: 0 }) !== 'too-large'
376 if (i.desktop) {
377 const svg = edgeSvg(rs, i.cell, lay.width, lay.height)
378 plan = edgePlan({ desktop: true, wholeClick, extraChars: i.extraChars, marks, cards: lay.cards.size + (model.root !== undefined ? 1 : 0), chips: lay.chips.length, cells: 0, runs: 0, svgChars: svg.length, svgWidth: px(lay.width, i.cell.w), svgHeight: px(lay.height, i.cell.h) })
379 if (plan === 'full' && Svg !== undefined) {
380 edges = h(Box, { key: 'edges', position: 'absolute', top: 0, left: 0 }, h(Svg, { source: svg, alt: 'dependency edges', width: px(lay.width, i.cell.w), height: px(lay.height, i.cell.h) }))
381 }
382 } else {
383 const cells = raster(rs)
384 const merged = runs(cells)
385 plan = edgePlan({ desktop: false, wholeClick, extraChars: i.extraChars, marks, cards: lay.cards.size + (model.root !== undefined ? 1 : 0), chips: lay.chips.length, cells: cells.size, runs: merged.length, svgChars: 0, svgWidth: 0, svgHeight: 0 })
386 if (plan === 'full') {
387 edges = h(
388 Box,
389 { key: 'edges', position: 'absolute', top: 0, left: 0, width: lay.width, height: lay.height },
390 ...[...cells].map(([k, c]) => {
391 const [x, y] = k.split('_').map(Number) as [number, number]
392
393 return h(Box, { key: `e:${k}`, position: 'absolute', top: y, left: x, width: 1, height: 1 }, h(Text, { color: c.color, dimColor: c.dim }, c.glyph))
394 }),
395 )
396 } else if (plan === 'runs') {
397 edges = h(
398 Box,
399 { key: 'edges', position: 'absolute', top: 0, left: 0, width: lay.width, height: lay.height },
400 ...merged.map(r => h(Box, { key: `e:${r.x}_${r.y}`, position: 'absolute', top: r.y, left: r.x, width: r.n, height: 1 }, h(Text, { color: r.color, dimColor: r.dim }, r.glyph.repeat(r.n)))),
401 )
402 }
403 }
404 if (plan === 'too-large') {
405 recordFocusOrder(null)
406
407 return [h(Text, { key: 'too-large' }, `Too large to draw (${i.count} items). Open a smaller scope.`)]
408 }
409 // Stripes go with the edges; with boxes only, when they still fit (they never refuse the scope).
410 const stripes = marks > 0 && marksFit({ desktop: i.desktop, wholeClick, extraChars: i.extraChars, marks, cards: cardCount, chips: lay.chips.length, cells: 0, runs: 0, svgChars: 0, svgWidth: 0, svgHeight: 0 }, plan)
411 /** Box ids in drawn order, for the one-stop-per-box focus. */
412 const order: string[] = []
413
414 const byId = new Map(model.cards.map(c => [c.id, c]))
415 const boxes: unknown[] = []
416 const root = model.root
417 if (root !== undefined) {
418 const words = root.kind === 'terminal' ? 'done' : root.kind
419 boxes.push(
420 boxOf(ui, { key: `card:${root.id}`, id: root.id, rect: lay.root, fill: KIND[root.kind], line1: `${root.glyph} [${root.label}] ${root.title}`, line3: i.overview ? `${words} · ${model.cards.length} children` : `${words} · ${model.cards.length} items · ${i.steps} steps`, warn: false, recent: false, remote: false }, i.open, wholeClick),
421 )
422 order.push(root.id)
423 }
424 for (const row of lay.rows) {
425 if (row.chip) {
426 const chip = lay.chips.find(c => c.step === row.step)
427 if (chip === undefined) continue
428 boxes.push(
429 h(
430 Box,
431 { key: `step:${chip.step}`, position: 'absolute', top: chip.rect.top, left: chip.rect.left, width: chip.rect.width, height: 1, backgroundColor: KIND.terminal },
432 h(Text, { color: '#ffffff', wrap: 'truncate-end' }, `✓ step ${chip.step} · ${chip.done} done${chip.cancelled > 0 ? `, ${chip.cancelled} cancelled` : ''}`),
433 ),
434 )
435 continue
436 }
437 for (const id of row.ids) {
438 const card = byId.get(id) as Card
439 const rect = lay.cards.get(id)
440 if (rect === undefined) continue
441 boxes.push(boxOf(ui, { key: `card:${id}`, id, rect, fill: card.ready ? READY : KIND[card.kind], line1: `${card.glyph} [${card.label}] ${card.title}`, line3: card.stateText, warn: card.warn, recent: card.recent, remote: card.remote }, i.open, wholeClick))
442 order.push(id)
443 // A sibling stripe on the card's left column: the card's own subtree is untouched (G5 keys hold).
444 if (stripes && critIds.has(id)) {
445 boxes.push(h(Box, { key: `crit:${id}`, position: 'absolute', top: rect.top, left: rect.left, width: 1, height: rect.height, backgroundColor: CRIT_COLOR }))
446 }
447 }
448 }
449 recordFocusOrder(order)
450
451 const omitted = plan === 'omit' ? [h(Text, { key: 'edges-omitted', dimColor: true }, 'Too many edges to draw here; boxes only.')] : []
452
453 return [...omitted, h(Box, { key: 'canvas', position: 'relative', width: lay.width, height: lay.height }, ...(edges !== null ? [edges] : []), ...boxes)]
454}
455
456export function registerGraphPane(on: On, options: PluginOptions = {}): void {
457 const cell = cellSize(options as { cellWidthPx?: unknown; cellHeightPx?: unknown })
458 const toasts = options.graphToasts === true
459
460 // Every unmatched hook of the session events is already taken (graph-data owns session.start, band
461 // owns prompt.submit/turn.complete, and the validator refuses a second one), so the command is
462 // registered from two matched hooks: the first snapshot landing (graph-data refreshes at session
463 // start) and the first AbovePrompt draw, whichever comes first. Idempotent.
464 on('state.set', { plugin: 'task-orchestrator-mod', key: 'graphSnapshot' }, async ($, e, next) => {
465 const wrote = await next(e)
466 await ensureCommand($)
467
468 return wrote
469 })
470
471 // The terminal redraws the prompt area at once and often, so this is the reliable early trigger.
472 on('ui.render', { component: 'AbovePrompt' }, async ($, e, next) => {
473 await ensureCommand($)
474
475 return next(e)
476 })
477
478 on('command.run', { command: 'to-graph' }, async ($, e) => {
479 const arg = parseScopeArg(e.args)
480 if (arg.kind === 'invalid') return { text: `Unknown argument "${arg.text}". Use /to-graph, /to-graph root, or /to-graph <item id>.` }
481 const scope = arg.kind === 'id' ? arg.id : arg.kind === 'root' ? null : await activeFeature($)
482
483 await update($, graphDetail, () => null)
484 await update($, graphScopeMode, arg.kind === 'active' ? auto : pinned)
485 await update($, graphScope, scopeTo(scope))
486 requestRefresh()
487 await $.ui.open({ id: PANE_ID, title: 'TO graph' })
488 // Counted once per open pane: a repeat /to-graph neither double-counts nor leaks a subscriber.
489 // Another session's transition in this project (graph-data/live.ts, debounced and quiet-gated): a $-free listener over this pane's own follow.
490 setRemoteAdvanceListener(itemId => void followActivity($, () => [itemId]))
491 if (!(await read($, graphPaneOpen))) {
492 await update($, graphPaneOpen, yes)
493 await update($, graphSubscribers, addSubscriber)
494 requestLiveSync(await read($, graphSubscribers))
495 }
496
497 return { text: scope === null ? 'Opened the work graph for the project root.' : `Opened the work graph for ${scope.slice(0, 8)}.` }
498 })
499
500 on('ui.close', { id: PANE_ID }, async ($, e, next) => {
501 const closed = await next(e)
502 setRemoteAdvanceListener(null)
503 if (await read($, graphPaneOpen)) {
504 await update($, graphPaneOpen, no)
505 await update($, graphSubscribers, removeSubscriber)
506 requestLiveSync(await read($, graphSubscribers))
507 }
508
509 return closed
510 })
511
512 // Remembers each spawned subagent's seat and model, so its TO calls can be attributed. Unmatched, and
513 // the only agent.spawn hook of the mod; the result is returned untouched.
514 on('agent.spawn', async ($, e, next) => {
515 const spawned = await next(e)
516 if (spawned.agentId !== undefined) {
517 const agentId = spawned.agentId
518 const info: GraphAgentInfo = { seat: e.subagentType.replace(/^[^:]+:/, ''), ...(spawned.model !== undefined ? { model: spawned.model } : {}) }
519 try {
520 await update($, graphAgents, cur => rememberAgent(cur, agentId, info))
521 } catch {
522 // bookkeeping only: the spawn is never affected
523 }
524 }
525 await followActivity($, rootId => {
526 const id = firstSpawnUuid(e.description, e.prompt, rootId)
527
528 return id === null ? [] : [id]
529 })
530
531 return spawned
532 })
533
534 // Who is working on which item, and what changed a moment ago. Matched (the mod's other tool.call hooks
535 // are matched too); the pane's own reads are skipped; the result is always returned as it came.
536 on('tool.call', { tool: TO_TOOL }, ($, e, next) => trackToolCall(activityIo($, toasts), e, next))
537
538 // One focus stop per box: a whole-click box's title lines hand the ring to its state line (focus.ts).
539 // Matcher spelled as a literal (an imported constant stays unresolved and never matches).
540 on('ui.focus', { requestId: 'to-graph' }, async ($, e, next) => redirectFocus(e, next))
541
542 on('ui.render', { component: 'Pane', requestId: PANE_ID }, async ($, e) => {
543 // Cross-session follow: (re)set the one-slot listener on every draw so a pane opened from the band or redrawn after a hot reload still follows (a draw after close is harmless: followActivity checks graphPaneOpen).
544 setRemoteAdvanceListener(itemId => void followActivity($, () => [itemId]))
545 const ui = $.ui.resolve(e)
546 const { Box, Text, Button, Svg } = ui
547 const snap = await read($, graphSnapshot)
548 const status = await read($, graphStatus)
549 const scope = await read($, graphScope)
550 const mode = await read($, graphScopeMode)
551 const detail = await read($, graphDetail)
552 const activity = await read($, graphActivity)
553 const showDone = await read($, graphShowDone)
554 const remote = (await read($, graphRemote)) ?? {}
555 const open = (id: string): Promise<void> => openDetail($, id)
556
557 const view: GraphView | null = snap
558 const switching = isSwitching(scope, snap)
559 const model = view === null ? null : cardsOf(view, activity.working, activity.changed, activity.blocked ?? {}, remote)
560 const steps = model === null ? null : stepsOf(model.cards)
561 const crit: CriticalPath = model === null || steps === null || view?.overview === true ? { ids: [], edges: new Set() } : criticalPath(model.cards, steps)
562 const lay = model === null || steps === null ? null : layoutTD(model.cards, steps, { bodyColumns: (e.props as { bodyColumns?: number }).bodyColumns, showDone: showDone || view?.overview === true, hasRoot: model.root !== undefined, wrap: view?.overview === true })
563
564 // Refresh is automatic (SSE or poll); Reconnect restarts the live source and shows only while degraded.
565 const controls = h(
566 Box,
567 { key: 'controls', flexDirection: 'row', columnGap: 1, marginBottom: 1 },
568 // Two-state scope toggle; the active scope draws as the primary button.
569 h(Button, {
570 key: 'scope-feature',
571 label: 'This feature',
572 ...(scope !== null ? { variant: 'primary' } : {}),
573 onPress: async () => {
574 const feature = await activeFeature($)
575 if (feature === null) $.ui.toast('No active feature in this project.')
576 else {
577 await update($, graphScopeMode, pinned)
578 await update($, graphScope, scopeTo(feature))
579 requestRefresh()
580 }
581 },
582 }),
583 h(Button, { key: 'scope-project', label: 'Whole project', ...(scope === null ? { variant: 'primary' } : {}), onPress: async () => {
584 await update($, graphScopeMode, pinned)
585 await update($, graphScope, scopeTo(null))
586 requestRefresh()
587 } }),
588 // Shown only while pinned: back to following this session's own TO activity (the bare /to-graph result).
589 ...(mode === 'pinned'
590 ? [h(Button, { key: 'scope-follow', label: 'Follow', onPress: async () => {
591 await update($, graphScopeMode, auto)
592 await update($, graphScope, scopeTo(await activeFeature($)))
593 requestRefresh()
594 } })]
595 : []),
596 ...(lay !== null && view?.overview !== true && lay.doneSteps.length > 0
597 ? [h(Button, { key: 'done-steps', label: showDone ? 'Hide done steps' : 'Show done steps', onPress: () => update($, graphShowDone, v => !v) })]
598 : []),
599 ...(isDegraded(status) ? [h(Button, { key: 'reconnect', label: 'Reconnect', onPress: async () => {
600 await update($, graphReconnectRequest, bump)
601 requestReconnect(await read($, graphSubscribers))
602 } })] : []),
603 )
604
605 if (view === null || model === null || steps === null || lay === null) {
606 return h(Box, { flexDirection: 'column' }, controls, h(Text, { key: 'loading' }, status.lastError ?? 'Loading the work graph…'))
607 }
608
609 // A scope switch is still loading: say so, and draw nothing of the old scope (the header would claim it).
610 if (switching) return h(Box, { flexDirection: 'column' }, controls, h(Text, { key: 'switching' }, `Loading ${loadingTitle(scope, snap)}…`))
611
612 // With a breadcrumb the scope is named there; the header keeps the counts and the live source.
613 const trail = view.trail ?? []
614 const header = [...(trail.length > 0 ? [] : [scopeHeader(view)]), summaryLine(view), `live: ${status.liveSource}${status.refreshing ? ' · refreshing' : ''}`].join(' · ')
615 const crumbs = trail.flatMap((c, i) => {
616 const last = i === trail.length - 1
617 const sep = i > 0 ? [h(Text, { key: `crumb-sep:${c.id}`, dimColor: true }, ' › ')] : []
618 if (last) return [...sep, h(Text, { key: `crumb:${c.id}`, bold: true }, cut(c.title, 40))]
619 const target = c.id === view.rootId ? null : c.id
620
621 return [
622 ...sep,
623 h(Button, {
624 key: `crumb:${c.id}`,
625 label: cut(c.title, 28),
626 plain: true,
627 onPress: async () => {
628 await update($, graphDetail, () => null)
629 await update($, graphScopeMode, pinned)
630 await update($, graphScope, scopeTo(target))
631 requestRefresh()
632 },
633 }),
634 ]
635 })
636 const body: unknown[] = []
637 if (view.truncated) body.push(h(Text, { key: 'truncated', dimColor: true }, view.overview === true ? 'More children than shown.' : 'Large subtree: showing only the shallowest 150 items.'))
638 if (view.error !== undefined) body.push(h(Text, { key: 'error', color: 'red' }, view.error))
639
640 if (view.nodes.length === 0) {
641 body.push(h(Text, { key: 'empty' }, 'No items in this scope.'))
642 } else {
643 body.push(
644 h(
645 Box,
646 { key: 'legend', flexDirection: 'row', flexWrap: 'wrap', rowGap: 1, marginBottom: 1 },
647 ...legendItems().map(item => h(Box, { key: `legend-${item.key}`, backgroundColor: item.color, paddingX: 1, marginRight: 1 }, h(Text, { color: '#ffffff' }, item.text))),
648 h(Box, { key: 'legend-ready', backgroundColor: READY, paddingX: 1, marginRight: 1 }, h(Text, { color: '#ffffff' }, '○ ready')),
649 h(Text, { key: 'legend-edges', dimColor: true }, '┄ contains ╌ open blocker ─ satisfied'),
650 ...(crit.ids.length > 0 ? [h(Text, { key: 'legend-crit', color: CRIT_COLOR }, ' ━ critical path')] : []),
651 ...(view.nodes.some(n => n.planLabel !== undefined) ? [h(Text, { key: 'legend-label', dimColor: true }, ' Tn = plan label')] : []),
652 ),
653 )
654 if (lay.tooWide !== null) {
655 body.push(h(Text, { key: 'too-wide', dimColor: true }, `Graph is ${lay.tooWide} columns wide; the pane shows ${lay.cols}. Widen the pane or open a smaller scope.`))
656 }
657 body.push(...canvasOf(ui, { model, lay, desktop: e.surface !== 'terminal' && Svg !== undefined, cell, steps: steps.max, open, count: view.nodes.length, overview: view.overview === true, extraChars: extrasChars(detail?.lines ?? null, trail.map(c => c.title)), crit }))
658 }
659
660 if (detail !== null) {
661 // The detail is its own bordered panel: title line bold, then facts, then a row of actions.
662 const facts = [
663 ...detail.lines.map((line, i) => h(Text, { key: `detail-${i}`, ...(i === 0 ? { bold: true } : {}) }, line)),
664 h(Text, { key: 'detail-uuid', dimColor: true }, `uuid: ${detail.itemId}`),
665 ]
666 const working = (activity.working[detail.itemId] ?? []).map(w =>
667 h(Text, { key: `working-${w.agentId}` }, `working: ${w.seat}${w.model !== undefined ? ` · ${w.model}` : ''} · ${w.agentId === 'main' ? 'main loop' : id8(w.agentId)}`),
668 )
669 const actions = [
670 h(Button, {
671 key: 'detail-copy',
672 label: 'Copy UUID',
673 onPress: async press => {
674 let copied = false
675 try {
676 const res = (await $.ui.copy({ text: detail.itemId, surface: press.surface })) as { isCopied?: boolean; value?: { isCopied?: boolean } }
677 copied = res.isCopied === true || res.value?.isCopied === true
678 } catch {
679 copied = false
680 }
681 // $.ui.copy has no path on a remote surface (the desktop app) yet: fall back to the OS clipboard tool.
682 if (!copied) copied = await hostCopy($, detail.itemId)
683 $.ui.toast(copied ? `Copied ${detail.itemId}` : `Copy unavailable here. UUID: ${detail.itemId}`)
684 },
685 }),
686 ...(detail.itemId !== view.scopeId
687 ? [
688 h(Button, {
689 key: 'detail-open-graph',
690 label: 'Open graph',
691 onPress: async () => {
692 await update($, graphDetail, () => null)
693 await update($, graphScopeMode, pinned)
694 await update($, graphScope, scopeTo(detail.itemId))
695 requestRefresh()
696 },
697 }),
698 ]
699 : []),
700 h(Button, { key: 'detail-close', label: 'Close detail', onPress: () => update($, graphDetail, () => null) }),
701 ]
702 body.push(
703 h(
704 Box,
705 { key: 'detail', flexDirection: 'column', borderStyle: 'round', borderColor: '#4b5563', paddingX: 1, marginTop: 1 },
706 ...facts,
707 ...working,
708 h(Box, { key: 'detail-actions', flexDirection: 'row', columnGap: 1, marginTop: 1 }, ...actions),
709 ),
710 )
711 }
712
713 return h(
714 Box,
715 { flexDirection: 'column', paddingX: 1 },
716 ...(crumbs.length > 0 ? [h(Box, { key: 'trail', flexDirection: 'row', flexWrap: 'wrap' }, ...crumbs)] : []),
717 h(Box, { key: 'header', marginBottom: 1 }, h(Text, { bold: crumbs.length === 0, dimColor: crumbs.length > 0 }, header)),
718 controls,
719 ...body,
720 )
721 })
722}
723src/phase-guard/index.ts 103 lines1// Phase guard + advisory hooks in-process (T6).
2// Owned by work item 26246238-681f-497d-ae4c-6dbafd2c1e12.
3//
4// In-process ports of four TO command hooks, semantics unchanged except that reads go through the
5// engine's own MCP connection (D2) and state lives in session-scoped `$.state` (D7), so no tmpdir
6// marker file is written while the mod runs:
7// - phase-guard-record.mjs -> classic.PostToolUse (advance_item), registered
8// - phase-guard.mjs -> classic.SubagentStop, returns { block }, registered
9// - skill-enforcement.mjs -> classic.PreToolUse (manage_notes), NOT registered, see below
10// - enforce-actor-attribution.mjs -> classic.PreToolUse (advance_item, manage_notes), NOT registered
11//
12// Dedupe (D3): after a clean run the hook forwards `next({ ...e, [MOD_ACTIVE_FLAG]: true })` so the
13// command hook, which checks `to_mod_active` on stdin, exits early. On an internal error it forwards
14// without the flag and the command hook acts as the fallback. T1 P3 proved the flag reaches stdin on
15// classic.PostToolUse; SubagentStop is unproven (if it does not reach stdin, the command phase-guard
16// finds no marker - recording was suppressed - and no-ops, so nothing is lost or doubled).
17//
18// The two PreToolUse ports (actorAttribution, skillEnforcement in core.ts) are kept and unit-tested but
19// have NO registration here, deliberately: verified live, on classic.PreToolUse the flag a mod adds lands
20// INSIDE tool_input, so it would be forwarded to the tool and the TO server as an argument. A PreToolUse
21// wrapper therefore must not inject the flag; the command hooks stay authoritative for those two events.
22//
23// This is the only file here that touches `$` (the validator follows `$` only into functions declared
24// in the same file); `guardIo` turns `$` into the plain functions the rest takes.
25import type { EngineInterface, On } from 'claude-code'
26
27import type { PhaseGuardEntry } from '../../types'
28import { MOD_ACTIVE_FLAG, PLUGIN, TO_SERVER, toToolName } from '../shared/constants.ts'
29import { CONFIG_PATH } from '../shared/config.ts'
30import { actorAttribution, checkStop, normalizeEntry, recordAdvance, skillEnforcement } from './core.ts'
31import type { GuardIo } from './core.ts'
32
33export { actorAttribution, checkStop, recordAdvance, skillEnforcement } from './core.ts'
34
35function parseToolResult(tool: string, result: { content: readonly unknown[]; isError?: boolean }): unknown {
36 let text = ''
37 for (const block of result.content) {
38 const b = block as { type?: string; text?: string }
39 if (b.type === 'text' && typeof b.text === 'string') {
40 text = b.text
41 break
42 }
43 }
44 if (result.isError) throw new Error(`${tool}: ${text || 'error result'}`)
45
46 return JSON.parse(text)
47}
48
49function guardIo($: EngineInterface): GuardIo {
50 return {
51 getEntry: async agentId => normalizeEntry((await $.state.get({ plugin: 'task-orchestrator-mod', key: 'phaseGuardAgents', id: agentId } as const)).value),
52 setEntry: async (agentId, entry: PhaseGuardEntry) => {
53 await $.state.set({ plugin: 'task-orchestrator-mod', key: 'phaseGuardAgents', id: agentId } as const, entry)
54 },
55 getWarned: async () => {
56 const { value } = await $.state.get({ plugin: 'task-orchestrator-mod', key: 'skillWarned' } as const)
57
58 return Array.isArray(value) ? value : []
59 },
60 setWarned: async pairs => {
61 await $.state.set({ plugin: 'task-orchestrator-mod', key: 'skillWarned' } as const, pairs)
62 },
63 callTo: async (tool, args) => parseToolResult(tool, await $.mcp.call(TO_SERVER, tool, args)),
64 readConfig: async () => {
65 try {
66 return await $.fs.read(CONFIG_PATH)
67 } catch {
68 return null
69 }
70 },
71 isHeadless: async () => (await $.env.get('TASK_ORCHESTRATOR_MODE')) === 'headless-iteration',
72 }
73}
74
75export function registerPhaseGuard(on: On): void {
76 // phase-guard-record.mjs: remember which items a subagent entered a phase for.
77 on('classic.PostToolUse', async ($, e, next) => {
78 // Not an advance_item call, or one this plugin's own `$` call raised: pass through untouched (D4).
79 if (toToolName(e.tool_name) !== 'advance_item' || next.origin?.plugin === PLUGIN) return next(e)
80 try {
81 await recordAdvance(guardIo($), e as never)
82 } catch {
83 return next(e)
84 }
85
86 return next({ ...e, [MOD_ACTIVE_FLAG]: true } as never)
87 })
88
89 // phase-guard.mjs: block a subagent that stops with its phase's required notes missing.
90 on('classic.SubagentStop', async ($, e, next) => {
91 if (next.origin?.plugin === PLUGIN) return next(e)
92 let reason: string | null
93 try {
94 reason = await checkStop(guardIo($), e as never)
95 } catch {
96 return next(e)
97 }
98 if (reason !== null) return { block: reason }
99
100 return next({ ...e, [MOD_ACTIVE_FLAG]: true } as never)
101 })
102}
103src/retro/index.ts 272 lines1// Retrospective dispatch in-process (T5).
2// Owned by work item 0f4fe024-9c3f-454c-83f7-ec13de173355.
3//
4// In-process port of retro-trigger.mjs (PostToolUse) and retro-backstop.mjs (Stop). State lives in
5// session-scoped `$.state` (it survives /compact and hot reload and is never keyed by rootId), and
6// the retrospective is spawned with `$.agent.spawn` at the run boundary instead of being handed to
7// the model as a directive.
8//
9// tool.call (advance_item / complete_tree) classify, update state, add a nudge or the one-line
10// "queued" notice to the result's context
11// tool.call (Bash running retro-ack.mjs) ack the mod's state
12// classic.PostToolUse forward `to_mod_retro` so retro-trigger.mjs steps aside
13// classic.Stop spawn the queued retrospective once nothing runs in the
14// background; else the nudge-only backstop
15//
16// Ownership: the mod owns retro handling only when the session is not a headless iteration AND the
17// cwd config is readable (the mod cannot see AGENT_CONFIG_DIR / main-checkout / user configs, and
18// must not override a config it cannot see). Otherwise every hook passes through untouched and the
19// command hooks act as before. The flag is feature-scoped (`to_mod_retro`), so it cannot silence
20// another feature's command hook.
21//
22// This is the only file here that touches `$` (the validator follows `$` only into functions
23// declared in the same file); logic.ts is pure.
24import type { EngineInterface, On } from 'claude-code'
25
26import type { RetroState } from '../../types'
27import { CONFIG_PATH } from '../shared/config.ts'
28import { PLUGIN, toToolName } from '../shared/constants.ts'
29import { readSection, scalar } from '../lib/yaml-lite.mjs'
30import {
31 ackState,
32 afterBackstop,
33 applyCall,
34 backstopRoots,
35 buildDispatch,
36 buildNudge,
37 buildSpawnPrompt,
38 classifyCall,
39 extractResponseJson,
40 holdsForTasks,
41 isAckCommand,
42 normalizeState,
43 parseRetrospectiveConfig,
44 QUEUED_NOTICE,
45} from './logic.ts'
46import type { RetroConfig } from './logic.ts'
47
48/** Added to a classic.PostToolUse / classic.Stop event so the retro command hooks exit early. */
49export const RETRO_OWNED_FLAG = 'to_mod_retro'
50
51const MAX_RETRO_AGENTS = 20
52
53/** The tools the retro hooks care about: advance_item, complete_tree, and Bash (the retro-ack.mjs call). Matchers keep this feature's hooks distinct from the other features' registrations. */
54const RETRO_TOOLS = /^(?:Bash|mcp__.*task-orchestrator.*__(?:advance_item|complete_tree))$/
55
56type Ownership = { owned: boolean; cfg: RetroConfig; ancestorId: string | null }
57type Held = { state: RetroState; version: number }
58
59async function ownership($: EngineInterface): Promise<Ownership> {
60 const none: Ownership = { owned: false, cfg: parseRetrospectiveConfig(null), ancestorId: null }
61 if ((await $.env.get('TASK_ORCHESTRATOR_MODE')) === 'headless-iteration') return none
62 // Same first step as the command hooks' config-locator: $AGENT_CONFIG_DIR's config wins when readable.
63 // The rest of the locator (walk-up, main checkout of a worktree, user-level config) is not mirrored:
64 // the mod reads only the cwd-relative config, so ownership is claimed only for what it can see.
65 let text: string | null = null
66 const dir = await $.env.get('AGENT_CONFIG_DIR')
67 if (typeof dir === 'string' && dir.length > 0) {
68 try {
69 text = await $.fs.read(`${dir.replace(/[\\/]+$/, '')}/${CONFIG_PATH}`)
70 } catch {
71 text = null
72 }
73 }
74 if (text === null) {
75 try {
76 text = await $.fs.read(CONFIG_PATH)
77 } catch {
78 return none
79 }
80 }
81 const cfg = parseRetrospectiveConfig(text)
82 const section = readSection(text, 'project', { blockOnly: true })
83 const id = section ? scalar(section.lines, 'rootId') : null
84
85 return { owned: cfg.mode !== 'off', cfg, ancestorId: typeof id === 'string' && id.length > 0 ? id : null }
86}
87
88async function load($: EngineInterface): Promise<Held> {
89 const got = await $.state.get({ plugin: 'task-orchestrator-mod', key: 'retro' } as const)
90
91 return { state: normalizeState(got.value), version: got.version }
92}
93
94async function save($: EngineInterface, state: RetroState, ifVersion?: number): Promise<boolean> {
95 const res = await $.state.set({ plugin: 'task-orchestrator-mod', key: 'retro' } as const, state, ifVersion === undefined ? undefined : { ifVersion })
96
97 return res.isSet
98}
99
100/** True for a retrospective agent the mod spawned and for any agent it descends from one. */
101async function isRetroAgent($: EngineInterface, agentId: string | undefined, state: RetroState): Promise<boolean> {
102 const mine = state.retroAgents ?? []
103 if (agentId === undefined || mine.length === 0) return false
104 if (mine.includes(agentId)) return true
105 const list = await $.agent.list()
106 let cur = agentId
107 for (let i = 0; i < 10; i++) {
108 const parent = list.find(a => a.id === cur)?.parentId
109 if (parent === undefined) return false
110 if (mine.includes(parent)) return true
111 cur = parent
112 }
113
114 return false
115}
116
117function parseResponse(r: { text?: string; result?: unknown }): Record<string, unknown> | null {
118 if (typeof r.text === 'string') {
119 try {
120 const parsed = JSON.parse(r.text)
121 if (parsed !== null && typeof parsed === 'object') return parsed as Record<string, unknown>
122 } catch {
123 // fall through to the structured result
124 }
125 }
126
127 return extractResponseJson(r.result)
128}
129
130function withContext<T extends object>(r: T, line: string): T {
131 const ctx = (r as { context?: readonly string[] }).context ?? []
132
133 return { ...r, context: [...ctx, line] }
134}
135
136function withBlock<T extends object>(r: T, text: string): T {
137 const prior = (r as { block?: string }).block
138
139 return { ...r, block: prior ? `${prior}\n\n${text}` : text }
140}
141
142export function registerRetro(on: On): void {
143 on('tool.call', { tool: RETRO_TOOLS }, async ($, e, next) => {
144 if (next.origin?.plugin === PLUGIN) return next(e)
145 const name = toToolName(e.tool)
146 const command = e.tool === 'Bash' ? (e as unknown as { command?: unknown }).command : undefined
147 const isAck = isAckCommand(command)
148 if (name !== 'advance_item' && name !== 'complete_tree' && !isAck) return next(e)
149
150 let own: Ownership
151 let held: Held
152 let retro: boolean
153 try {
154 own = await ownership($)
155 if (!own.owned) return next(e)
156 held = await load($)
157 retro = await isRetroAgent($, e.agentId, held.state)
158 } catch {
159 return next(e)
160 }
161
162 if (isAck) {
163 const r = await next(e)
164 // A failed or denied ack run changed nothing on the command side; do not ack the mod's state.
165 if (r.deny !== undefined || r.isError === true) return r
166 try {
167 const now = await $.clock.now()
168 const cur = await load($)
169 await save($, ackState(cur.state, now, !retro))
170 } catch {
171 // an ack that cannot be recorded leaves the cooldown to run out
172 }
173
174 return r
175 }
176 if (retro) return next(e)
177
178 const r = await next(e)
179 if (r.deny !== undefined || r.isError === true) return r
180 try {
181 const input = e as unknown as Record<string, unknown>
182 const classified = classifyCall(name as string, input, parseResponse(r))
183 if (classified === null) return r
184 const now = await $.clock.now()
185 const cur = await load($)
186 const out = applyCall(cur.state, classified, own.cfg, now, e.agentId !== undefined)
187 if (out.state !== cur.state) await save($, out.state)
188 if (out.action === 'queue') return withContext(r, QUEUED_NOTICE)
189 if (out.action === 'nudge') return withContext(r, buildNudge(out.roots))
190 } catch {
191 return r
192 }
193
194 return r
195 })
196
197 on('classic.PostToolUse', { tool_name: RETRO_TOOLS }, async ($, e, next) => {
198 const name = toToolName(e.tool_name)
199 if ((name !== 'advance_item' && name !== 'complete_tree') || next.origin?.plugin === PLUGIN) return next(e)
200 try {
201 if (!(await ownership($)).owned) return next(e)
202 } catch {
203 return next(e)
204 }
205
206 return next({ ...e, [RETRO_OWNED_FLAG]: true } as never)
207 })
208
209 on('classic.Stop', async ($, e, next) => {
210 if (next.origin?.plugin === PLUGIN) return next(e)
211 let own: Ownership
212 try {
213 own = await ownership($)
214 if (!own.owned) return next(e)
215 } catch {
216 return next(e)
217 }
218 // Other plugins' Stop hooks run first, whatever this one decides.
219 const r = await next({ ...e, [RETRO_OWNED_FLAG]: true } as never)
220 if (e.stop_hook_active) return r
221
222 try {
223 const { state, version } = await load($)
224 const now = await $.clock.now()
225
226 const roots = state.pendingDispatch ?? []
227 if (roots.length > 0) {
228 // Hold while work runs in the background; the Stop that ends the wake-up turn re-checks.
229 const busy = e.background_tasks !== undefined ? holdsForTasks(e.background_tasks) : holdsForTasks(await $.agent.list())
230 if (busy) return r
231 // Claim it: only the Stop whose write lands at `version` spawns.
232 if (!(await save($, { ...state, pendingDispatch: null }, version))) return r
233
234 let agentId: string | undefined
235 let spawnedOk = false
236 try {
237 const spawned = await $.agent.spawn({
238 subagentType: 'general-purpose',
239 model: 'sonnet',
240 description: 'Session retrospective',
241 prompt: buildSpawnPrompt(roots, own.ancestorId),
242 })
243 // Core sets agentId on a started agent; a denial carries `deny` and starts none.
244 spawnedOk = spawned.deny === undefined
245 agentId = spawned.deny === undefined ? spawned.agentId : undefined
246 } catch {
247 spawnedOk = false
248 }
249 if (!spawnedOk) return withBlock(r, buildDispatch(roots, own.ancestorId))
250
251 const cur = (await load($)).state
252 await save($, {
253 ...cur,
254 handledAt: now,
255 dispatched: { ...(agentId !== undefined && { agentId }), roots: [...roots], at: now },
256 ...(agentId !== undefined && { retroAgents: [...(cur.retroAgents ?? []), agentId].slice(-MAX_RETRO_AGENTS) }),
257 })
258
259 return r
260 }
261
262 const nudgeRoots = backstopRoots(state, own.cfg, now)
263 if (nudgeRoots === null) return r
264 await save($, afterBackstop(state, now))
265
266 return withBlock(r, buildNudge(nudgeRoots))
267 } catch {
268 return r
269 }
270 })
271}
272src/band/model.ts 113 lines1// Pure selection and text for the band and the status line (T4 ec1e2f91). No `$`, no atoms: the
2// hooks in index.ts feed it a snapshot and draw what it returns.
3import type { GateInfo, GraphNode, GraphSnapshot } from '../../types'
4import { cardsOf, num } from '../graph-pane/model.ts'
5
6export interface BandModel {
7 /** The band's line; absent when hidden or nothing is in flight. */
8 band?: string
9 /** The compact status-line text; absent when nothing is in flight (clears the line). */
10 status?: string
11}
12
13/** Columns the Hide button and its spacing take beside the band text. */
14const BUTTON_COLUMNS = 10
15const MIN_TITLE = 8
16
17const ACTIVE = new Set(['work', 'review'])
18
19/** The in-flight leaves, in the order the band picks from: work before review, deeper first, snapshot order. */
20export function inFlight(snapshot: GraphSnapshot | null): GraphNode[] {
21 if (!snapshot) return []
22 const active = snapshot.nodes.filter(node => ACTIVE.has(node.role))
23 const byId = new Map(snapshot.nodes.map(node => [node.id, node]))
24 // A node is a container when an active descendant exists: mark every ancestor of every active node.
25 const containers = new Set<string>()
26 for (const node of active) {
27 let parent = node.parentId === null ? undefined : byId.get(node.parentId)
28 for (let guard = 0; parent && guard < 1000; guard += 1) {
29 containers.add(parent.id)
30 parent = parent.parentId === null ? undefined : byId.get(parent.parentId)
31 }
32 }
33 const rank = (node: GraphNode): number => (node.role === 'work' ? 0 : 1)
34
35 return active
36 .map((node, index) => ({ node, index }))
37 .filter(({ node }) => !containers.has(node.id))
38 .sort((a, b) => rank(a.node) - rank(b.node) || b.node.depth - a.node.depth || a.index - b.index)
39 .map(({ node }) => node)
40}
41
42/** `2/3 work notes`, `work` (nothing required), or the role (no gate entry); ` ✓` when it can advance. */
43export function gateText(node: GraphNode, gate: GateInfo | undefined): string {
44 if (!gate) return node.role
45 const base = gate.required > 0 ? `${gate.filled}/${gate.required} ${gate.phase} notes` : gate.phase
46
47 return gate.canAdvance ? `${base} ✓` : base
48}
49
50/** `title` cut to `room` characters with a trailing ellipsis. */
51export function truncate(title: string, room: number): string {
52 if (title.length <= room) return title
53 if (room <= 1) return '…'
54
55 return `${title.slice(0, room - 1)}…`
56}
57
58/** Most labels shown per segment of the smart line; the rest read ` +k`. */
59export const SEGMENT_MAX = 3
60
61const byPlanNumber = (a: string, b: string): number => num(a) - num(b) || (a < b ? -1 : a > b ? 1 : 0)
62
63function segment(name: string, labels: string[]): string | undefined {
64 if (labels.length === 0) return undefined
65 const sorted = [...labels].sort(byPlanNumber)
66 const more = sorted.length > SEGMENT_MAX ? ` +${sorted.length - SEGMENT_MAX}` : ''
67
68 return `${name}: ${sorted.slice(0, SEGMENT_MAX).join(', ')}${more}`
69}
70
71/**
72 * The band line while the graph is scoped to an item (D1 a): `◉ <id8> · <d>/<n> done · ready: … · waiting: …`.
73 * n counts the scope's cards (not the scope itself), d the terminal ones (cancelled counts as done);
74 * ready are queue cards with every blocker satisfied, waiting the open cards with an open blocker.
75 * Undefined on the root overview, with no snapshot, or with no cards (the in-flight line then shows).
76 */
77export function smartLine(snapshot: GraphSnapshot | null, bodyColumns = 80): string | undefined {
78 if (!snapshot || snapshot.overview === true || snapshot.scopeId === null) return undefined
79 const { cards } = cardsOf(snapshot)
80 if (cards.length === 0) return undefined
81 const done = cards.filter(c => c.role === 'terminal').length
82 const parts = [
83 `◉ ${snapshot.scopeId.slice(0, 8)}`,
84 `${done}/${cards.length} done`,
85 segment('ready', cards.filter(c => c.ready).map(c => c.label)),
86 segment('waiting', cards.filter(c => c.role !== 'terminal' && c.openBlockers.length > 0).map(c => c.label)),
87 ].filter((p): p is string => p !== undefined)
88
89 return truncate(parts.join(' · '), Math.max(MIN_TITLE, bodyColumns - BUTTON_COLUMNS))
90}
91
92/** What the band draws: the smart line when scoped to an item, else the in-flight line (undefined: nothing). */
93export function bandLine(snapshot: GraphSnapshot | null, bodyColumns = 80): string | undefined {
94 return smartLine(snapshot, bodyColumns) ?? bandModel(snapshot, false, bodyColumns).band
95}
96
97export function bandModel(snapshot: GraphSnapshot | null, hidden: boolean, bodyColumns = 80): BandModel {
98 const list = inFlight(snapshot)
99 const first = list[0]
100 if (!snapshot || !first) return {}
101
102 const id8 = first.id.slice(0, 8)
103 const more = list.length > 1 ? ` (+${list.length - 1})` : ''
104 const gate = gateText(first, snapshot.gates[first.id])
105 const status = `TO ◉ ${id8} ${gate.replace(' notes', '')}${more}`
106 if (hidden) return { status }
107
108 const fixed = `◉ ${id8} — ${gate}${more}`.length
109 const room = Math.max(MIN_TITLE, bodyColumns - BUTTON_COLUMNS - fixed)
110
111 return { band: `◉ ${id8} ${truncate(first.title, room)} — ${gate}${more}`, status }
112}
113src/shared/constants.ts 23 lines1// Names shared by every feature of the mod. Facts behind them were verified live in spike
2// 11960f01 (see its implementation-notes).
3
4/** This plugin's name: `next.origin.plugin` equals it when a dispatch was raised by our own `$` call. */
5export const PLUGIN = 'task-orchestrator-mod'
6
7/** The TO MCP server's name as `/mcp` lists it; what `$.mcp.call` takes. */
8export const TO_SERVER = 'mcp-task-orchestrator'
9
10/** Matches the tool names of the TO MCP server as the engine spells them (`mcp__<server>__<tool>`). */
11export const TO_TOOL = /^mcp__.*task-orchestrator.*__/
12
13/**
14 * Field a `classic.<Event>` hook adds to the event before `next`, so a TO command hook reading it
15 * on stdin can tell the mod already handled that event and exit early (selective dedupe).
16 */
17export const MOD_ACTIVE_FLAG = 'to_mod_active'
18
19/** The bare TO tool name (`advance_item`) of an engine tool name, or null when it is not a TO tool. */
20export function toToolName(tool: string): string | null {
21 return TO_TOOL.test(tool) ? tool.replace(TO_TOOL, '') : null
22}
23src/shared/config.ts 18 lines1// Project config parsing. Pure: the caller reads the file itself with `$.fs.read(CONFIG_PATH)` in
2// the module that owns `$` — the plugin validator follows `$` only within one file, so a helper
3// that takes `$` across an import is refused. src/lib/yaml-lite.mjs is a copy of
4// claude-plugins/task-orchestrator/hooks/yaml-lite.mjs — keep the two in sync.
5import { readSection, scalar } from '../lib/yaml-lite.mjs'
6
7/** The project config, relative to the session's working directory. */
8export const CONFIG_PATH = '.taskorchestrator/config.yaml'
9
10/** `project.rootId` from config text, or null when there is none. */
11export function parseProjectRootId(configText: string | null): string | null {
12 if (!configText) return null
13 const section = readSection(configText, 'project', { blockOnly: true })
14 const rootId = section ? scalar(section.lines, 'rootId') : null
15
16 return typeof rootId === 'string' && rootId.length > 0 ? rootId : null
17}
18src/call-shape/rewrite.ts 75 lines1// Pure call-shape repair for TO tool calls (T7). No `$`, so it is unit-testable: the caller reads
2// the project rootId and hands it in.
3
4export type Rewrite = { input: Record<string, unknown>; changes: string[] }
5
6type Input = Record<string, unknown>
7
8const absent = (v: unknown): boolean => v === undefined || v === null
9
10/**
11 * A `tags` or `type` filter of any value means the caller is looking up a specific kind of item
12 * (process-global containers, trends, observations, retrospectives, proposals, personal roots),
13 * which may live outside the project root; such listings stay unscoped. An explicit
14 * `ancestorId: ""` also passes through untouched (the server reads blank as unscoped).
15 */
16function hasKindFilter(input: Input): boolean {
17 return !absent(input.tags) || !absent(input.type)
18}
19
20/** Whether rule A (inject `ancestorId`) is in play for this tool and input, before the skip conditions. */
21function isScopable(tool: string, input: Input): boolean {
22 switch (tool) {
23 case 'query_items':
24 return input.operation === 'search' && input.query === undefined && input.depth !== 0
25 case 'get_next_item':
26 case 'get_blocked_items':
27 return true
28 case 'get_context':
29 return input.itemId === undefined && input.mode !== 'item'
30 default:
31 return false
32 }
33}
34
35/** True when the tool's own caller is this plugin (its `$.mcp.call` reads), which are never repaired. */
36export function isOwnCall(originPlugin: string | undefined, plugin: string): boolean {
37 return originPlugin === plugin
38}
39
40/**
41 * Repairs one TO tool call. `tool` is the bare TO tool name. Returns the (possibly new) input and
42 * one change label per applied rewrite; `input` is the same object when nothing changed.
43 */
44export function rewriteCall(tool: string, input: Input, rootId: string | null): Rewrite {
45 const changes: string[] = []
46 let out = input
47
48 if (rootId !== null && isScopable(tool, input) && input.ancestorId === undefined && absent(input.parentId) && !hasKindFilter(input)) {
49 out = { ...out, ancestorId: rootId }
50 changes.push(`${tool} +ancestorId=${rootId.slice(0, 8)}`)
51 }
52
53 if (tool === 'manage_notes' && input.operation === 'upsert' && Array.isArray(input.notes)) {
54 const actor = input.actor
55 if (typeof actor === 'object' && actor !== null && !Array.isArray(actor)) {
56 const filled: number[] = []
57 const notes = input.notes.map((n: unknown, i: number) => {
58 if (typeof n === 'object' && n !== null && !Array.isArray(n) && (n as Input).actor === undefined) {
59 filled.push(i)
60
61 return { ...(n as Input), actor: { ...(actor as Input) } }
62 }
63
64 return n
65 })
66 if (filled.length > 0) {
67 out = { ...out, notes }
68 changes.push(`manage_notes actor -> notes[${filled.join(',')}]`)
69 }
70 }
71 }
72
73 return { input: out, changes }
74}
75src/graph-pane/activity.ts 207 lines1// Who is working on which item, and which items changed a moment ago. Pure: index.ts owns the
2// atoms and the hooks, and the render never reads the clock (prune-on-write plus timers expire entries).
3import type { GraphActivity, GraphAgentInfo, GraphGateBlock, GraphWorker } from '../../types'
4import { toToolName } from '../shared/constants.ts'
5
6/** A worker stays on an item this long after its last call about it. */
7export const WORKING_TTL_MS = 120_000
8/** An item reads as recently changed this long after a TO write touched it. */
9export const RECENT_MS = 30_000
10/** A repeat call from the same worker inside this window writes nothing. */
11export const ACTIVITY_REFRESH_MS = 10_000
12/** Most workers kept per item, and most agents remembered. */
13export const MAX_WORKERS = 8
14export const MAX_AGENTS = 50
15
16export const emptyActivity = (): GraphActivity => ({ working: {}, changed: {} })
17
18type Obj = Record<string, unknown>
19const isObj = (v: unknown): v is Obj => typeof v === 'object' && v !== null && !Array.isArray(v)
20const text = (v: unknown): string | undefined => (typeof v === 'string' && v.length > 0 ? v : undefined)
21
22/** An item id a TO call named, with the actor id on that element (or the call). */
23export interface Touch {
24 id: string
25 actor?: string
26}
27
28const actorId = (v: unknown): string | undefined => (isObj(v) ? text(v.id) : undefined)
29
30/**
31 * The ids a TO tool call names: top-level `itemId`, and `itemId`/`id` of each element of
32 * transitions, claims, notes and items. Unresolved (may be prefixes). Deduplicated, call order.
33 */
34export function touchedIds(input: unknown): Touch[] {
35 if (!isObj(input)) return []
36 const top = actorId(input.actor)
37 const out = new Map<string, Touch>()
38 const add = (id: string | undefined, actor: string | undefined): void => {
39 if (id === undefined || out.has(id)) return
40 out.set(id, actor !== undefined ? { id, actor } : { id })
41 }
42 add(text(input.itemId), top)
43 for (const key of ['transitions', 'claims', 'notes', 'items']) {
44 const list = input[key]
45 if (!Array.isArray(list)) continue
46 for (const el of list) {
47 if (isObj(el)) add(text(el.itemId) ?? text(el.id), actorId(el.actor) ?? top)
48 }
49 }
50
51 return [...out.values()]
52}
53
54/** Prefers the seat agent.spawn recorded, then the actor id's prefix before `:`, then `agent` / `main`. */
55export function seatOf(agentId: string | undefined, agents: Record<string, GraphAgentInfo>, actor?: string): string {
56 const known = agentId !== undefined ? agents[agentId]?.seat : undefined
57 if (known !== undefined && known.length > 0) return known
58 const prefix = actor !== undefined ? actor.split(':')[0] : undefined
59 if (prefix !== undefined && prefix.length > 0) return prefix
60
61 return agentId !== undefined ? 'agent' : 'main'
62}
63
64/** `claude-opus-5-5-20260101` -> `opus-5-5`; an alias such as `opus` stays as it is. */
65export function shortModel(model: string | undefined): string | undefined {
66 if (model === undefined || model.length === 0) return undefined
67 const short = model.replace(/^claude-/, '').replace(/-\d{8}$/, '').replace(/\[.*\]$/, '')
68
69 return short.length > 0 ? short : model
70}
71
72/** Resolves a call's id against the snapshot's ids: exact, or a prefix (>= 4 chars) that matches exactly one. */
73export function resolveId(raw: string, known: readonly string[]): string | null {
74 if (known.includes(raw)) return raw
75 if (raw.length < 4) return null
76 const hits = known.filter(id => id.startsWith(raw))
77
78 return hits.length === 1 ? (hits[0] as string) : null
79}
80
81/** A gate-blocked advance_item marks its item this long (or until a later TO write on it). */
82export const GATE_BLOCK_MS = 600_000
83
84/** One gate-blocked transition of an advance_item result. */
85export interface GateBlockRow {
86 itemId: string
87 targetRole?: string
88 missing: string[]
89}
90
91/**
92 * The gate-blocked rows of an advance_item result as the model reads it (`ran.text`, the JSON text):
93 * rows with `applied: false` and `errorCode: 'gate_blocked'`. Any other tool, or text that is not
94 * that JSON, gives []. `missingNotes` entries may be `{ key }` objects or plain strings.
95 */
96export function gateBlocks(tool: string, text: unknown): GateBlockRow[] {
97 if (toToolName(tool) !== 'advance_item' || typeof text !== 'string') return []
98 let parsed: unknown
99 try {
100 parsed = JSON.parse(text)
101 } catch {
102 return []
103 }
104 const rows = isObj(parsed) && Array.isArray(parsed.results) ? parsed.results : []
105 const out: GateBlockRow[] = []
106 for (const row of rows) {
107 if (!isObj(row) || row.applied !== false || row.errorCode !== 'gate_blocked') continue
108 const itemId = typeof row.itemId === 'string' ? row.itemId : undefined
109 if (itemId === undefined) continue
110 const missing = Array.isArray(row.missingNotes)
111 ? row.missingNotes.map(m => (typeof m === 'string' ? m : isObj(m) && typeof m.key === 'string' ? m.key : undefined)).filter((k): k is string => k !== undefined)
112 : []
113 const target = typeof row.targetRole === 'string' ? row.targetRole : undefined
114 out.push({ itemId, missing, ...(target !== undefined ? { targetRole: target } : {}) })
115 }
116
117 return out
118}
119
120/** Drops workers silent for > WORKING_TTL_MS, changes older than RECENT_MS and gate blocks older than GATE_BLOCK_MS. Same object when nothing expired. */
121export function pruneActivity(act: GraphActivity, now: number): GraphActivity {
122 let dirty = false
123 const working: Record<string, GraphWorker[]> = {}
124 for (const [id, list] of Object.entries(act.working)) {
125 const live = list.filter(w => now - w.at <= WORKING_TTL_MS)
126 if (live.length !== list.length) dirty = true
127 if (live.length > 0) working[id] = live
128 }
129 const changed: Record<string, number> = {}
130 for (const [id, at] of Object.entries(act.changed)) {
131 if (now - at <= RECENT_MS) changed[id] = at
132 else dirty = true
133 }
134 // A value written before gate blocks existed has no `blocked`: tolerated, and left without one.
135 const blocked: Record<string, GraphGateBlock> = {}
136 for (const [id, b] of Object.entries(act.blocked ?? {})) {
137 if (now - b.at <= GATE_BLOCK_MS) blocked[id] = b
138 else dirty = true
139 }
140
141 return dirty ? withBlocked({ working, changed }, blocked) : act
142}
143
144const withBlocked = (act: GraphActivity, blocked: Record<string, GraphGateBlock>): GraphActivity =>
145 Object.keys(blocked).length > 0 ? { ...act, blocked } : { working: act.working, changed: act.changed }
146
147/**
148 * Records this session's gate-blocked rows on their items (prunes first). Returns the SAME object when
149 * every row is already recorded with the same keys and target under ACTIVITY_REFRESH_MS ago.
150 */
151export function recordGateBlocks(act: GraphActivity, rows: readonly { id: string; missing: string[]; target?: string }[], now: number): GraphActivity {
152 const base = pruneActivity(act, now)
153 const blocked = { ...(base.blocked ?? {}) }
154 let dirty = false
155 for (const r of rows) {
156 const cur = blocked[r.id]
157 if (cur !== undefined && cur.target === r.target && cur.missing.join('\u0000') === r.missing.join('\u0000') && now - cur.at < ACTIVITY_REFRESH_MS) continue
158 blocked[r.id] = { at: now, missing: [...r.missing], ...(r.target !== undefined ? { target: r.target } : {}) }
159 dirty = true
160 }
161
162 return dirty ? withBlocked(base, blocked) : base
163}
164
165/**
166 * Records one call: the worker on the item and, when `changed`, the change time. Prunes first. Returns
167 * the SAME object when the call adds nothing new (same worker, same seat and model, last seen under
168 * ACTIVITY_REFRESH_MS ago), so the caller can skip the write.
169 */
170export function recordActivity(
171 act: GraphActivity,
172 entry: { id: string; agentId: string; seat: string; model?: string; changed: boolean },
173 now: number,
174): GraphActivity {
175 const base = pruneActivity(act, now)
176 const list = base.working[entry.id] ?? []
177 const mine = list.find(w => w.agentId === entry.agentId)
178 const workerStale = mine === undefined || mine.seat !== entry.seat || mine.model !== entry.model || now - mine.at >= ACTIVITY_REFRESH_MS
179 const lastChange = base.changed[entry.id]
180 const changeStale = entry.changed && (lastChange === undefined || now - lastChange >= ACTIVITY_REFRESH_MS)
181 // A later TO write on a gate-blocked item clears its mark (the hook records a call's own blocks after this).
182 const unblock = entry.changed && base.blocked?.[entry.id] !== undefined
183 if (!workerStale && !changeStale && !unblock) return base
184
185 const worker: GraphWorker = { agentId: entry.agentId, seat: entry.seat, at: now, ...(entry.model !== undefined ? { model: entry.model } : {}) }
186 const next = workerStale ? [...list.filter(w => w.agentId !== entry.agentId), worker].slice(-MAX_WORKERS) : list
187 const { [entry.id]: _cleared, ...restBlocked } = base.blocked ?? {}
188
189 return withBlocked(
190 {
191 working: workerStale ? { ...base.working, [entry.id]: next } : base.working,
192 changed: changeStale ? { ...base.changed, [entry.id]: now } : base.changed,
193 },
194 unblock ? restBlocked : (base.blocked ?? {}),
195 )
196}
197
198/** Remembers a spawned agent, oldest dropped past MAX_AGENTS. Same object when unchanged. */
199export function rememberAgent(agents: Record<string, GraphAgentInfo>, agentId: string, info: GraphAgentInfo): Record<string, GraphAgentInfo> {
200 const cur = agents[agentId]
201 if (cur !== undefined && cur.seat === info.seat && cur.model === info.model) return agents
202 const { [agentId]: _drop, ...rest } = agents
203 const entries = Object.entries({ ...rest, [agentId]: info })
204
205 return Object.fromEntries(entries.slice(-MAX_AGENTS))
206}
207