Live terminal visualization of Omnigraph operations


omnigraph-viz is a Claude Code mod that shows, live in your terminal, what an agent does to an Omnigraph graph. Every omnigraph command Claude runs lands in a pane beside the conversation: a browser that follows the agent through the graph, and a feed of every operation.
Epic › Offline editing & sync › contains › ….● label, links as thin dotted lines.◆ query ready · 5).NEW.enter shows the command and its output.◉ og · main · 14 ops · +3n./og replay <ledger> plays a recorded session offline, with no Omnigraph needed: a demo that can't fail on stage.The mod only watches. It never changes the agent's command or what the agent reads back, and a bug in the viz never turns into a tool error.
omnigraph CLI on your PATH, or its path in the omnigraphBin option.In Claude Code:
/plugin marketplace add pronskiy/omnigraph-claude-mod
/plugin install omnigraph-viz@omnigraph-claude-mod
Or run it from a clone:
git clone https://github.com/pronskiy/omnigraph-claude-mod
cd omnigraph-claude-mod
claude --plugin-dir ./omnigraph-viz
Open the pane with /og, then work with Omnigraph as usual: ask Claude to query or change a graph. The mod records operations from the first omnigraph command whether or not the pane is open.
/og | toggle the pane (it opens with the keyboard) |
/og replay <ledger.jsonl> [speed] | replay a recorded session; 0.5 is half speed, 2 double |
/og stop | end a replay, back to the live session |
/og clear | reset the Feed (the graph stays) |
g / f | Graph / Feed tab |
| click the graph | hand it the keys (once; until then the arrows move between the pane's buttons) |
| ↑ / ↓ | select in the rightmost column |
| → or ⏎ or ⌥↓ | open the selection as a new column |
| ← or ⌫ or ⌥↑ | back one column |
l | follow the agent again |
enter (Feed) | show an operation's command and output; select it to see it in the browser |
esc | keys back to the prompt |
For Finder's ⌘↓ / ⌘↑ in iTerm2, map them to send ⌥↓ / ⌥↑ (Profile → Keys → Key Mappings → Send Escape Sequence [1;3B / [1;3A).
From a clone, start Claude Code with the plugin and replay the recorded demo, an agent triaging a bug bash on the dev-graph cookbook:
/og replay demos/dev-graph-triage.jsonl
To run it live with a real agent, see the runbook in docs/demo.md: scripts/setup-demo-graph.sh builds the store, and demos/dev-graph-triage.prompt.txt is the prompt.
| Option | Default | What it does |
|---|---|---|
omnigraphBin | omnigraph | Name or absolute path of the CLI. When the bare name isn't on Claude Code's PATH, the mod also asks your login shell and tries the usual install folders. |
maxGhostNodes | 2000 | Largest graph loaded whole; a bigger one keeps its best-connected nodes |
ledgerDir | .omnigraph-viz | Where session ledgers go, relative to the session directory |
enrich | true | Run omnigraph … --json in the background for graph detail (export once, commit changes after writes) |
debug | false | Log viz errors as dim lines in the transcript |
With --plugin-dir, pass options as settings:
claude --plugin-dir ./omnigraph-viz \
--settings '{"pluginConfigs":{"omnigraph-viz":{"options":{"omnigraphBin":"/Users/you/.local/bin/omnigraph"}}}}'
tool.call hook spots omnigraph invocations in Bash commands (through cd … &&, env prefixes and pipes). It lets each one run untouched and parses its human or --json output.omnigraph calls on the same store: export once to load the graph, and commit list / commit changes after each write. These give the exact nodes and edges that changed. Reads light the nodes their result rows name..omnigraph-viz/. The mod adds a .gitignore there, because ledgers hold command output.npm install
claude --plugin-dir ./omnigraph-viz # once: Claude Code lays the mod API types that tsc needs next to the plugin
npm run check # tsc, claude plugin validate, claude plugin test
Hot reload: with claude --plugin-dir ./omnigraph-viz, saving a file reloads the mod.
.claude-plugin/marketplace.json the marketplace (one plugin)
omnigraph-viz/
hooks/register.tsx wiring: hooks, commands, the pane, state
hooks/explore.tsx the browser's Client: draws rows, posts keys
src/capture/ command tokenizer, classifier, output parser
src/enrich/ background omnigraph calls, read highlighting
src/model/ ledger, reducer, view state
src/graph/ graph model, effects, palette
src/browse/ browser model and star layout (pure)
src/render/ canvas, browser drawing, pane, feed
src/replay/ /og replay player
tests/ claude plugin test suites
demos/ a recorded demo ledger and its prompt
scripts/ dev-only Node and shell scripts
The tests stand in for the engine: they answer every $ call themselves and mount the pane with $.ui.mount. They need no file system, network or omnigraph.
hooks/register.tsx 602 lines1import { atom, read, update } from 'claude-code'
2import type { EngineInterface, Register } from 'claude-code'
3
4import { beginOps, endLines } from '../src/capture/capture'
5import { parseOutput } from '../src/capture/parse-output'
6import { findInvocations } from '../src/capture/tokenize'
7import {
8 commitChangesArgv, commitListArgv, exportArgv, parseChanges, parseCommits, pickCommit, schemaArgv, schemaSourceOf, targetKey,
9} from '../src/enrich/cli'
10import { buildMatchIndex, matchCells } from '../src/enrich/highlight'
11import { binCandidates, isPlainName, pathFromShell } from '../src/enrich/locate'
12import { JobQueue } from '../src/enrich/queue'
13import { foldExport, parseEdgeEnds } from '../src/graph/ghost'
14import { LedgerWriter, decodeLedger, ledgerPathFor } from '../src/model/ledger'
15import type { LedgerIO } from '../src/model/ledger'
16import { emptyViz, reduce } from '../src/model/reduce'
17import { INITIAL_UI, INITIAL_VIEW, toView } from '../src/model/state'
18import { BrailleCanvas } from '../src/render/braille'
19import { statusText } from '../src/render/format'
20import { browserRowsFor } from '../src/render/graph-view'
21import type { GraphViewData } from '../src/render/graph-view'
22import { drawPane } from '../src/render/pane'
23import { actionFor } from '../src/render/explore'
24import { opName, Scene } from '../src/render/scene'
25import { OG_HINT, parseOgArgs } from '../src/replay/command'
26import { Player, readReplay } from '../src/replay/player'
27import { basename, resolvePath } from '../src/util/path'
28import type { GraphState, LedgerLine, OpEvent, OpIntent, UiState, VizState } from '../types'
29
30const MOD_VERSION = '0.1.0'
31const PANE = 'omnigraph'
32const FRAME_MS = 33
33
34interface Config {
35 bin: string
36 ledgerDir: string
37 maxGhostNodes: number
38 isEnriching: boolean
39 isDebug: boolean
40}
41
42/** What the pane can show: a reduced state and the scene drawn from the same lines. */
43interface Show {
44 viz: VizState
45 scene: Scene
46}
47
48interface Session extends Show {
49 cwd: string
50 /** Created by the first op, so a session without omnigraph leaves no files. */
51 ledger: LedgerWriter | undefined
52 /** After `session.end`: background jobs may still finish, but the status line stays clear. */
53 isEnded?: boolean
54}
55
56/** A recorded ledger playing in the pane instead of the live session (which carries on, unseen). */
57interface Replay extends Show {
58 name: string
59 player: Player
60 hasGraph: boolean
61 timer: { cancel: () => void } | undefined
62}
63
64// $.state: small, survives hot reload, redraws its readers. Declared here so the validator can see them.
65const uiAtom = atom({ plugin: 'omnigraph-viz', key: 'ui' } as const, INITIAL_UI)
66const viewAtom = atom({ plugin: 'omnigraph-viz', key: 'view' } as const, INITIAL_VIEW)
67const ledgerPathAtom = atom({ plugin: 'omnigraph-viz', key: 'ledgerPath' } as const, null)
68
69// Module memory: rebuilt from the ledger on every (re)load.
70let config: Config = { bin: 'omnigraph', ledgerDir: '.omnigraph-viz', maxGhostNodes: 2000, isEnriching: true, isDebug: false }
71let session: Session | undefined
72let replay: Replay | undefined
73let queue: JobQueue | undefined
74/** Bootstrap progress for the graph `graphKey` names; `ready` once a graph line exists. */
75let graphState: GraphState = 'none'
76let graphMessage: string | null = null
77let graphKey: string | null = null
78/** The browser's canvas while the Graph tab is on screen; the animation loop redraws it. */
79let canvas: BrailleCanvas | null = null
80let frameTimer: { cancel: () => void } | undefined
81let isFrameBusy = false
82/** Where `omnigraph` really is, once a bare name failed to start on Claude Code's PATH. */
83let resolvedBin: string | undefined
84/** HOME, read once: `null` until then. */
85let homeCache: string | undefined | null = null
86
87function io($: EngineInterface): LedgerIO {
88 return {
89 write: (path, text) => $.fs.write(path, text),
90 onError: err => debug($, err),
91 }
92}
93
94function debug($: EngineInterface, err: unknown): void {
95 if (config.isDebug) $.ui.log(`omnigraph-viz: ${err instanceof Error ? err.message : String(err)}`)
96}
97
98/** The pane shows the replay while one is on, else the live session. */
99function shown(): Show | undefined {
100 return replay ?? session
101}
102
103/** The one path every line takes, live or replayed: the reducer, then the scene. */
104function feed(show: Show, line: LedgerLine): void {
105 const before = show.viz.ghost
106 show.viz = reduce(show.viz, line)
107 show.scene.onLine(line, before, show.viz.ghost)
108}
109
110/** Every live change goes through here: ledger, then reducer and scene, then animation. */
111function apply($: EngineInterface, s: Session, line: LedgerLine): void {
112 s.ledger?.append(line, io($))
113 if (line.t === 'graph') graphKey = targetKey(line.target, line.cwd)
114 feed(s, line)
115 if (shown() === s && (line.t === 'graph' || line.t === 'diff' || line.t === 'touch')) startAnimation($)
116}
117
118/** The view and status line of whatever the pane shows. */
119async function publish($: EngineInterface): Promise<void> {
120 const r = replay
121 const view = r
122 ? toView(r.viz, r.viz.ghost ? 'ready' : r.hasGraph ? 'loading' : 'none', null, {
123 speed: r.player.speed,
124 progress: r.player.progress,
125 isDone: r.player.isDone,
126 })
127 : session
128 ? toView(session.viz, graphState, graphMessage)
129 : INITIAL_VIEW
130 await update($, viewAtom, () => view)
131 $.ui.status(view.counters.ops > 0 && !session?.isEnded ? statusText(view) : undefined)
132}
133
134/** On the first op: the ledger file, a `.gitignore` beside it, and the session line. */
135async function ensureLedger($: EngineInterface, s: Session, now: number): Promise<void> {
136 if (s.ledger) return
137 const path = ledgerPathFor(s.cwd, config.ledgerDir, now)
138 const ignore = `${path.slice(0, path.lastIndexOf('/'))}/.gitignore`
139 if (!(await $.fs.exists(ignore))) await $.fs.write(ignore, '*\n')
140 await update($, ledgerPathAtom, () => path)
141 s.ledger = new LedgerWriter(path)
142 apply($, s, { v: 1, t: 'session', at: now, cwd: s.cwd, mod: MOD_VERSION })
143}
144
145/**
146 * Rehydrates from the ledger `$.state` remembers (a hot reload), else starts
147 * a new one. The scene is laid out once from the final graph: replaying every
148 * diff's layout would make a reload as slow as the session was long.
149 */
150async function openSession($: EngineInterface, cwd: string): Promise<Session> {
151 const s: Session = { cwd, ledger: undefined, viz: emptyViz(), scene: new Scene() }
152 const previous = await read($, ledgerPathAtom)
153 if (previous === null) return s
154 try {
155 const { lines } = decodeLedger(await $.fs.read(previous))
156 let graphLine: Extract<LedgerLine, { t: 'graph' }> | undefined
157 for (const line of lines) {
158 s.viz = reduce(s.viz, line)
159 if (line.t === 'graph') graphLine = line
160 }
161 if (graphLine && s.viz.ghost) {
162 graphKey = targetKey(graphLine.target, graphLine.cwd)
163 graphState = 'ready'
164 s.scene.onLine({ ...graphLine, at: await $.clock.now(), graph: s.viz.ghost }, null, s.viz.ghost)
165 }
166 s.ledger = new LedgerWriter(previous, lines)
167 return s
168 } catch {
169 // file gone: start over
170 return { cwd, ledger: undefined, viz: emptyViz(), scene: new Scene() }
171 }
172}
173
174/**
175 * Finds `omnigraph` the way the person's terminal does: their login shell
176 * (whose profile usually adds ~/.local/bin and the like), then the usual
177 * install folders. Claude Code's own PATH may lack what its Bash tool has.
178 */
179async function locateOmnigraph($: EngineInterface): Promise<string | undefined> {
180 if (!isPlainName(config.bin)) return undefined
181 const shell = (await $.env.get('SHELL')) ?? '/bin/sh'
182 const viaShell = await $.process
183 .run([shell, '-lic', `command -v ${config.bin}`], { timeoutMs: 5000 })
184 .then(r => pathFromShell(r.stdout, config.bin))
185 .catch(() => undefined)
186 if (viaShell !== undefined) return viaShell
187 for (const path of binCandidates(await $.env.get('HOME'), config.bin)) {
188 if (await $.fs.exists(path)) return path
189 }
190 return undefined
191}
192
193/** Runs `omnigraph <args>`; if the binary does not start, locates it once and retries. */
194async function runOmnigraph($: EngineInterface, args: readonly string[], cwd: string) {
195 const run = (bin: string) => $.process.run([bin, ...args], { cwd, timeoutMs: 20_000 })
196 if (resolvedBin !== undefined) return run(resolvedBin)
197 try {
198 return await run(config.bin)
199 } catch (err) {
200 const found = await locateOmnigraph($)
201 if (found === undefined) {
202 throw new Error(
203 `can't find \`${config.bin}\` on Claude Code's PATH: set the omnigraphBin option to its full path (\`which ${config.bin}\`)`,
204 )
205 }
206 resolvedBin = found
207 return run(found)
208 }
209}
210
211function enqueue($: EngineInterface, at: number, run: () => Promise<void>): void {
212 queue ??= new JobQueue(() => $.clock.now(), err => debug($, err))
213 queue.push({ at, run })
214}
215
216/** Schema (for edge endpoint types) and export of the graph the op ran against → a `graph` line. */
217async function bootstrap($: EngineInterface, s: Session, intent: OpIntent): Promise<void> {
218 const run = (argv: string[]) => runOmnigraph($, argv.slice(1), intent.cwd)
219 try {
220 const schema = await run(schemaArgv(config.bin, intent.target)).catch(() => undefined)
221 const edgeEnds = parseEdgeEnds(schemaSourceOf(schema?.stdout ?? '') ?? '')
222 const out = await run(exportArgv(config.bin, intent.target))
223 if (out.exitCode !== 0) {
224 const failure = parseOutput(out.stdout, out.stderr, true)
225 throw new Error(failure.kind === 'error' ? failure.message : `exit ${out.exitCode}`)
226 }
227 // Over `$.process.run`'s 4 MiB cap the export is cut: start touched-only instead.
228 const graph = out.isStdoutTruncated
229 ? { nodes: [], edges: [], truncated: true, edgeEnds }
230 : foldExport(out.stdout, edgeEnds, config.maxGhostNodes)
231 apply($, s, { v: 1, t: 'graph', at: await $.clock.now(), target: intent.target, cwd: intent.cwd, graph })
232 graphState = 'ready'
233 graphMessage = null
234 } catch (err) {
235 graphState = 'failed'
236 graphMessage = (err instanceof Error ? err.message : String(err)).slice(0, 160)
237 }
238 await publish($)
239}
240
241/** The commit a write op made → its changes → a `diff` line. */
242async function enrichWrite($: EngineInterface, s: Session, op: OpEvent, branch: string, endedAt: number): Promise<void> {
243 if (!s.viz.ghost) return
244 const run = (argv: string[]) => runOmnigraph($, argv.slice(1), op.intent.cwd)
245 const list = await run(commitListArgv(config.bin, op.intent.target, branch))
246 const consumed = new Set((s.ledger?.lines ?? []).flatMap(l => (l.t === 'diff' ? [l.commitId] : [])))
247 const commitId = pickCommit(parseCommits(list.stdout), op.startedAt, endedAt, consumed)
248 if (commitId === undefined) return
249 const changes = await run(commitChangesArgv(config.bin, op.intent.target, commitId))
250 const parsed = parseChanges(changes.stdout)
251 if (parsed.changes.length === 0) return
252 apply($, s, { v: 1, t: 'diff', at: await $.clock.now(), opId: op.id, commitId, ...parsed })
253}
254
255/** After an op ends: bootstrap the graph once, light a read's nodes, fetch a write's diff. */
256function enrich($: EngineInterface, s: Session, op: OpEvent, end: LedgerLine): void {
257 if (!config.isEnriching || end.t !== 'op.end') return
258 const key = targetKey(op.intent.target, op.intent.cwd)
259 const o = end.outcome
260 const isWrite = end.status === 'ok' && (o.kind === 'change' || o.kind === 'load')
261 // Bootstrap once; after a failure, retry for another graph or once a write on this one succeeded.
262 if (graphState === 'none' || (graphState === 'failed' && (key !== graphKey || isWrite))) {
263 graphState = 'loading'
264 graphKey = key
265 enqueue($, end.at, () => bootstrap($, s, op.intent))
266 }
267 // Only the session's graph is drawn: another graph's ops stay in the Feed alone.
268 if (end.status !== 'ok' || key !== graphKey) return
269 if (o.kind === 'read' && o.cells.length > 0) {
270 enqueue($, end.at, async () => {
271 if (!s.viz.ghost) return
272 const { nodeIds, isExact } = matchCells(o.cells, buildMatchIndex(s.viz.ghost))
273 if (nodeIds.length > 0) apply($, s, { v: 1, t: 'touch', at: await $.clock.now(), opId: op.id, nodeIds, isExact })
274 })
275 }
276 const branch = o.kind === 'change' || o.kind === 'load' ? o.branch : o.kind === 'branch' && o.action === 'merged' ? o.into : undefined
277 if (branch) enqueue($, end.at, () => enrichWrite($, s, op, branch, end.at))
278}
279
280async function startCapture($: EngineInterface, command: string, toolUseId: string, agentId: string | undefined): Promise<OpEvent[]> {
281 const s = session
282 // Tokenize before any $ call: most Bash commands are not omnigraph.
283 if (!s || findInvocations(command, '/', config.bin).length === 0) return []
284 const now = await $.clock.now()
285 await ensureLedger($, s, now)
286 if (homeCache === null) homeCache = await $.env.get('HOME')
287 const home = homeCache
288 const ops = beginOps({
289 command,
290 toolUseId,
291 cwd: await $.session.cwd(),
292 bin: config.bin,
293 now,
294 ...(agentId !== undefined ? { agentId } : {}),
295 ...(home !== undefined ? { home } : {}),
296 })
297 for (const op of ops) apply($, s, { v: 1, t: 'op.start', at: op.startedAt, op })
298 if (ops.length > 0) await publish($)
299 return ops
300}
301
302async function finishCapture($: EngineInterface, ops: readonly OpEvent[], r: Parameters<typeof endLines>[1]): Promise<void> {
303 const s = session
304 if (!s || ops.length === 0) return
305 const ends = endLines(ops, r, await $.clock.now())
306 for (const line of ends) apply($, s, line)
307 ops.forEach((op, i) => enrich($, s, op, ends[i]!))
308 await publish($)
309}
310
311function stopAnimation(): void {
312 frameTimer?.cancel()
313 frameTimer = undefined
314}
315
316/** Runs frames while the Graph tab is on screen and something moves; idles otherwise. */
317function startAnimation($: EngineInterface): void {
318 if (frameTimer || !canvas) return
319 frameTimer = $.clock.every(FRAME_MS, () => void frame($))
320}
321
322async function frame($: EngineInterface): Promise<void> {
323 if (isFrameBusy) return
324 const c = canvas
325 if (!c) {
326 stopAnimation()
327 return
328 }
329 isFrameBusy = true
330 try {
331 const now = await $.clock.now()
332 // Replay lines due by now play first, so a frame never steps the layout past one.
333 const isNewLines = playDue(now)
334 const show = shown()
335 if (!show) {
336 stopAnimation()
337 return
338 }
339 show.scene.tick(now)
340 // The browser is a Client: its pane redraws it, and the pane's draw renders the frame.
341 $.ui.invalidate('ui.render')
342 if (isNewLines) await publish($)
343 if (!show.scene.isAnimating(now)) stopAnimation()
344 } catch (err) {
345 debug($, err)
346 stopAnimation()
347 } finally {
348 isFrameBusy = false
349 }
350}
351
352/** The selection link: the op's nodes light up again on the graph and keep their labels. */
353async function selectOp($: EngineInterface, id: string): Promise<void> {
354 const ui = await update($, uiAtom, u => ({ ...u, selectedOpId: u.selectedOpId === id ? null : id }))
355 const show = shown()
356 const ghost = show?.viz.ghost
357 if (!show || ui.selectedOpId === null || !ghost) {
358 if (show && ghost) {
359 show.scene.unfocus(ghost, await $.clock.now())
360 startAnimation($)
361 } else if (show) show.scene.pinned = []
362 return
363 }
364 const op = show.viz.ops.find(o => o.id === ui.selectedOpId)
365 show.scene.focus(show.viz.touches[ui.selectedOpId] ?? [], ghost, await $.clock.now(), op ? opName(op.intent) : undefined)
366 startAnimation($)
367}
368
369/** The Graph tab's browser: a canvas sized to the pane, drawn now. */
370async function graphView($: EngineInterface, show: Show | undefined, columns: number, bodyRows: number, ghost: GraphViewData['ghost']): Promise<GraphViewData> {
371 // A truncated graph says so in a row of its own under the browser.
372 const rows = browserRowsFor(bodyRows) - (ghost.truncated ? 1 : 0)
373 let runs: GraphViewData['runs'] = null
374 if (show?.viz.ghost) {
375 if (!canvas || canvas.cols !== columns || canvas.rows !== rows) canvas = new BrailleCanvas(columns, rows)
376 const now = await $.clock.now()
377 canvas.clear()
378 show.scene.drawFrame(canvas, show.viz.ghost, now)
379 runs = canvas.runs()
380 // Only while something moves: each frame redraws the pane, and this draw must not restart a loop that just stopped.
381 if (show.scene.isAnimating(now)) startAnimation($)
382 } else {
383 canvas = null
384 }
385 return { ghost, columns, rows, runs, isBrowsing: show ? !show.scene.isFollowing : false, recent: show?.viz.ops.slice(-3) ?? [] }
386}
387
388/** Plays every replay line now due; true when any did. */
389function playDue(now: number): boolean {
390 const r = replay
391 if (!r) return false
392 const lines = r.player.due(now)
393 for (const line of lines) feed(r, line)
394 return lines.length > 0
395}
396
397/** The replay's clock: plays what is due, publishes, and sleeps until the next line. A failure ends the replay. */
398async function pumpReplay($: EngineInterface, r: Replay): Promise<void> {
399 r.timer = undefined
400 if (replay !== r) return
401 try {
402 const now = await $.clock.now()
403 if (playDue(now)) {
404 startAnimation($)
405 await publish($)
406 }
407 const next = r.player.nextAt
408 if (replay === r && next !== undefined) r.timer = $.clock.after(Math.max(0, next - now), () => void pumpReplay($, r))
409 } catch (err) {
410 debug($, err)
411 if (replay === r) {
412 stopReplay()
413 await publish($).catch(() => undefined)
414 }
415 }
416}
417
418function stopReplay(): void {
419 replay?.timer?.cancel()
420 replay = undefined
421}
422
423/** `/og replay <path> [speed]`: reads and checks the ledger, then plays it in the pane's Graph tab. */
424async function startReplay($: EngineInterface, path: string, speed: number): Promise<string> {
425 if (homeCache === null) homeCache = await $.env.get('HOME')
426 const file = resolvePath(session?.cwd ?? (await $.session.cwd()), path, homeCache)
427 let text: string
428 try {
429 text = await $.fs.read(file)
430 } catch (err) {
431 return `Can't read ${path}: ${err instanceof Error ? err.message : String(err)}`
432 }
433 const loaded = readReplay(text)
434 if ('error' in loaded) return `Can't replay ${path}: ${loaded.error}`
435 // Everything that can fail before the switch happens first, so a running replay is never left half-stopped.
436 const r: Replay = {
437 name: basename(file),
438 player: new Player(loaded.lines, speed, Math.round(await $.clock.now())),
439 hasGraph: loaded.lines.some(line => line.t === 'graph'),
440 viz: emptyViz(),
441 scene: new Scene(),
442 timer: undefined,
443 }
444 stopReplay()
445 if (session) session.scene.pinned = []
446 replay = r
447 try {
448 await update($, uiAtom, (u): UiState => ({ ...u, tab: 'graph', selectedOpId: null }))
449 await $.ui.open({ id: PANE, title: 'omnigraph', focus: true })
450 await pumpReplay($, r)
451 await publish($)
452 } catch (err) {
453 debug($, err)
454 if (replay === r) stopReplay()
455 await publish($).catch(() => undefined)
456 return `Couldn't replay ${path}: ${err instanceof Error ? err.message : String(err)}`
457 }
458 return `Replaying ${r.name} at ${speed}×, about ${Math.round(r.player.durationMs / 1000)} s. /og stop ends it.`
459}
460
461export const register: Register = (on, options) => {
462 config = {
463 bin: typeof options.omnigraphBin === 'string' ? options.omnigraphBin : 'omnigraph',
464 ledgerDir: typeof options.ledgerDir === 'string' ? options.ledgerDir : '.omnigraph-viz',
465 maxGhostNodes: typeof options.maxGhostNodes === 'number' ? options.maxGhostNodes : 2000,
466 isEnriching: options.enrich !== false,
467 isDebug: options.debug === true,
468 }
469
470 on('session.start', async ($, e, next) => {
471 try {
472 await $.command.register({
473 name: 'og',
474 description: 'Omnigraph viz: toggle the pane, replay a recorded session, or reset the view',
475 argumentHint: OG_HINT,
476 })
477 stopReplay()
478 graphState = 'none'
479 graphMessage = null
480 graphKey = null
481 session = await openSession($, e.cwd)
482 await publish($)
483 } catch (err) {
484 debug($, err)
485 }
486 return next(e)
487 })
488
489 on('session.end', async ($, e, next) => {
490 // `/clear` raises session.end but no session.start follows: the session goes on.
491 if (e.reason !== 'clear') {
492 if (session) session.isEnded = true
493 stopReplay()
494 }
495 stopAnimation()
496 $.ui.status(undefined)
497 return next(e)
498 })
499
500 on('command.run', { command: 'og' }, async ($, e) => {
501 const cmd = parseOgArgs(e.args)
502 if (cmd.kind === 'usage') return { text: cmd.message }
503 if (cmd.kind === 'replay') return { text: await startReplay($, cmd.path, cmd.speed) }
504 if (cmd.kind === 'stop') {
505 if (!replay) return { text: 'No replay is running.' }
506 stopReplay()
507 if (session) session.scene.pinned = []
508 await update($, uiAtom, u => ({ ...u, selectedOpId: null }))
509 await publish($)
510 return { text: 'Replay stopped; back to the live session.' }
511 }
512 if (cmd.kind === 'clear') {
513 stopReplay()
514 if (session) {
515 // The Feed starts over; the graph stays, so the Graph tab keeps working.
516 session.viz = { ...emptyViz(), ghost: session.viz.ghost }
517 }
518 await publish($)
519 return { text: 'Omnigraph viz cleared.' }
520 }
521 const isShown = (await $.ui.panes()).some(p => p.id === PANE && p.isShown)
522 if (isShown) {
523 await $.ui.close({ id: PANE })
524 return { text: 'Omnigraph viz closed.' }
525 }
526 // Focused, so g / f / Tab / Enter work at once; Esc hands the keys back to the prompt.
527 await $.ui.open({ id: PANE, title: 'omnigraph', focus: true })
528 return { text: 'Omnigraph viz opened.' }
529 })
530
531 // Observer only: the command and the result pass through untouched. Every
532 // Bash call comes through (a custom omnigraphBin need not say "omnigraph");
533 // startCapture tokenizes before touching $.
534 on('tool.call', { tool: 'Bash' }, async ($, e, next) => {
535 let ops: OpEvent[] = []
536 try {
537 ops = await startCapture($, e.command, e.tool_use_id, e.agentId)
538 } catch (err) {
539 debug($, err)
540 }
541 const r = await next(e)
542 try {
543 await finishCapture($, ops, r)
544 } catch (err) {
545 debug($, err)
546 }
547 return r
548 }).catch(($, e, next) => next(e))
549
550 // The browser Client's keys: select, open, go back, follow live, switch tabs. The answer is its next rows.
551 on('ui.message', { requestId: 'omnigraph' }, async ($, e) => {
552 const show = shown()
553 const ghost = show?.viz.ghost
554 const data = e.data as { key?: unknown; meta?: unknown; shift?: unknown; ctrl?: unknown } | null
555 if (!show || !ghost || !canvas || e.element !== 'explore' || typeof data?.key !== 'string') return {}
556 const action = actionFor({ key: data.key, meta: data.meta === true, shift: data.shift === true, ctrl: data.ctrl === true })
557 if (action.kind === 'tab') {
558 await update($, uiAtom, (u): UiState => ({ ...u, tab: action.tab }))
559 return {}
560 }
561 const wasFollowing = show.scene.isFollowing
562 const now = await $.clock.now()
563 show.scene.key(action, ghost, now)
564 // The tab row's hint changes when browsing starts or stops: redraw the pane for it.
565 if (show.scene.isFollowing !== wasFollowing) $.ui.invalidate('ui.render')
566 if (show.scene.isAnimating(now)) startAnimation($)
567 canvas.clear()
568 show.scene.drawFrame(canvas, ghost, now)
569 return { props: { rows: canvas.runs() } }
570 })
571
572 on('ui.render', { component: 'Pane', requestId: 'omnigraph' }, async ($, e, next) => {
573 if (e.surface !== 'terminal') return next(e)
574 const [view, ui] = await Promise.all([read($, viewAtom), read($, uiAtom)])
575 const show = shown()
576 const selectedOutput = show?.viz.ops.find(op => op.id === ui.selectedOpId)?.output
577 const columns = e.props.bodyColumns
578 const rows = e.props.scroll.bodyRows
579 let graph: GraphViewData | undefined
580 if (ui.tab !== 'feed') {
581 graph = await graphView($, show, columns, rows, view.ghost)
582 } else {
583 canvas = null
584 }
585 return drawPane(
586 $.ui.resolve(e),
587 {
588 view,
589 ui,
590 columns,
591 rows,
592 ...(selectedOutput !== undefined ? { selectedOutput } : {}),
593 ...(graph !== undefined ? { graph } : {}),
594 },
595 {
596 setTab: tab => void update($, uiAtom, u => ({ ...u, tab })),
597 select: id => void selectOp($, id),
598 },
599 )
600 })
601}
602src/capture/capture.ts 130 lines1import type { LedgerLine, OpEvent, OpOutcome, Verb } from '../../types'
2import { classify } from './classify'
3import { parseOutput, segmentOutput, stripAnsi } from './parse-output'
4import { findInvocations } from './tokenize'
5
6/** The parts of a `tool.call` result the capture reads. */
7export interface ToolResultLike {
8 deny?: string
9 isError?: boolean
10 result?: unknown
11 text?: string
12}
13
14export function bashOutput(r: ToolResultLike): { stdout: string; stderr: string; isError: boolean } {
15 if (r.deny !== undefined) return { stdout: '', stderr: r.deny, isError: true }
16 const result = r.result
17 if (typeof result === 'object' && result !== null && 'stdout' in result) {
18 const { stdout, stderr } = result as { stdout?: unknown; stderr?: unknown }
19 return {
20 stdout: typeof stdout === 'string' ? stdout : '',
21 stderr: typeof stderr === 'string' ? stderr : '',
22 isError: r.isError === true,
23 }
24 }
25 return { stdout: r.text ?? '', stderr: '', isError: r.isError === true }
26}
27
28export const OUTPUT_LINES = 40
29export const OUTPUT_CHARS = 4000
30
31/** The first 40 lines / 4000 chars: enough for the Feed's detail, small enough for the ledger. */
32export function clipOutput(text: string): string {
33 const lines = text.split('\n')
34 const head = lines.slice(0, OUTPUT_LINES).join('\n')
35 const clipped = head.length > OUTPUT_CHARS ? head.slice(0, OUTPUT_CHARS) : head
36 return clipped.length < text.length ? `${clipped.trimEnd()}\n…` : clipped
37}
38
39export interface BeginArgs {
40 command: string
41 toolUseId: string
42 cwd: string
43 bin: string
44 now: number
45 agentId?: string
46 /** HOME, to expand `~` the way the shell did. */
47 home?: string
48}
49
50export function beginOps(args: BeginArgs): OpEvent[] {
51 return findInvocations(args.command, args.cwd, args.bin, args.home).map(inv => {
52 const op: OpEvent = {
53 id: `${args.toolUseId}#${inv.index}`,
54 startedAt: args.now,
55 intent: classify(inv),
56 status: 'running',
57 }
58 if (args.agentId !== undefined) op.agentId = args.agentId
59 return op
60 })
61}
62
63/** Which verbs can have produced each outcome kind. */
64const PRODUCED_BY: Readonly<Record<OpOutcome['kind'], readonly Verb[] | 'any'>> = {
65 read: ['query', 'alias'],
66 change: ['mutate'],
67 branch: ['branch', 'mutate'],
68 load: ['load'],
69 error: 'any',
70 other: 'any',
71}
72
73const isResult = (o: OpOutcome | undefined) => o !== undefined && o.kind !== 'error' && o.kind !== 'other'
74
75interface Slot {
76 outcome: OpOutcome
77 output: string
78}
79
80/**
81 * One outcome per op. A single call parses the whole output. Several calls
82 * get one segment each when the counts line up; otherwise each segment goes
83 * to the next op whose verb can have produced it, so one result never counts
84 * twice. A failed exit with no error found blames the last op with no result.
85 */
86function slotsFor(ops: readonly OpEvent[], stdout: string, stderr: string, isError: boolean): Slot[] {
87 if (ops.length === 1) {
88 return [{ outcome: parseOutput(stdout, stderr, isError), output: isError && stderr !== '' ? stderr : stdout }]
89 }
90 const segments = segmentOutput(stdout)
91 const slots: (Slot | undefined)[] = ops.map(() => undefined)
92 if (segments.length === ops.length) {
93 segments.forEach((segment, i) => {
94 slots[i] = { outcome: parseOutput(segment, '', false), output: segment }
95 })
96 } else {
97 let from = 0
98 for (const segment of segments) {
99 const outcome = parseOutput(segment, '', false)
100 const verbs = PRODUCED_BY[outcome.kind]
101 const at = ops.findIndex((op, i) => i >= from && (verbs === 'any' || verbs.includes(op.intent.verb)))
102 if (at < 0) continue
103 slots[at] = { outcome, output: segment }
104 from = at + 1
105 }
106 }
107 if (isError && !slots.some(slot => slot?.outcome.kind === 'error')) {
108 let at = -1
109 slots.forEach((slot, i) => {
110 if (!isResult(slot?.outcome)) at = i
111 })
112 if (at >= 0) slots[at] = { outcome: parseOutput(stdout, stderr, true), output: stderr !== '' ? stderr : stdout }
113 }
114 return slots.map(slot => slot ?? { outcome: { kind: 'other', firstLine: '' }, output: '' })
115}
116
117export function endLines(ops: readonly OpEvent[], r: ToolResultLike, now: number): LedgerLine[] {
118 const { stdout, stderr, isError } = bashOutput(r)
119 return slotsFor(ops, stdout, stderr, isError).map(({ outcome, output }, i) => ({
120 v: 1,
121 t: 'op.end',
122 at: now,
123 id: ops[i]!.id,
124 status: outcome.kind === 'error' ? 'error' : 'ok',
125 durationMs: Math.max(0, now - ops[i]!.startedAt),
126 outcome,
127 output: clipOutput(stripAnsi(output)),
128 }))
129}
130src/capture/parse-output.ts 286 lines1import type { MergeKind, OpOutcome } from '../../types'
2
3// Human formats from crates/omnigraph-cli/src/output.rs and read_format.rs.
4const READ_HEADER = /^(\d+) rows(?: from (branch|snapshot) (\S+))? via (\S+)$/m
5const CHANGE = /^changed (\S+) via (\S+): (\d+) nodes, (\d+) edges$/m
6const BRANCH_CREATED = /^created branch (\S+) from (\S+)$/m
7const BRANCH_DELETED = /^deleted branch (\S+)$/m
8const BRANCH_MERGED = /^merged (\S+) into (\S+): (already_up_to_date|fast_forward|merged)$/m
9const LOADED =
10 /^loaded (\S+) on branch (\S+) with (\w+): (\d+) entities across (\d+) node types and (\d+) edge types$/m
11const LOAD_BRANCH = /^branch (\S+) created from (\S+)$/m
12/** `snapshot`'s human output starts with this line. */
13const SNAPSHOT_HEAD = /^graph_branch: \S+$/
14
15/** The target echo a write prints on stderr; Claude's Bash merges it into stdout. */
16const TARGET_ECHO = /^omnigraph \S+ → .+ \((?:direct|served)[^)]*\)$/
17/**
18 * clap's `error: …` and color-eyre's `Error:` (message on the next line, as ` 0: …`).
19 * Both print at column 0; an indented `error:` is data (a schema property, a cell).
20 */
21const ERROR_START = /^error:/i
22/**
23 * Terminal escapes a shell can wrap output in: CSI (colors, `ESC[?25h`), OSC
24 * (titles, links), charset picks such as `ESC(B` (what `tput sgr0` prints;
25 * seen before every result when Claude Code's sandbox is off) and the other
26 * two-byte escapes.
27 */
28const ANSI = /\x1b(?:\[[0-?]*[ -/]*[@-~]|\][^\x07\x1b]*(?:\x07|\x1b\\)|[ -/]+[0-~]|[0-~])/g
29
30/** The first line of any result the CLI prints; used to split multi-invocation output. */
31const RESULT_LINE = new RegExp(
32 [READ_HEADER, CHANGE, BRANCH_CREATED, BRANCH_DELETED, BRANCH_MERGED, LOADED, SNAPSHOT_HEAD]
33 .map(re => re.source)
34 .join('|'),
35)
36
37export const MAX_ROWS = 200
38
39type Json = Record<string, unknown>
40
41const isRecord = (v: unknown): v is Json => typeof v === 'object' && v !== null && !Array.isArray(v)
42const str = (v: unknown): string | undefined => (typeof v === 'string' ? v : undefined)
43const num = (v: unknown): number | undefined => (typeof v === 'number' ? v : undefined)
44
45function lastLine(text: string): string | undefined {
46 const lines = text.split('\n').map(l => l.trim()).filter(Boolean)
47 return lines[lines.length - 1]
48}
49
50function firstLine(text: string): string {
51 return (text.split('\n').find(l => l.trim() !== '') ?? '').trim().slice(0, 200)
52}
53
54/** As read_format.rs `stringify_value` prints a cell. */
55export function cellText(v: unknown): string {
56 if (v === null || v === undefined) return 'null'
57 if (typeof v === 'string') return v
58 if (typeof v === 'number' || typeof v === 'boolean') return String(v)
59 return JSON.stringify(v)
60}
61
62/** Control characters left after the escapes (tmux's sgr0 ends with SI); tabs and line breaks stay. */
63const CONTROL = /[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]/g
64
65export function stripAnsi(text: string): string {
66 return text.replace(ANSI, '').replace(CONTROL, '')
67}
68
69/** No colors, no CRLF, no target echo lines. */
70function clean(text: string): string {
71 return stripAnsi(text)
72 .replace(/\r\n/g, '\n')
73 .split('\n')
74 .filter(line => !TARGET_ECHO.test(line))
75 .join('\n')
76}
77
78/** The message of the first error block in `text`, if there is one. */
79function errorIn(text: string): string | undefined {
80 const lines = text.split('\n')
81 const at = lines.findIndex(line => ERROR_START.test(line))
82 if (at < 0) return undefined
83 const inline = lines[at]!.replace(/^error:\s*/i, '').trim()
84 if (inline !== '') return inline.slice(0, 200)
85 const next = lines.slice(at + 1).find(line => line.trim() !== '')
86 return next?.trim().replace(/^\d+:\s*/, '').slice(0, 200)
87}
88
89function errorMessage(stdout: string, stderr: string): string {
90 return (errorIn(stderr) ?? errorIn(stdout) ?? lastLine(stderr) ?? lastLine(stdout) ?? 'command failed').slice(0, 200)
91}
92
93function parseTable(text: string, m: RegExpMatchArray): OpOutcome {
94 const where =
95 m[2] === 'branch' ? { branch: m[3]! } : m[2] === 'snapshot' ? { snapshot: m[3]! } : {}
96 const base = { kind: 'read' as const, queryName: m[4]!, rowCount: Number(m[1]), ...where }
97 const lines = text.slice((m.index ?? 0) + m[0].length).split('\n').slice(1)
98 const head = lines[0]
99 if (head === undefined || head.trim() === '' || head.trim() === '(no rows)') {
100 return { ...base, columns: [], cells: [] }
101 }
102 const columns = head.split(' | ').map(c => c.trim())
103 const firstRow = /^[-+]+$/.test(lines[1]?.trim() ?? '') ? 2 : 1
104 const cells: string[][] = []
105 for (const line of lines.slice(firstRow)) {
106 if (line.trim() === '') break
107 const parts = line.split(' | ')
108 if (parts.length !== columns.length) break
109 const row = parts.map(c => c.trimEnd())
110 const prev = cells[cells.length - 1]
111 // Wrap layout: a continuation line repeats the row with blanks where nothing wrapped.
112 if (prev && row[0] === '' && row.some(c => c !== '')) {
113 row.forEach((c, i) => {
114 prev[i] = `${prev[i] ?? ''}${c}`
115 })
116 continue
117 }
118 if (cells.length >= MAX_ROWS) break
119 cells.push(row)
120 }
121 return { ...base, columns, cells }
122}
123
124function parseJsonish(text: string): unknown[] | undefined {
125 try {
126 return [JSON.parse(text)]
127 } catch {
128 // maybe JSONL: metadata first, then one row per line
129 }
130 const out: unknown[] = []
131 for (const line of text.split('\n')) {
132 if (line.trim() === '') continue
133 try {
134 out.push(JSON.parse(line))
135 } catch {
136 return undefined
137 }
138 }
139 return out.length > 0 ? out : undefined
140}
141
142function parseJson(text: string): OpOutcome | undefined {
143 const values = parseJsonish(text)
144 const head = values?.[0]
145 if (!values || !isRecord(head)) return undefined
146
147 if ('row_count' in head || 'rows' in head) {
148 const rows = Array.isArray(head.rows) ? head.rows : values.slice(1)
149 const first = rows[0]
150 const columns = Array.isArray(head.columns)
151 ? head.columns.filter((c): c is string => typeof c === 'string')
152 : isRecord(first)
153 ? Object.keys(first)
154 : []
155 const target = isRecord(head.target) ? head.target : {}
156 const branch = str(target.branch)
157 const snapshot = str(target.snapshot)
158 return {
159 kind: 'read',
160 queryName: str(head.query_name) ?? '',
161 rowCount: num(head.row_count) ?? rows.length,
162 ...(branch !== undefined ? { branch } : {}),
163 ...(snapshot !== undefined ? { snapshot } : {}),
164 columns,
165 cells: rows.slice(0, MAX_ROWS).map(r => columns.map(c => cellText(isRecord(r) ? r[c] : undefined))),
166 }
167 }
168
169 const outcome = head.outcome
170 if (isRecord(outcome) && typeof outcome.kind === 'string') {
171 if (outcome.kind === 'created') {
172 return { kind: 'branch', action: 'created', name: str(outcome.name) ?? '', from: str(outcome.from) ?? '' }
173 }
174 if (outcome.kind === 'deleted') return { kind: 'branch', action: 'deleted', name: str(outcome.name) ?? '' }
175 if (outcome.kind === 'merged') {
176 return {
177 kind: 'branch', action: 'merged', name: str(outcome.source) ?? '',
178 into: str(outcome.target) ?? '', merge: (str(outcome.merge) ?? 'merged') as MergeKind,
179 }
180 }
181 }
182 if ('affected_nodes' in head) {
183 return {
184 kind: 'change', branch: str(head.branch) ?? '', queryName: str(head.query_name) ?? '',
185 nodes: num(head.affected_nodes) ?? 0, edges: num(head.affected_edges) ?? 0,
186 }
187 }
188 if (typeof head.source === 'string' && typeof head.target === 'string' && typeof outcome === 'string') {
189 return { kind: 'branch', action: 'merged', name: head.source, into: head.target, merge: outcome as MergeKind }
190 }
191 if ('total_entities' in head) {
192 const base = str(head.base_branch)
193 return {
194 kind: 'load', branch: str(head.branch) ?? '', mode: str(head.mode) ?? '',
195 total: num(head.total_entities) ?? 0,
196 nodeTypes: Array.isArray(head.nodes) ? head.nodes.length : 0,
197 edgeTypes: Array.isArray(head.edges) ? head.edges.length : 0,
198 ...(head.branch_created === true && base !== undefined ? { branchCreatedFrom: base } : {}),
199 }
200 }
201 return { kind: 'other', firstLine: firstLine(text) }
202}
203
204/** A read envelope cut off mid-way (long output gets truncated): keep the header facts. */
205function parseCutJson(text: string): OpOutcome | undefined {
206 const queryName = text.match(/"query_name"\s*:\s*"([^"]*)"/)?.[1]
207 const rowCount = text.match(/"row_count"\s*:\s*(\d+)/)?.[1]
208 if (queryName === undefined || rowCount === undefined) return undefined
209 const branch = text.match(/"branch"\s*:\s*"([^"]*)"/)?.[1]
210 return {
211 kind: 'read', queryName, rowCount: Number(rowCount),
212 ...(branch !== undefined ? { branch } : {}),
213 columns: [], cells: [],
214 }
215}
216
217export function parseOutput(stdout: string, stderr: string, isError: boolean): OpOutcome {
218 const text = clean(stdout)
219 const errors = clean(stderr)
220 if (isError) return { kind: 'error', message: errorMessage(text, errors) }
221 const trimmed = text.trim()
222 if (trimmed.startsWith('{')) {
223 const fromJson = parseJson(trimmed) ?? parseCutJson(trimmed)
224 if (fromJson) return fromJson
225 }
226 let m = text.match(READ_HEADER)
227 if (m) return parseTable(text, m)
228 if ((m = text.match(CHANGE))) {
229 return { kind: 'change', branch: m[1]!, queryName: m[2]!, nodes: Number(m[3]), edges: Number(m[4]) }
230 }
231 if ((m = text.match(BRANCH_CREATED))) return { kind: 'branch', action: 'created', name: m[1]!, from: m[2]! }
232 if ((m = text.match(BRANCH_DELETED))) return { kind: 'branch', action: 'deleted', name: m[1]! }
233 if ((m = text.match(BRANCH_MERGED))) {
234 return { kind: 'branch', action: 'merged', name: m[1]!, into: m[2]!, merge: m[3] as MergeKind }
235 }
236 if ((m = text.match(LOADED))) {
237 const created = text.match(LOAD_BRANCH)
238 return {
239 kind: 'load', branch: m[2]!, mode: m[3]!, total: Number(m[4]),
240 nodeTypes: Number(m[5]), edgeTypes: Number(m[6]),
241 ...(created ? { branchCreatedFrom: created[2]! } : {}),
242 }
243 }
244 // No result line, but an error block: a failure the shell's exit status hid (`a ; b`).
245 const failure = errorIn(text) ?? errorIn(errors)
246 if (failure !== undefined) return { kind: 'error', message: failure }
247 return { kind: 'other', firstLine: firstLine(trimmed) }
248}
249
250type Boundary = 'echo' | 'result' | 'error'
251
252function boundaryOf(line: string): Boundary | undefined {
253 if (TARGET_ECHO.test(line)) return 'echo'
254 if (RESULT_LINE.test(line)) return 'result'
255 if (ERROR_START.test(line)) return 'error'
256 return undefined
257}
258
259/**
260 * Cuts one command's stdout into the outputs of the omnigraph calls in it.
261 * A call's output starts at a write's target echo, a result line or an error
262 * block; a result or error right after an echo belongs to that echo. Lines
263 * before the first boundary join the first segment; no boundary, one segment.
264 */
265export function segmentOutput(stdout: string): string[] {
266 const lines = stripAnsi(stdout).replace(/\r\n/g, '\n').split('\n')
267 const starts: number[] = []
268 let isEchoOpen = false
269 lines.forEach((line, i) => {
270 const boundary = boundaryOf(line)
271 if (boundary === undefined) return
272 if (boundary !== 'echo' && isEchoOpen) {
273 isEchoOpen = false
274 return
275 }
276 starts.push(i)
277 isEchoOpen = boundary === 'echo'
278 })
279 if (starts.length === 0) return [stdout]
280 starts[0] = 0
281 return starts.map((start, k) => {
282 const end = starts[k + 1] ?? lines.length
283 return lines.slice(start, end).join('\n') + (end < lines.length ? '\n' : '')
284 })
285}
286src/capture/tokenize.ts 199 lines1import type { Invocation } from '../../types'
2import { basename, expandHome, resolvePath } from '../util/path'
3
4type Token = { kind: 'word'; value: string } | { kind: 'op'; value: string }
5
6const WRAPPERS = new Set(['env', 'time', 'command', 'exec', 'nohup'])
7/** Words that open a compound command's body: `do omnigraph …`, `then omnigraph …`. */
8const KEYWORDS = new Set(['do', 'then', 'else', 'elif', 'if', 'while', 'until', '!', '{'])
9/** `timeout` options that take a value. */
10const TIMEOUT_VALUE_FLAGS = new Set(['-s', '-k', '--signal', '--kill-after'])
11const ASSIGNMENT = /^[A-Za-z_][A-Za-z0-9_]*=/
12
13/**
14 * A small POSIX-ish lexer: quotes, escapes, comments, `&&` `||` `;` `|` `&`,
15 * newlines and parentheses as separators, and redirections (`>`, `>>`, `<`,
16 * `2>`, `2>&1`, `&>`), whose targets are dropped.
17 */
18export function lex(src: string): Token[] {
19 const tokens: Token[] = []
20 let word = ''
21 let isInWord = false
22 const flush = () => {
23 if (isInWord) tokens.push({ kind: 'word', value: word })
24 word = ''
25 isInWord = false
26 }
27 let i = 0
28 while (i < src.length) {
29 const c = src[i]!
30 if (c === "'") {
31 const end = src.indexOf("'", i + 1)
32 const stop = end < 0 ? src.length : end
33 word += src.slice(i + 1, stop)
34 isInWord = true
35 i = stop + 1
36 continue
37 }
38 if (c === '"') {
39 isInWord = true
40 i++
41 while (i < src.length && src[i] !== '"') {
42 const next = src[i + 1]
43 if (src[i] === '\\' && next !== undefined && '"\\$`'.includes(next)) {
44 word += next
45 i += 2
46 continue
47 }
48 word += src[i]
49 i++
50 }
51 i++
52 continue
53 }
54 if (c === '\\') {
55 const next = src[i + 1]
56 if (next !== undefined && next !== '\n') {
57 word += next
58 isInWord = true
59 }
60 i += 2
61 continue
62 }
63 if (c === ' ' || c === '\t') {
64 flush()
65 i++
66 continue
67 }
68 if (c === '#' && !isInWord) {
69 while (i < src.length && src[i] !== '\n') i++
70 continue
71 }
72 if (c === '\n' || c === ';' || c === '(' || c === ')') {
73 flush()
74 tokens.push({ kind: 'op', value: ';' })
75 i++
76 continue
77 }
78 if (c === '&' || c === '|') {
79 flush()
80 const two = src.slice(i, i + 2)
81 if (two === '&&' || two === '||') {
82 tokens.push({ kind: 'op', value: two })
83 i += 2
84 continue
85 }
86 if (two === '&>') {
87 tokens.push({ kind: 'op', value: '>' })
88 i += 2
89 continue
90 }
91 tokens.push({ kind: 'op', value: c })
92 i++
93 continue
94 }
95 if (c === '>' || c === '<') {
96 if (isInWord && /^\d+$/.test(word)) {
97 word = ''
98 isInWord = false
99 } else {
100 flush()
101 }
102 let j = i + 1
103 if (src[j] === c) j++
104 if (src[j] === '&') {
105 j++
106 while (j < src.length && /[\d-]/.test(src[j]!)) j++
107 i = j
108 continue
109 }
110 tokens.push({ kind: 'op', value: '>' })
111 i = j
112 continue
113 }
114 word += c
115 isInWord = true
116 i++
117 }
118 flush()
119 return tokens
120}
121
122/**
123 * Strips compound-command keywords (`do`, `then`, `{`, …), `VAR=x`, and the
124 * `env`, `time`, `command`, `exec`, `nohup`, `timeout <duration>` wrappers;
125 * undefined for `command -v x`.
126 */
127function stripPrefixes(words: string[]): string[] | undefined {
128 let k = 0
129 for (;;) {
130 const word = words[k]
131 if (word === undefined) break
132 if (KEYWORDS.has(word) || ASSIGNMENT.test(word)) {
133 k++
134 continue
135 }
136 if (word === 'timeout') {
137 k++
138 while (words[k]?.startsWith('-')) k += TIMEOUT_VALUE_FLAGS.has(words[k]!) ? 2 : 1
139 k++ // the duration
140 continue
141 }
142 if (WRAPPERS.has(word)) {
143 k++
144 if (word === 'command' && (words[k] === '-v' || words[k] === '-V')) return undefined
145 while (k < words.length && (ASSIGNMENT.test(words[k]!) || words[k]!.startsWith('-'))) k++
146 continue
147 }
148 break
149 }
150 return words.slice(k)
151}
152
153/**
154 * Every `omnigraph` invocation in a Bash command line, with the directory it
155 * runs in (`cd` segments move it). `bin` is the configured binary; matching is
156 * by basename, so `/usr/local/bin/omnigraph` matches `omnigraph`. With `home`,
157 * `~` in `cd` targets and leading `~/` in arguments expand as the shell would.
158 */
159export function findInvocations(command: string, cwd: string, bin: string, home?: string): Invocation[] {
160 const want = basename(bin)
161 const found: Invocation[] = []
162 let dir = cwd
163 let words: string[] = []
164 let isSkippingTarget = false
165
166 const endSegment = () => {
167 const segment = stripPrefixes(words)
168 words = []
169 if (!segment || segment.length === 0) return
170 if (segment[0] === 'cd') {
171 const to = segment[1]
172 if (to !== undefined && to !== '-') dir = resolvePath(dir, to, home)
173 return
174 }
175 if (basename(segment[0]!) === want) {
176 // The shell expands a leading `~`: record what omnigraph actually received.
177 found.push({ argv: segment.map(word => expandHome(word, home)), cwd: dir, index: found.length })
178 }
179 }
180
181 for (const token of lex(command)) {
182 if (token.kind === 'word') {
183 if (isSkippingTarget) {
184 isSkippingTarget = false
185 continue
186 }
187 words.push(token.value)
188 continue
189 }
190 if (token.value === '>') {
191 isSkippingTarget = true
192 continue
193 }
194 endSegment()
195 }
196 endSegment()
197 return found
198}
199src/enrich/cli.ts 135 lines1import type { EntityChange, Target } from '../../types'
2import { resolvePath } from '../util/path'
3
4/** Global flags that choose the graph, passed through as the agent gave them. */
5export function targetFlags(t: Target): string[] {
6 const flags: string[] = []
7 for (const name of ['server', 'graph', 'store', 'profile', 'cluster'] as const) {
8 const value = t[name]
9 if (value !== undefined) flags.push(`--${name}`, value)
10 }
11 return flags
12}
13
14const uri = (t: Target): string[] => (t.uri !== undefined ? [t.uri] : [])
15
16export const schemaArgv = (bin: string, t: Target): string[] => [bin, 'schema', 'show', ...uri(t), '--json', ...targetFlags(t)]
17
18export const exportArgv = (bin: string, t: Target): string[] => [bin, 'export', ...uri(t), ...targetFlags(t)]
19
20export const commitListArgv = (bin: string, t: Target, branch: string): string[] => [
21 bin, 'commit', 'list', ...uri(t), '--branch', branch, '--json', ...targetFlags(t),
22]
23
24export const commitChangesArgv = (bin: string, t: Target, commitId: string): string[] => [
25 bin, 'commit', 'changes', commitId, ...(t.uri !== undefined ? ['--uri', t.uri] : []), '--json', ...targetFlags(t),
26]
27
28const isUrl = (p: string) => /^[a-z][a-z0-9+.-]*:\/\//i.test(p)
29const absolute = (cwd: string, p: string | undefined) => (p === undefined ? '' : isUrl(p) ? p : resolvePath(cwd, p))
30
31/**
32 * Which graph an op ran against. Paths resolve against the op's directory,
33 * so `cd dev && … --store graphs/x` and `… --store dev/graphs/x` are one
34 * graph; with no path at all the directory itself (its config) decides.
35 */
36export function targetKey(t: Target, cwd: string): string {
37 const hasPath = t.store !== undefined || t.uri !== undefined || t.cluster !== undefined
38 return JSON.stringify([
39 hasPath ? '' : cwd,
40 t.server ?? '', t.graph ?? '', t.profile ?? '',
41 absolute(cwd, t.store), absolute(cwd, t.cluster), absolute(cwd, t.uri),
42 ])
43}
44
45type Json = Record<string, unknown>
46const isRecord = (v: unknown): v is Json => typeof v === 'object' && v !== null && !Array.isArray(v)
47
48export function schemaSourceOf(stdout: string): string | undefined {
49 try {
50 const parsed: unknown = JSON.parse(stdout)
51 return isRecord(parsed) && typeof parsed.schema_source === 'string' ? parsed.schema_source : undefined
52 } catch {
53 return undefined
54 }
55}
56
57export interface CommitRef {
58 id: string
59 /** µs since the epoch. */
60 createdAt: number
61}
62
63/** `commit list --json`: newest first. */
64export function parseCommits(stdout: string): CommitRef[] {
65 try {
66 const parsed: unknown = JSON.parse(stdout)
67 const list = isRecord(parsed) && Array.isArray(parsed.commits) ? parsed.commits : Array.isArray(parsed) ? parsed : []
68 return list.flatMap(c =>
69 isRecord(c) && typeof c.graph_commit_id === 'string' && typeof c.created_at === 'number'
70 ? [{ id: c.graph_commit_id, createdAt: c.created_at }]
71 : [],
72 )
73 } catch {
74 return []
75 }
76}
77
78/** Clock slack between the agent's shell and the store, in ms. */
79export const COMMIT_SLACK_MS = 5000
80
81/**
82 * The op's commit: the oldest one made while it ran (with slack) that no
83 * earlier diff took. Chained writes in one Bash call share a window, so each
84 * takes the next commit in order instead of all taking the newest.
85 */
86export function pickCommit(
87 commits: readonly CommitRef[],
88 opStartedAtMs: number,
89 opEndedAtMs: number,
90 consumed: ReadonlySet<string> = new Set(),
91): string | undefined {
92 const inWindow = commits.filter(c => {
93 const ms = c.createdAt / 1000
94 return ms >= opStartedAtMs - COMMIT_SLACK_MS && ms <= opEndedAtMs + COMMIT_SLACK_MS && !consumed.has(c.id)
95 })
96 return inWindow[inWindow.length - 1]?.id // newest first: the oldest is last
97}
98
99export const MAX_CHANGES = 500
100
101const labelIn = (image: unknown): string | undefined => {
102 if (!isRecord(image) || !isRecord(image.properties)) return undefined
103 for (const field of ['name', 'title', 'slug']) {
104 const value = image.properties[field]
105 if (typeof value === 'string' && value !== '') return value
106 }
107 return undefined
108}
109
110/** `commit changes --json` → entity changes (at most 500) and how many more there were. */
111export function parseChanges(stdout: string): { changes: EntityChange[]; more: number } {
112 let parsed: unknown
113 try {
114 parsed = JSON.parse(stdout)
115 } catch {
116 return { changes: [], more: 0 }
117 }
118 const raw = isRecord(parsed) && Array.isArray(parsed.changes) ? parsed.changes : []
119 const changes: EntityChange[] = []
120 for (const c of raw) {
121 if (!isRecord(c) || !isRecord(c.type) || typeof c.type.name !== 'string' || typeof c.id !== 'string') continue
122 if (c.kind !== 'node' && c.kind !== 'edge') continue
123 if (c.op !== 'insert' && c.op !== 'update' && c.op !== 'delete') continue
124 const image = c.after ?? c.before
125 const change: EntityChange = { kind: c.kind, type: c.type.name, id: c.id, op: c.op }
126 const ends = isRecord(image) && isRecord(image.endpoints) ? image.endpoints : undefined
127 if (typeof ends?.from === 'string') change.from = ends.from
128 if (typeof ends?.to === 'string') change.to = ends.to
129 const label = labelIn(image)
130 if (label !== undefined) change.label = label
131 changes.push(change)
132 }
133 return { changes: changes.slice(0, MAX_CHANGES), more: Math.max(0, changes.length - MAX_CHANGES) }
134}
135src/enrich/highlight.ts 79 lines1import type { GhostGraph } from '../../types'
2
3export interface MatchIndex {
4 byKey: Map<string, string[]>
5 byLabel: Map<string, string[]>
6}
7
8export const MAX_MATCHES = 50
9
10function add(map: Map<string, string[]>, text: string, id: string): void {
11 const ids = map.get(text)
12 if (!ids) map.set(text, [id])
13 else if (!ids.includes(id)) ids.push(id)
14}
15
16export function buildMatchIndex(graph: GhostGraph): MatchIndex {
17 const byKey = new Map<string, string[]>()
18 const byLabel = new Map<string, string[]>()
19 for (const node of graph.nodes) {
20 add(byKey, node.key, node.id)
21 add(byLabel, node.label, node.id)
22 }
23 return { byKey, byLabel }
24}
25
26/** A `return { $p }` cell is the node as JSON: `{"@id": "…", …}`. */
27function exactKey(cell: string): string | undefined {
28 if (!cell.startsWith('{')) return undefined
29 try {
30 const parsed: unknown = JSON.parse(cell)
31 if (typeof parsed !== 'object' || parsed === null) return undefined
32 const id = (parsed as Record<string, unknown>)['@id']
33 return typeof id === 'string' ? id : undefined
34 } catch {
35 return undefined
36 }
37}
38
39function prefixMatch(map: Map<string, string[]>, prefix: string): string[] | undefined {
40 let found: string[] | undefined
41 for (const [text, ids] of map) {
42 if (!text.startsWith(prefix)) continue
43 if (found) return undefined
44 found = ids
45 }
46 return found
47}
48
49/**
50 * Nodes a read's result cells name: exact via `@id`, else a cell equal to a
51 * node key or label, else a truncated cell (`…`) matching one key or label.
52 * Numbers, booleans, `null` and short cells are ignored; at most 50 nodes.
53 */
54export function matchCells(cells: readonly string[][], index: MatchIndex): { nodeIds: string[]; isExact: boolean } {
55 const found = new Set<string>()
56 let isExact = true
57 for (const row of cells) {
58 for (const cell of row) {
59 if (found.size >= MAX_MATCHES) break
60 const text = cell.trim()
61 if (text.length < 3 || text === 'null' || text === 'true' || text === 'false' || /^-?\d+(\.\d+)?$/.test(text)) continue
62 const key = exactKey(text)
63 let ids: string[] | undefined
64 if (key !== undefined) {
65 ids = index.byKey.get(key)
66 } else {
67 ids = index.byKey.get(text) ?? index.byLabel.get(text)
68 if (!ids && text.endsWith('…') && text.length > 3) {
69 const prefix = text.slice(0, -1)
70 ids = prefixMatch(index.byKey, prefix) ?? prefixMatch(index.byLabel, prefix)
71 }
72 if (ids) isExact = false
73 }
74 for (const id of ids ?? []) if (found.size < MAX_MATCHES) found.add(id)
75 }
76 }
77 return { nodeIds: [...found], isExact: found.size > 0 && isExact }
78}
79src/enrich/locate.ts 18 lines1/** A bare command name, safe to hand to a shell as `command -v <name>`. */
2export function isPlainName(bin: string): boolean {
3 return /^[A-Za-z0-9._-]+$/.test(bin)
4}
5
6/** Where installers usually put a CLI, for when neither PATH nor the shell knows it. */
7export function binCandidates(home: string | undefined, bin: string): string[] {
8 const homeDirs = home !== undefined ? [`${home}/.local/bin`, `${home}/.cargo/bin`] : []
9 return [...homeDirs, '/opt/homebrew/bin', '/usr/local/bin'].map(dir => `${dir}/${bin}`)
10}
11
12/** The absolute path `command -v <bin>` printed, past any shell start-up noise. */
13export function pathFromShell(stdout: string, bin: string): string | undefined {
14 const lines = stdout.split('\n').map(l => l.trim()).filter(Boolean)
15 const last = lines[lines.length - 1]
16 return last !== undefined && last.startsWith('/') && last.endsWith(`/${bin}`) ? last : undefined
17}
18src/enrich/queue.ts 46 lines1export interface Job {
2 run: () => Promise<void>
3 /** When it was queued, in the same clock as `now`. */
4 at: number
5}
6
7/**
8 * Background work, one job at a time, in order. A job that waited longer
9 * than `maxAgeMs` is dropped: by then the moment it would have shown is gone.
10 */
11export class JobQueue {
12 private jobs: Job[] = []
13 private isRunning = false
14
15 constructor(
16 private readonly now: () => Promise<number>,
17 private readonly onError: (err: unknown) => void,
18 private readonly maxAgeMs = 60_000,
19 ) {}
20
21 get size(): number {
22 return this.jobs.length + (this.isRunning ? 1 : 0)
23 }
24
25 push(job: Job): void {
26 this.jobs.push(job)
27 if (!this.isRunning) void this.drain()
28 }
29
30 private async drain(): Promise<void> {
31 this.isRunning = true
32 try {
33 for (let job = this.jobs.shift(); job; job = this.jobs.shift()) {
34 if ((await this.now()) - job.at > this.maxAgeMs) continue
35 try {
36 await job.run()
37 } catch (err) {
38 this.onError(err)
39 }
40 }
41 } finally {
42 this.isRunning = false
43 }
44 }
45}
46src/graph/ghost.ts 126 lines1import type { EntityChange, GhostEdge, GhostGraph, GhostNode } from '../../types'
2
3export const nodeId = (type: string, key: string): string => `${type}:${key}`
4
5/** `edge Name: From -> To` declarations of a `.pg` schema. */
6export function parseEdgeEnds(schemaSource: string): Record<string, [string, string]> {
7 const ends: Record<string, [string, string]> = {}
8 for (const m of schemaSource.matchAll(/^\s*edge\s+(\w+)\s*:\s*(\w+)\s*->\s*(\w+)/gm)) {
9 ends[m[1]!] = [m[2]!, m[3]!]
10 }
11 return ends
12}
13
14export function labelOf(data: Record<string, unknown>, key: string): string {
15 for (const field of ['name', 'title', 'slug']) {
16 const value = data[field]
17 if (typeof value === 'string' && value !== '') return value
18 }
19 return key
20}
21
22/** The node id a raw key names; with several types sharing it, `type` picks one. */
23export function resolveKey(graph: GhostGraph, key: string, type?: string): string | undefined {
24 if (type !== undefined) {
25 const id = nodeId(type, key)
26 return graph.nodes.some(n => n.id === id) ? id : undefined
27 }
28 const found = graph.nodes.filter(n => n.key === key)
29 return found.length === 1 ? found[0]!.id : undefined
30}
31
32function resolveEnd(byKey: Map<string, GhostNode[]>, key: string, type: string | undefined): string | undefined {
33 const candidates = byKey.get(key) ?? []
34 if (type !== undefined) return candidates.find(n => n.type === type)?.id ?? (candidates.length === 1 ? candidates[0]!.id : undefined)
35 return candidates.length === 1 ? candidates[0]!.id : undefined
36}
37
38function keyIndex(nodes: readonly GhostNode[]): Map<string, GhostNode[]> {
39 const byKey = new Map<string, GhostNode[]>()
40 for (const node of nodes) byKey.set(node.key, [...(byKey.get(node.key) ?? []), node])
41 return byKey
42}
43
44/**
45 * Folds `omnigraph export` JSONL (`{"type","id","data"}` nodes and
46 * `{"edge","id","from","to","data"}` edges) into a ghost graph. Above
47 * `maxNodes`, the best-connected nodes are kept (ties in export order).
48 */
49export function foldExport(jsonl: string, edgeEnds: Record<string, [string, string]>, maxNodes: number): GhostGraph {
50 const nodes: GhostNode[] = []
51 const rawEdges: { id: string; type: string; from: string; to: string }[] = []
52 for (const line of jsonl.split('\n')) {
53 if (line.trim() === '') continue
54 let row: unknown
55 try {
56 row = JSON.parse(line)
57 } catch {
58 continue
59 }
60 if (typeof row !== 'object' || row === null) continue
61 const r = row as Record<string, unknown>
62 const data = typeof r.data === 'object' && r.data !== null ? (r.data as Record<string, unknown>) : {}
63 if (typeof r.edge === 'string' && typeof r.from === 'string' && typeof r.to === 'string') {
64 rawEdges.push({ id: typeof r.id === 'string' ? r.id : `${r.edge}|${r.from}|${r.to}`, type: r.edge, from: r.from, to: r.to })
65 } else if (typeof r.type === 'string' && typeof r.id === 'string') {
66 nodes.push({ id: nodeId(r.type, r.id), type: r.type, key: r.id, label: labelOf(data, r.id) })
67 }
68 }
69 const byKey = keyIndex(nodes)
70 let edges: GhostEdge[] = []
71 for (const e of rawEdges) {
72 const ends = edgeEnds[e.type]
73 const from = resolveEnd(byKey, e.from, ends?.[0])
74 const to = resolveEnd(byKey, e.to, ends?.[1])
75 if (from !== undefined && to !== undefined) edges.push({ id: e.id, type: e.type, from, to })
76 }
77 if (nodes.length <= maxNodes) return { nodes, edges, truncated: false, edgeEnds }
78
79 const degree = new Map<string, number>()
80 for (const e of edges) {
81 degree.set(e.from, (degree.get(e.from) ?? 0) + 1)
82 degree.set(e.to, (degree.get(e.to) ?? 0) + 1)
83 }
84 const ranked = nodes
85 .map((node, order) => ({ node, order, degree: degree.get(node.id) ?? 0 }))
86 .sort((a, b) => b.degree - a.degree || a.order - b.order)
87 .slice(0, maxNodes)
88 .sort((a, b) => a.order - b.order)
89 const kept = new Set(ranked.map(r => r.node.id))
90 edges = edges.filter(e => kept.has(e.from) && kept.has(e.to))
91 return { nodes: ranked.map(r => r.node), edges, truncated: true, edgeEnds }
92}
93
94/** A commit's changes applied to the graph (a new graph; the input is untouched). */
95export function applyChanges(graph: GhostGraph, changes: readonly EntityChange[]): GhostGraph {
96 let nodes = graph.nodes.slice()
97 let edges = graph.edges.slice()
98 for (const c of changes) {
99 if (c.kind === 'node') {
100 const id = nodeId(c.type, c.id)
101 const at = nodes.findIndex(n => n.id === id)
102 if (c.op === 'delete') {
103 if (at >= 0) nodes.splice(at, 1)
104 edges = edges.filter(e => e.from !== id && e.to !== id)
105 } else if (at >= 0) {
106 if (c.label !== undefined) nodes[at] = { ...nodes[at]!, label: c.label }
107 } else {
108 nodes.push({ id, type: c.type, key: c.id, label: c.label ?? c.id })
109 }
110 continue
111 }
112 if (c.op === 'delete') {
113 edges = edges.filter(e => e.id !== c.id)
114 continue
115 }
116 if (c.from === undefined || c.to === undefined || edges.some(e => e.id === c.id)) continue
117 const byKey = keyIndex(nodes)
118 const ends = graph.edgeEnds[c.type]
119 const from = resolveEnd(byKey, c.from, ends?.[0])
120 const to = resolveEnd(byKey, c.to, ends?.[1])
121 if (from !== undefined && to !== undefined) edges.push({ id: c.id, type: c.type, from, to })
122 }
123 nodes = nodes.slice()
124 return { ...graph, nodes, edges }
125}
126src/model/ledger.ts 96 lines1import { resolvePath } from '../util/path'
2import type { LedgerLine } from '../../types'
3
4/** The `$` calls the writer needs, bound to whichever hook is appending. */
5export interface LedgerIO {
6 write: (path: string, text: string) => Promise<void>
7 onError?: (err: unknown) => void
8}
9
10export function encodeLedger(lines: readonly LedgerLine[]): string {
11 return lines.map(line => `${JSON.stringify(line)}\n`).join('')
12}
13
14function isLedgerLine(v: unknown): v is LedgerLine {
15 if (typeof v !== 'object' || v === null) return false
16 const { v: version, t } = v as { v?: unknown; t?: unknown }
17 return version === 1 && typeof t === 'string'
18}
19
20/** Parses a ledger; a cut-off last line (a crash mid-write) is dropped silently. */
21export function decodeLedger(text: string): { lines: LedgerLine[]; errors: string[] } {
22 const raw = text.split('\n')
23 let last = raw.length - 1
24 while (last >= 0 && raw[last]!.trim() === '') last--
25 const lines: LedgerLine[] = []
26 const errors: string[] = []
27 raw.forEach((row, i) => {
28 if (row.trim() === '') return
29 let value: unknown
30 try {
31 value = JSON.parse(row)
32 } catch {
33 if (i !== last) errors.push(`line ${i + 1}: not JSON`)
34 return
35 }
36 if (isLedgerLine(value)) lines.push(value)
37 else errors.push(`line ${i + 1}: not a v1 ledger line`)
38 })
39 return { lines, errors }
40}
41
42export function ledgerFileName(nowMs: number): string {
43 const d = new Date(nowMs)
44 const pad = (n: number, width = 2) => String(n).padStart(width, '0')
45 return (
46 `${d.getUTCFullYear()}${pad(d.getUTCMonth() + 1)}${pad(d.getUTCDate())}-` +
47 `${pad(d.getUTCHours())}${pad(d.getUTCMinutes())}${pad(d.getUTCSeconds())}-` +
48 `${pad(d.getUTCMilliseconds(), 3)}.jsonl`
49 )
50}
51
52/** Absolute: `$.fs` resolves relative paths against the process, not the session. */
53export function ledgerPathFor(cwd: string, ledgerDir: string, nowMs: number): string {
54 return resolvePath(cwd, `${ledgerDir.replace(/\/+$/, '')}/${ledgerFileName(nowMs)}`)
55}
56
57/**
58 * Keeps every line in memory and writes through: each append rewrites the
59 * file, and appends made while a write is in flight ride on the next one. No
60 * timer, so a hot reload (which cancels pending waits) loses nothing.
61 */
62export class LedgerWriter {
63 readonly lines: LedgerLine[] = []
64 private isDirty = false
65 private isWriting = false
66
67 constructor(
68 readonly path: string,
69 initial: readonly LedgerLine[] = [],
70 ) {
71 this.lines.push(...initial)
72 }
73
74 append(line: LedgerLine, io: LedgerIO): void {
75 this.lines.push(line)
76 this.isDirty = true
77 if (!this.isWriting) void this.drain(io)
78 }
79
80 private async drain(io: LedgerIO): Promise<void> {
81 this.isWriting = true
82 try {
83 while (this.isDirty) {
84 this.isDirty = false
85 try {
86 await io.write(this.path, encodeLedger(this.lines))
87 } catch (err) {
88 io.onError?.(err)
89 }
90 }
91 } finally {
92 this.isWriting = false
93 }
94 }
95}
96src/model/reduce.ts 104 lines1import { applyChanges, resolveKey } from '../graph/ghost'
2import { basename } from '../util/path'
3import type { Counters, EntityChange, GhostGraph, LedgerLine, OpOutcome, Target, VizState } from '../../types'
4
5export function emptyViz(): VizState {
6 return {
7 ops: [],
8 counters: { ops: 0, errors: 0, nodes: 0, edges: 0 },
9 graph: null,
10 branch: null,
11 ghost: null,
12 touches: {},
13 }
14}
15
16/** A short name for where an op ran: the graph id, else the store/uri/cluster folder, else the server. */
17export function targetLabel(t: Target): string | null {
18 if (t.graph) return t.graph
19 const where = t.store ?? t.uri ?? t.cluster
20 if (where) return basename(where.replace(/\/+$/, ''))
21 return t.server ?? null
22}
23
24function outcomeBranch(o: OpOutcome): string | null {
25 switch (o.kind) {
26 case 'read':
27 return o.branch ?? null
28 case 'change':
29 case 'load':
30 return o.branch || null
31 case 'branch':
32 return o.action === 'merged' ? (o.into ?? null) : null
33 default:
34 return null
35 }
36}
37
38function addCounts(c: Counters, o: OpOutcome): Counters {
39 if (o.kind === 'error') return { ...c, errors: c.errors + 1 }
40 if (o.kind === 'change') return { ...c, nodes: c.nodes + o.nodes, edges: c.edges + o.edges }
41 if (o.kind === 'load') return { ...c, nodes: c.nodes + o.total }
42 return c
43}
44
45/** Node ids a diff touched: changed nodes and both ends of changed edges. */
46function diffTouches(graph: GhostGraph, changes: readonly EntityChange[]): string[] {
47 const ids = new Set<string>()
48 for (const c of changes) {
49 if (c.kind === 'node') {
50 ids.add(`${c.type}:${c.id}`)
51 continue
52 }
53 const ends = graph.edgeEnds[c.type]
54 for (const [key, type] of [[c.from, ends?.[0]], [c.to, ends?.[1]]] as const) {
55 const id = key === undefined ? undefined : resolveKey(graph, key, type)
56 if (id !== undefined) ids.add(id)
57 }
58 }
59 return [...ids]
60}
61
62export function reduce(s: VizState, line: LedgerLine): VizState {
63 if (line.t === 'graph') return { ...s, ghost: line.graph }
64 if (line.t === 'touch') return { ...s, touches: { ...s.touches, [line.opId]: line.nodeIds } }
65 if (line.t === 'diff') {
66 if (!s.ghost) return s
67 const ghost = applyChanges(s.ghost, line.changes)
68 return { ...s, ghost, touches: { ...s.touches, [line.opId]: diffTouches(ghost, line.changes) } }
69 }
70 if (line.t === 'op.start') {
71 return {
72 ...s,
73 ops: [...s.ops, line.op],
74 counters: { ...s.counters, ops: s.counters.ops + 1 },
75 graph: targetLabel(line.op.intent.target) ?? s.graph,
76 branch: line.op.intent.branch ?? s.branch,
77 }
78 }
79 if (line.t === 'op.end') {
80 const i = s.ops.findIndex(o => o.id === line.id)
81 const prev = s.ops[i]
82 if (!prev) return s
83 const ops = s.ops.slice()
84 ops[i] = {
85 ...prev,
86 status: line.status,
87 durationMs: line.durationMs,
88 outcome: line.outcome,
89 output: line.output,
90 }
91 return {
92 ...s,
93 ops,
94 counters: addCounts(s.counters, line.outcome),
95 branch: outcomeBranch(line.outcome) ?? s.branch,
96 }
97 }
98 return s
99}
100
101export function reduceAll(lines: readonly LedgerLine[], from: VizState = emptyViz()): VizState {
102 return lines.reduce(reduce, from)
103}
104src/model/state.ts 44 lines1import type { GraphState, OpEvent, ReplayView, UiState, ViewState, VizState } from '../../types'
2
3export const RECENT_OPS = 200
4
5export const INITIAL_UI: UiState = { tab: 'graph', selectedOpId: null }
6
7export const INITIAL_VIEW: ViewState = {
8 recentOps: [],
9 counters: { ops: 0, errors: 0, nodes: 0, edges: 0 },
10 graph: null,
11 branch: null,
12 ghost: { state: 'none', nodes: 0, edges: 0, truncated: false, message: null },
13 replay: null,
14}
15
16/** What the Feed lists: no raw output, no read cells (those stay in module memory). */
17function slim(op: OpEvent): OpEvent {
18 const { output: _output, ...rest } = op
19 if (rest.outcome?.kind === 'read') return { ...rest, outcome: { ...rest.outcome, cells: [] } }
20 return rest
21}
22
23export function toView(
24 viz: VizState,
25 state: GraphState = 'none',
26 message: string | null = null,
27 replay: ReplayView | null = null,
28): ViewState {
29 return {
30 ghost: {
31 state: viz.ghost ? 'ready' : state,
32 nodes: viz.ghost?.nodes.length ?? 0,
33 edges: viz.ghost?.edges.length ?? 0,
34 truncated: viz.ghost?.truncated ?? false,
35 message,
36 },
37 recentOps: viz.ops.slice(-RECENT_OPS).map(slim),
38 counters: viz.counters,
39 graph: viz.graph,
40 branch: viz.branch,
41 replay,
42 }
43}
44