SLOPSHOPPER

omnigraph-viz

Live terminal visualization of Omnigraph operations

newpaneguardcommandstatusprocess
A shopper browsing a rack in a slop shop
Preview · a replayed session in a sandbox
claude · ~/work/app · omnigraph-viz
│ ┃ omnigraph ✕ › fix the failing auth test and add an audit log call │ ┃ ◉ omnigraph · no graph yet · main 0 ops · │ ┃ [ Graph ] [ Feed ] following ops · click the ⏺ Read(src/auth.ts) │ ┃ Waiting for omnigraph operations… the graph ⎿ Read 6 lines │ ┃ after the first command that names a graph. ⏺ Update(src/auth.ts) │ ⎿ Added 2 lines, removed 1 line │ ⏺ Bash(bun test) │ ⎿ 3 pass, 1 fail │ │ ● Done. refresh now rejects expired claims and logs an audit event. │ │ ✻ Worked for 42s · done 4:20 PM │ │ › /og │ ⎿ omnigraph-viz: Omnigraph viz opened. │ │ ────────────────────────────────────────────────────────────────────────────────────────────────────────────────────── › ? for shortcuts

Draws

Pane · omnigraph
◉ omnigraph · no graph yet · main 0 ops · +0n +0e · 0… [ Graph ] [ Feed ] following ops · click the graph to browse Waiting for omnigraph operations… the graph appears after the first command that names a graph.
README

omnigraph-claude-mod

Claude Code with the omnigraph-viz pane: the agent's bug-bash triage on the left; on the right, the graph browser on the presence component, its neighbours around it and Finder-style columns below

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.

What you get

  • Graph tab (default): a graph browser.
  • Breadcrumb: reads as a sentence, e.g. Epic › Offline editing & sync › contains › ….
  • Picture: the selected node in the middle, its neighbours around it as ● label, links as thin dotted lines.
  • Finder-style columns: node types → a type's nodes → a node's links grouped by relation → the next hop.
  • Following: it follows the agent as it works.
  • A read lights its results around the query itself (◆ query ready · 5).
  • A new node flashes in, marked NEW.
  • A new link draws in from its source.
  • Feed tab: one line per operation, with counts and timings; enter shows the command and its output.
  • Status line: ◉ og · main · 14 ops · +3n.
  • Replay: /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.

Requirements

  • Claude Code with mods (plugins of function hooks). Built and tested against Claude Code 2.1.291.
  • The omnigraph CLI on your PATH, or its path in the omnigraphBin option.
  • A truecolor terminal and a font with braille glyphs (Menlo, SF Mono, JetBrains Mono).
  • A wide window (144+ columns) so the pane sits beside the conversation.

Install

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

Use

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.

/ogtoggle the pane (it opens with the keyboard)
/og replay <ledger.jsonl> [speed]replay a recorded session; 0.5 is half speed, 2 double
/og stopend a replay, back to the live session
/og clearreset the Feed (the graph stays)
g / fGraph / Feed tab
click the graphhand 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
lfollow the agent again
enter (Feed)show an operation's command and output; select it to see it in the browser
esckeys back to the prompt

For Finder's ⌘↓ / ⌘↑ in iTerm2, map them to send ⌥↓ / ⌥↑ (Profile → Keys → Key Mappings → Send Escape Sequence [1;3B / [1;3A).

Try it in a minute

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.

Options

OptionDefaultWhat it does
omnigraphBinomnigraphName 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.
maxGhostNodes2000Largest graph loaded whole; a bigger one keeps its best-connected nodes
ledgerDir.omnigraph-vizWhere session ledgers go, relative to the session directory
enrichtrueRun omnigraph … --json in the background for graph detail (export once, commit changes after writes)
debugfalseLog 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"}}}}'

How it works

  1. Capture. A 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.
  2. Enrich. In the background, the mod runs its own 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.
  3. Ledger. Every operation, export and diff is appended to a JSONL ledger in .omnigraph-viz/. The mod adds a .gitignore there, because ledgers hold command output.
  4. One path for live and replay. Every visible change is a ledger line folded by one reducer. A replay feeds a recorded ledger through the same path, timed by the recording, so it draws the same frames every time.

Development

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.

License

MIT

Source 29 files
hooks/register.tsx 602 lines
1import { 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}
602
src/capture/capture.ts 130 lines
1import 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}
130
src/capture/parse-output.ts 286 lines
1import 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}
286
src/capture/tokenize.ts 199 lines
1import 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}
199
src/enrich/cli.ts 135 lines
1import 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}
135
src/enrich/highlight.ts 79 lines
1import 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}
79
src/enrich/locate.ts 18 lines
1/** 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}
18
src/enrich/queue.ts 46 lines
1export 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}
46
src/graph/ghost.ts 126 lines
1import 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}
126
src/model/ledger.ts 96 lines
1import { 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}
96
src/model/reduce.ts 104 lines
1import { 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}
104
src/model/state.ts 44 lines
1import 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