Live metro map of a running Nextflow pipeline, inside Claude Code

A Claude Code mod: when Claude runs a Nextflow pipeline, a pane draws the pipeline as a live metro map.
FETCHNGS · stupefied_becquerel
●━━━━━━━━━━━━━━●━━━━━━┳╌╌╌╌╌╌╌◌╌╌╌╌╌╌┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━●
SRA_IDS_TO_ SRA_RUNINFO_ ┃ ASPERA_CLI ┃ ╎ MULTIQC_
RUNINFO TO_FTP ┃ ┃ ╎ MAPPIN…ONFIG
┣╌╌╌╌╌╌╌◌╌╌╌╌╌╌┫ ╎
┃ SRA_FASTQ_ ┃ ╎
┃ FTP ┃ ╎
┣━━━━━━━●━━━━━━┫ ╎
╎ FASTQDL ╎ ╎
╌╌╌╌╌╌╌╌◌╌╌╌╌╌╌╌╌╌╌╌╌╌╌◌╌╌╌╌╌╌╌╌╌╌╌╌╌╌◌╌╌╌╌╌╌╌
CUSTOM_ SRATOOLS_ SRATOOLS_
SRATOO…TINGS PREFETCH FASTERQDUMP
╰ FASTQ_DOWNLOAD_PREFETCH_FASTERQDUMP_SR─╯
○ pending · ◉ running · ● done · ✖ failed · ◌ never ran
claude plugin marketplace add camlloyd/metro-map-mod
claude plugin install metro-map-mod@metro-map-mod
Needs Nextflow and python3. /metro-map opens the pane.
Try it: ask Claude to run
echo DRR028935 > ids.csv
nextflow run nf-core/fetchngs -profile test,docker --input ids.csv --download_method fastq-dl --outdir results
-c <weblog config> to each nextflow run Claude starts (-with-weblog when the command uses -C); a local listener turns the events into station states.-preview -with-dag for the graph, in .nextflow/metro-map/ with its own NXF_CACHE_DIR, so your history is untouched.-resume, cached tasks come from .nextflow.log.◀ / ▶ beside the title, shown only when there's more that way, move it (click, or ctrl+x tab then Enter).fetchngs done in 2m 3s · 6 succeeded, or fetchngs failed after 40s at SRA_FASTQ_FTP.claude --plugin-dir .
claude plugin test .
hooks/register.tsx 201 lines1import { atom, read, update } from 'claude-code'
2import type { EngineInterface, Register, Timer } from 'claude-code'
3
4import type { Dag, Run } from '../types'
5import { metro, parseDot, pipelineOf, STATION } from './dag'
6import { applyCached, applyEvent, dagConfig, endToast, previewArgv, statusOf, tally, WEBLOG_CONFIG, withWeblog } from './weblog'
7import type { WeblogEvent } from './weblog'
8
9const PANE = 'metro-map'
10const openPane = ($: EngineInterface) => $.ui.open({ id: PANE, title: 'Metro map' })
11const run = atom({ plugin: 'metro-map-mod', key: 'run' } as const, null)
12const dag = atom({ plugin: 'metro-map-mod', key: 'dag' } as const, null)
13const lastPort = atom({ plugin: 'metro-map-mod', key: 'port' } as const, null)
14// Columns the person moved the map from where it follows the run by itself.
15const pan = atom({ plugin: 'metro-map-mod', key: 'pan' } as const, 0)
16
17// A local endpoint for Nextflow's weblog: writes a config file pointing at it, prints its port and that file, then each
18// event posted to it as one JSON line. It takes the port it had before a reload (argv[1]) if that frees up within ~2s,
19// so runs already reporting keep their map.
20const LISTENER = `
21import http.server, os, sys, tempfile, time
22class H(http.server.BaseHTTPRequestHandler):
23 def do_POST(self):
24 body = self.rfile.read(int(self.headers.get('Content-Length', 0)))
25 sys.stdout.write(body.decode('utf-8', 'replace').replace('\\n', ' ') + '\\n'); sys.stdout.flush()
26 self.send_response(200); self.end_headers()
27 def log_message(self, *a): pass
28srv = None
29for _ in range(20):
30 try: srv = http.server.ThreadingHTTPServer(('127.0.0.1', int(sys.argv[1])), H); break
31 except (OSError, ValueError, IndexError): time.sleep(0.1)
32srv = srv or http.server.ThreadingHTTPServer(('127.0.0.1', 0), H)
33port = srv.server_address[1]
34config = os.path.join(tempfile.mkdtemp(prefix='metro-map-'), '${WEBLOG_CONFIG}')
35with open(config, 'w') as f: f.write("weblog { enabled = true; url = 'http://127.0.0.1:%d/events' }\\n" % port)
36print('PORT', port, config, flush=True)
37srv.serve_forever()
38`
39
40let weblogUrl: string | undefined
41let weblogConfig: string | undefined
42let cachedTimer: Timer | undefined
43
44// The session's weblog listener; a reload restarts it (session.start fires again) and kills the old one.
45async function listen($: EngineInterface) {
46 let buffer = ''
47 try {
48 const before = String(await read($, lastPort) ?? 0)
49 for await (const { stream, text } of $.process.spawn({ argv: ['python3', '-c', LISTENER, before] })) {
50 if (stream !== 'stdout') continue
51 buffer += text
52 let nl: number
53 while ((nl = buffer.indexOf('\n')) >= 0) {
54 const line = buffer.slice(0, nl)
55 buffer = buffer.slice(nl + 1)
56 if (line.startsWith('PORT ')) {
57 const [, port, config] = line.trim().split(' ')
58 weblogUrl = `http://127.0.0.1:${port}/events`
59 weblogConfig = config
60 await update($, lastPort, () => Number(port))
61 }
62 else if (line.startsWith('{')) await onEvent($, JSON.parse(line) as WeblogEvent).catch(() => undefined)
63 }
64 }
65 } finally {
66 weblogUrl = undefined
67 $.ui.status('metro map: weblog listener stopped (is python3 on PATH?)')
68 }
69}
70
71async function onEvent($: EngineInterface, ev: WeblogEvent) {
72 const before = await read($, run)
73 const applied = applyEvent(before, ev)
74 if (!applied || applied === before) return
75 const next = ev.event === 'started' ? { ...applied, startedAt: Date.now() } : applied
76 await update($, run, () => next)
77
78 if (ev.event === 'started') {
79 await update($, dag, () => null)
80 await update($, pan, () => 0)
81 void openPane($)
82 void preview($, next)
83 cachedTimer?.cancel()
84 // ponytail: assumes the default .nextflow.log in the launch dir; a `-log elsewhere` resume shows cached stations as skipped
85 if (next.isResume) cachedTimer = $.clock.every(2000, () => void fillCached($, next.id))
86 }
87 if (next.status !== 'running') {
88 cachedTimer?.cancel()
89 if (next.isResume) await fillCached($, next.id)
90 const d = await read($, dag)
91 const pipeline = (d?.runId === next.id && pipelineOf(d.nodes).title.toLowerCase()) || next.name || 'Nextflow run'
92 $.ui.toast(endToast((await read($, run)) ?? next, pipeline, next.startedAt && Date.now() - next.startedAt))
93 }
94}
95
96async function fillCached($: EngineInterface, id: string) {
97 const r = await read($, run)
98 if (!r || r.id !== id) return
99 const text = await $.fs.read(`${r.dir}/.nextflow.log`).catch(() => '')
100 await update($, run, cur => (cur && cur.id === id ? applyCached(cur, text) : cur))
101}
102
103// Re-runs the command with -preview for the DAG, since -with-dag only writes it when a run ends. NXF_CACHE_DIR keeps
104// the preview out of the run's .nextflow/history.
105async function preview($: EngineInterface, r: Run) {
106 const tmp = `${r.dir}/.nextflow/metro-map`
107 const argv = previewArgv(r.command, tmp)
108 const fail = (error: string) => update($, dag, () => ({ runId: r.id, nodes: [], edges: [], error }) satisfies Dag)
109 if (!argv) return fail('no `nextflow run` in the command line')
110 try {
111 await $.process.run(['mkdir', '-p', tmp])
112 await $.fs.write(`${tmp}/dag.config`, dagConfig(tmp))
113 // ponytail: ${tmp}/cache grows ~12 KB per preview and is never pruned
114 const ran = await $.process.run(argv, { cwd: r.dir, env: { NXF_CACHE_DIR: `${tmp}/cache` }, timeoutMs: 5 * 60_000 })
115 if (ran.exitCode !== 0) return fail(`preview exited ${ran.exitCode}: ${(ran.stderr || ran.stdout).trim().split('\n').pop() ?? ''}`)
116 const parsed = parseDot(await $.fs.read(`${tmp}/dag.dot`))
117 await update($, dag, () => ({ runId: r.id, ...parsed }) satisfies Dag)
118 } catch (err) {
119 await fail(String(err))
120 }
121}
122
123export const register: Register = on => {
124 on('session.start', async ($, e, next) => {
125 await $.command.register({ name: 'metro-map', description: 'Show the metro map of the current Nextflow run' })
126 const started = await next(e)
127 void listen($)
128 return started
129 })
130
131 on('command.run', { command: 'metro-map' }, async $ => {
132 await openPane($)
133 return { text: 'Metro map opened.' }
134 })
135
136 // Every `nextflow run` the agent starts reports to the listener.
137 on('tool.call', { tool: 'Bash' }, async ($, e, next) => {
138 const command = weblogUrl && weblogConfig ? withWeblog(e.command, weblogUrl, weblogConfig) : undefined
139 if (!command) return next(e)
140 await update($, run, () => ({ id: '', name: '', dir: '', command: '', isResume: false, status: 'waiting', procs: [] }) satisfies Run)
141 void openPane($)
142 return next({ ...e, command })
143 })
144
145 on('ui.render', { component: 'Pane', requestId: PANE }, async ($, e) => {
146 const { Box, Button, Text } = $.ui.resolve(e)
147 const r = await read($, run)
148 if (!r) return <Text dimColor>No Nextflow run yet. The map opens when one starts.</Text>
149
150 const room = Math.max(3, (e.viewport?.rows ?? 30) - 5)
151 const { failed } = tally(r.procs)
152 const head = r.status === 'waiting' ? 'waiting for Nextflow to start…' : `${r.status}${failed ? ` · ${failed} failed` : ''}`
153 const headColor = r.status === 'failed' ? 'red' : r.status === 'done' ? 'green' : undefined
154
155 // The metro map once the -preview DAG is in.
156 const d = await read($, dag)
157 if (d && d.runId === r.id && d.nodes.length > 0) {
158 const byName = new Map(r.procs.map(p => [p.name, p]))
159 const { rows, title, from, follow, count, of } = metro(d.nodes, d.edges, n => {
160 const st = statusOf(byName.get(n))
161 return st === 'pending' && r.status !== 'running' ? 'skipped' : st
162 }, e.props.bodyColumns ?? e.viewport?.columns ?? 100, await read($, pan))
163 // A page less one column per press, so one station stays in view; set from where it shows, so an edge doesn't stick.
164 const step = Math.max(1, count - 1)
165 const move = (by: number) => () => void update($, pan, () => from + by - follow)
166 return (
167 <Box flexDirection="column">
168 <Text bold color={headColor}>{head}</Text>
169 <Box flexDirection="row">
170 {from > 0 && <Button plain label="◀" onPress={move(-step)} />}
171 <Text dimColor>{from > 0 ? ' ' : ''}{[title, r.name].filter(Boolean).join(' · ')}{from + count < of ? ' ' : ''}</Text>
172 {from + count < of && <Button plain label="▶" onPress={move(step)} />}
173 </Box>
174 {rows.slice(0, room).map(row => (
175 <Text wrap="truncate">
176 {row.map(s => <Text color={s.color} bold={s.bold} dimColor={s.dim}>{s.text}</Text>)}
177 </Text>
178 ))}
179 </Box>
180 )
181 }
182
183 // Until then (or if the preview failed): the processes seen so far.
184 return (
185 <Box flexDirection="column">
186 <Text bold color={headColor}>{head}</Text>
187 {r.dir !== '' && <Text dimColor>{[r.name, d?.error ? `DAG preview failed: ${d.error}` : 'building the map…'].filter(Boolean).join(' · ')}</Text>}
188 {r.procs.slice(-room).map(p => {
189 const s = STATION[statusOf(p)]
190 return (
191 <Text wrap="truncate">
192 <Text color={s.color} bold={s.bold} dimColor={s.dim}>{s.text}</Text> {p.name}
193 {p.failed > 0 && <Text dimColor> ✖{p.failed}</Text>}
194 </Text>
195 )
196 })}
197 </Box>
198 )
199 })
200}
201hooks/dag.ts 244 lines1import type { Dag } from '../types'
2
3/** Reads Nextflow's `-with-dag x.dot` into process nodes and process->process edges (operators folded away). */
4export function parseDot(dot: string): Pick<Dag, 'nodes' | 'edges'> {
5 const label = new Map<string, string>() // vN -> process name; operators/channels have none
6 const out = new Map<string, string[]>()
7 for (const line of dot.split('\n')) {
8 const node = line.match(/^\s*(v\d+) \[(.*)\];\s*$/)
9 if (node && !/shape=/.test(node[2]!)) {
10 const name = node[2]!.match(/label="([^"]*)"/)?.[1]
11 if (name) label.set(node[1]!, name)
12 }
13 const edge = line.match(/^\s*(v\d+) -> (v\d+)/)
14 if (edge) out.set(edge[1]!, [...(out.get(edge[1]!) ?? []), edge[2]!])
15 }
16 const nodes = [...new Set(label.values())]
17 const edges: [string, string][] = []
18 for (const [v, name] of label) {
19 // Walk through operator nodes to the next processes.
20 const seen = new Set<string>(), stack = [...(out.get(v) ?? [])], hit = new Set<string>()
21 while (stack.length) {
22 const w = stack.pop()!
23 if (seen.has(w)) continue
24 seen.add(w)
25 const to = label.get(w)
26 if (to) hit.add(to)
27 else stack.push(...(out.get(w) ?? []))
28 }
29 for (const to of hit) if (to !== name) edges.push([name, to])
30 }
31 return { nodes, edges }
32}
33
34export type Status = 'pending' | 'running' | 'done' | 'failed' | 'skipped'
35export type Seg = { text: string; color?: string; bold?: boolean; dim?: boolean }
36
37// Transit-map line colours, one per route.
38export const PALETTE = ['#e4002b', '#0098d4', '#00a651', '#f3a900', '#9b5ba5', '#ef7b10', '#00afad', '#b26300']
39// Track not taken (yet): never ran, or not reached.
40export const DIM = '#555555'
41export const STATION: Record<Status, Seg> = {
42 skipped: { text: '◌', dim: true },
43 pending: { text: '○', dim: true },
44 running: { text: '◉', color: '#ffd700', bold: true },
45 done: { text: '●', color: '#ffffff' },
46 failed: { text: '✖', color: '#ff4040', bold: true },
47}
48const LABEL: Record<Status, Omit<Seg, 'text'>> = {
49 pending: { dim: true }, skipped: { dim: true }, running: { color: '#ffd700', bold: true }, done: {}, failed: { color: '#ff4040' },
50}
51// Heavy box glyph by which sides connect: up, down, left, right.
52const BOX: Record<string, string> = { // horizontal-only falls back to ━
53 '1100': '┃', '1000': '┃', '0100': '┃',
54 '0101': '┏', '0110': '┓', '1001': '┗', '1010': '┛', '0111': '┳', '1011': '┻',
55 '1101': '┣', '1110': '┫', '1111': '╋',
56}
57export const short = (n: string) => n.split(':').pop() ?? n
58
59/** Every process shares the pipeline's own workflow path (NFCORE_X:X): its last part is the title, `depth` its length. */
60export function pipelineOf(nodes: readonly string[]): { title: string; depth: number } {
61 const paths = nodes.map(n => n.split(':').slice(0, -1))
62 let depth = 0
63 while (paths.length && paths.every(p => p.length > depth && p[depth] === paths[0]![depth])) depth++
64 return { title: paths[0]?.slice(0, depth).at(-1) ?? '', depth }
65}
66
67/**
68 * Draws the DAG as an nf-core-style metro map, left to right: each analysis route a coloured line through
69 * its stations (processes) by graph depth, forks and merges as vertical links, subworkflows as bracketed sections.
70 * Returns rows of styled segments sized to `width` columns, and which columns they show
71 * (`count` of `of` from `from`; `follow` is where the window sits unshifted).
72 */
73export function metro(nodes: readonly string[], edges: readonly (readonly [string, string])[],
74 statusOf: (n: string) => Status, width: number, shift = 0): { rows: Seg[][]; title: string; from: number; follow: number; count: number; of: number } {
75 if (nodes.length === 0) return { rows: [], title: '', from: 0, follow: 0, count: 0, of: 0 }
76 const { title, depth } = pipelineOf(nodes)
77 const section = (n: string) => n.split(':').slice(depth, -1)[0] ?? ''
78 const kids = new Map<string, string[]>(), parents = new Map<string, string[]>(), indeg = new Map(nodes.map(n => [n, 0]))
79 for (const [a, b] of edges) {
80 kids.set(a, [...(kids.get(a) ?? []), b])
81 parents.set(b, [...(parents.get(b) ?? []), a])
82 indeg.set(b, (indeg.get(b) ?? 0) + 1)
83 }
84 // Kahn's sort, ties broken by the DOT's own order. ponytail: O(n²), fine for a few hundred processes.
85 const order: string[] = [], pending = [...nodes]
86 while (pending.length) {
87 const i = Math.max(0, pending.findIndex(n => (indeg.get(n) ?? 0) === 0))
88 const n = pending.splice(i, 1)[0]!
89 order.push(n)
90 for (const k of kids.get(n) ?? []) indeg.set(k, (indeg.get(k) ?? 1) - 1)
91 }
92 const layer = new Map<string, number>()
93 for (const n of order) layer.set(n, Math.max(0, ...(parents.get(n) ?? []).map(p => (layer.get(p) ?? 0) + 1)))
94 const L = Math.max(...layer.values()) + 1
95
96 // Routes: greedy path cover in topo order, each path one track.
97 const lineOf = new Map<string, number>(), lines: string[][] = []
98 for (const n of order) {
99 if (lineOf.has(n)) continue
100 const path = [n]
101 lineOf.set(n, lines.length)
102 for (let cur = n; ;) {
103 const next = (kids.get(cur) ?? []).filter(k => !lineOf.has(k)).sort((a, b) => layer.get(a)! - layer.get(b)!)[0]
104 if (!next) break
105 lineOf.set(next, lines.length)
106 path.push(next)
107 cur = next
108 }
109 lines.push(path)
110 }
111 const R = lines.length // x: even = gap before column x/2, odd = station column
112 const span = lines.map(p => [2 * layer.get(p[0]!)! + 1, 2 * layer.get(p[p.length - 1]!)! + 1] as [number, number])
113 const taken = (n: string) => statusOf(n) !== 'pending' && statusOf(n) !== 'skipped'
114 const links: { a: string; b: string; x: number; r1: number; r2: number; color: string; isLit: boolean; ra: number; rb: number }[] = []
115 for (const [a, b] of edges) {
116 const ra = lineOf.get(a)!, rb = lineOf.get(b)!
117 if (ra === rb) continue
118 const isBranch = lines[rb]![0] === b, isJoin = lines[ra]!.at(-1) === a
119 const x = isJoin && !isBranch ? 2 * layer.get(a)! + 2 : 2 * layer.get(b)!
120 links.push({ a, b, x, r1: Math.min(ra, rb), r2: Math.max(ra, rb), ra, rb, isLit: taken(a) && taken(b),
121 color: PALETTE[(isBranch || !isJoin ? rb : ra) % PALETTE.length]! })
122 for (const r of [ra, rb]) span[r] = [Math.min(span[r]![0], x), Math.max(span[r]![1], x)]
123 }
124
125 // Lit stretches per track: where each edge that ran is drawn. A link runs along its source's track to its x,
126 // then along its target's track.
127 const litSpans: [number, number][][] = lines.map(() => [])
128 for (const [a, b] of edges) {
129 if (!taken(a) || !taken(b)) continue
130 const ra = lineOf.get(a)!, rb = lineOf.get(b)!, xa = 2 * layer.get(a)! + 1, xb = 2 * layer.get(b)! + 1
131 const l = links.find(k => k.a === a && k.b === b)
132 if (!l) { litSpans[ra]!.push([xa, xb]); continue }
133 litSpans[ra]!.push([Math.min(xa, l.x), Math.max(xa, l.x)])
134 litSpans[rb]!.push([Math.min(l.x, xb), Math.max(l.x, xb)])
135 }
136 const isLit = (r: number, x: number, side: 'left' | 'right') =>
137 litSpans[r]!.some(([lo, hi]) => (side === 'left' ? lo < x && x <= hi : lo <= x && x < hi))
138
139 // Fit: shrink station columns, then window around the first column still running or pending, moved `shift` columns.
140 const fitW = (cols: number) => Math.floor((width - 3 * (cols + 1)) / cols)
141 const W = Math.max(10, Math.min(14, fitW(L)))
142 const k = Math.max(1, Math.min(L, Math.floor((width - 3) / (W + 3))))
143 const active = order.find(n => statusOf(n) !== 'done')
144 const window = (at: number) => Math.max(0, Math.min(L - k, at))
145 const follow = window((active ? layer.get(active)! : L) - 1), c0 = window(follow + shift)
146 const x0 = 2 * c0, x1 = 2 * (c0 + k)
147 const cellW = (x: number) => (x % 2 === 0 ? 3 : W)
148 const station = new Map<string, string>() // "r,x" -> node
149 for (const [n, l] of layer) station.set(`${lineOf.get(n)},${2 * l + 1}`, n)
150
151 const rows: Seg[][] = []
152 const push = (row: Seg[], seg: Seg) => {
153 const last = row.at(-1)
154 if (last && last.color === seg.color && last.bold === seg.bold && last.dim === seg.dim) last.text += seg.text
155 else row.push({ ...seg })
156 }
157 // Too long for one row: break after the last `_` that fits (GATK4_ / MARKDUPLICATES).
158 const wrap = (name: string, w: number): [string, string] => {
159 if (name.length <= w) return [name, '']
160 const cut = name.lastIndexOf('_', w - 1)
161 return cut > 0 ? [name.slice(0, cut + 1), name.slice(cut + 1)] : [name.slice(0, w), name.slice(w)]
162 }
163 const centred = (text: string, w: number) => {
164 // Cut in the middle: the end often tells names apart (…_R1_FQ / …_R2_FQ).
165 const head = Math.ceil((w - 1) / 2)
166 const t = text.length > w ? text.slice(0, head) + '…' + text.slice(text.length - (w - 1 - head)) : text
167 const left = Math.floor((w - t.length) / 2)
168 return ' '.repeat(left) + t + ' '.repeat(w - t.length - left)
169 }
170
171 // The link passing down from track r at x, preferring a lit one.
172 const linkBelow = (r: number, x: number) => {
173 const here = links.filter(l => l.x === x && l.r1 <= r && r < l.r2)
174 return here.find(l => l.isLit) ?? here[0]
175 }
176
177 // A label or bracket cell: blank, or the link passing down through it.
178 const under = (row: Seg[], r: number, x: number) => {
179 const w = cellW(x), mid = Math.floor(w / 2), below = linkBelow(r, x)
180 if (!below) return push(row, { text: ' '.repeat(w) })
181 push(row, { text: ' '.repeat(mid) })
182 push(row, { text: '┃', color: below.isLit ? below.color : DIM })
183 push(row, { text: ' '.repeat(w - mid - 1) })
184 }
185
186 for (let r = 0; r < R; r++) {
187 const color = PALETTE[r % PALETTE.length]!
188 const track: Seg[] = [], label: Seg[] = [], label2: Seg[] = []
189 for (let x = x0; x <= x1; x++) {
190 const w = cellW(x), mid = Math.floor(w / 2)
191 const [lo, hi] = span[r]!
192 const left = x > lo && x <= hi, right = x >= lo && x < hi
193 const vs = links.filter(l => l.x === x && l.r1 <= r && r <= l.r2)
194 const up = vs.some(l => l.r1 < r), down = vs.some(l => l.r2 > r)
195 const node = station.get(`${r},${x}`)
196 const fill = (on: boolean, side: 'left' | 'right', n: number): Seg =>
197 on ? { text: '━'.repeat(n), color: isLit(r, x, side) ? color : DIM } : { text: ' '.repeat(n) }
198 push(track, fill(left, 'left', mid))
199 if (node) push(track, STATION[statusOf(node)])
200 else {
201 const key = `${+up}${+down}${+left}${+right}`
202 const g = BOX[key] ?? (left || right ? '━' : ' ')
203 const lit = vs.find(l => l.isLit)
204 const litL = isLit(r, x, 'left'), litR = isLit(r, x, 'right')
205 push(track, { text: g, color: lit ? lit.color : litL && litR ? color : up || down || !(litL || litR) ? DIM : color })
206 }
207 push(track, fill(right, 'right', w - mid - 1))
208 // Label rows: station names, and links passing down to the next track.
209 if (node) {
210 const [a, b] = wrap(short(node), w)
211 push(label, { text: centred(a, w), ...LABEL[statusOf(node)] })
212 push(label2, { text: centred(b, w), ...LABEL[statusOf(node)] })
213 } else for (const row of [label, label2]) under(row, r, x)
214 }
215 rows.push(track, label)
216 if (lines[r]!.some(n => { const x = 2 * layer.get(n)! + 1; return x >= x0 && x <= x1 && short(n).length > W })) rows.push(label2)
217
218 // Section brackets: runs of this line's stations inside one subworkflow, under their labels.
219 const runs: { from: number; to: number; name: string }[] = []
220 for (const n of lines[r]!) {
221 const x = 2 * layer.get(n)! + 1, name = section(n), last = runs.at(-1)
222 if (name && last?.name === name) last.to = x
223 else if (name) runs.push({ from: x, to: x, name })
224 }
225 if (runs.length === 0) continue
226 const bracket: Seg[] = []
227 for (let x = x0; x <= x1; x++) {
228 const run = runs.find(u => u.from <= x && x <= u.to)
229 if (run) {
230 let total = 0
231 for (let y = Math.max(run.from, x0); y <= Math.min(run.to, x1); y++) total += cellW(y)
232 if (x === Math.max(run.from, x0)) {
233 const name = ` ${run.name} `.slice(0, Math.max(0, total - 3))
234 push(bracket, { text: total < 3 ? ' '.repeat(total) : '╰' + name + '─'.repeat(total - 2 - name.length) + '╯', dim: true })
235 }
236 continue
237 }
238 under(bracket, r, x)
239 }
240 rows.push(bracket)
241 }
242 return { rows, title, from: c0, follow, count: k, of: L }
243}
244hooks/weblog.ts 146 lines1import type { Proc, Run } from '../types'
2import type { Status } from './dag'
3
4/** One event as Nextflow posts it to `weblog.url`. */
5export type WeblogEvent = {
6 event: string // started | process_submitted | process_started | process_completed | error | completed
7 runId: string
8 runName?: string
9 metadata?: { workflow?: { commandLine?: string; launchDir?: string; resume?: boolean; success?: boolean } }
10 trace?: { process: string; name?: string; status: string; exit?: number }
11}
12
13const blank = (name: string): Proc => ({ name, submitted: 0, completed: 0, failed: 0, cached: 0 })
14
15function bump(run: Run, name: string, change: (p: Proc) => Proc): Run {
16 const at = run.procs.findIndex(p => p.name === name)
17 const procs = [...run.procs]
18 if (at < 0) procs.push(change(blank(name)))
19 else procs[at] = change(procs[at]!)
20 return { ...run, procs }
21}
22
23/** Folds one weblog event into the run it belongs to; `started` begins a new run. */
24export function applyEvent(run: Run | null, e: WeblogEvent): Run | null {
25 const wf = e.metadata?.workflow
26 if (e.event === 'started') {
27 return { id: e.runId, name: e.runName ?? '', dir: wf?.launchDir ?? '', command: wf?.commandLine ?? '', isResume: wf?.resume === true, status: 'running', procs: [] }
28 }
29 if (!run || run.id !== e.runId) return run
30 const process = e.trace?.process
31 if (e.event === 'process_submitted' && process) return bump(run, process, p => ({ ...p, submitted: p.submitted + 1 }))
32 if (e.event === 'process_completed' && process) {
33 const isOk = e.trace!.status === 'COMPLETED'
34 return bump(run, process, p => (isOk ? { ...p, completed: p.completed + 1 } : { ...p, failed: p.failed + 1 }))
35 }
36 // A failure sends `error` then `completed`: the first one ends the run, so it toasts once.
37 if ((e.event === 'error' || e.event === 'completed') && run.status !== 'running') return run
38 if (e.event === 'error') return { ...run, status: 'failed' }
39 if (e.event === 'completed') return { ...run, status: wf?.success === false ? 'failed' : 'done' }
40 return run
41}
42
43/** The weblog sends nothing for tasks a `-resume` skips; their "Cached process" log lines fill them in. */
44export function applyCached(run: Run, logText: string): Run {
45 const counts = new Map<string, number>()
46 for (const m of logText.matchAll(/Cached process > (.+?)(?:\s+\(.*\))?\s*$/gm)) counts.set(m[1]!, (counts.get(m[1]!) ?? 0) + 1)
47 let next = run
48 for (const [name, cached] of counts) next = bump(next, name, p => ({ ...p, cached }))
49 return next
50}
51
52/** The command with what's inside quotes and heredoc bodies replaced by spaces (newlines kept), same length. */
53function textBlanked(command: string): string {
54 const spaces = (s: string) => s.replace(/[^\n]/g, ' ')
55 return command
56 .replace(/(<<-?\s*(['"]?)(\w+)\2[^\n]*\n)([\s\S]*?\n)(\s*\3)(?=\n|$)/g,
57 (_, head: string, _q, _tag, body: string, end: string) => head + spaces(body) + end)
58 .replace(/'[^']*'|"(?:\\.|[^"\\])*"/g, spaces)
59}
60
61/** The config file the listener writes, holding `weblog.url`; its name marks it in a command line. */
62export const WEBLOG_CONFIG = 'metro-map-weblog.config'
63
64/** The Bash command with its `nextflow run` reporting to the listener (`-c <config>`), or undefined when it runs none
65 * (or already reports somewhere). `-C` drops every `-c`, so there it falls back to the deprecated `-with-weblog <url>`. */
66export function withWeblog(command: string, url: string, config: string): string | undefined {
67 if (/-with-weblog\b/.test(command) || command.includes(WEBLOG_CONFIG)) return undefined
68 // `nextflow` where a command starts (after `;`, `&&`, `|`, `(` or a newline; `time`, `nohup`, `env`, `VAR=…`
69 // prefixes; any path), searched in a copy with quoted text and heredoc bodies blanked, so the same index fits both.
70 // ponytail: no `bash -c "…"` (quoted, so skipped) and no escaped quotes outside strings
71 const run = textBlanked(command).match(
72 /((?:^|[;&|(\n])\s*(?:(?:\w+=\S*|time|nohup|env|exec)\s+)*(?:[^\s;&|()]*\/)?nextflow)((?:\s+-\S+(?:\s+[^-\s]\S*)?)*?)\s+run\b/)
73 if (!run) return undefined
74 const end = run.index! + run[0].length
75 if (/\s-C\s/.test(run[2]!)) return `${command.slice(0, end)} -with-weblog ${url}${command.slice(end)}`
76 const bin = run.index! + run[1]!.length
77 return `${command.slice(0, bin)} -c ${config}${command.slice(bin)}`
78}
79
80/** Task counts over processes: cached tasks count as done. */
81export const tally = (procs: Proc[]) => {
82 let done = 0, total = 0, failed = 0, running = 0
83 for (const p of procs) {
84 done += p.completed + p.cached
85 total += p.submitted + p.cached
86 failed += p.failed
87 running += Math.max(0, p.submitted - p.completed - p.failed)
88 }
89 return { done, total, failed, running }
90}
91
92/** How long a run took, Nextflow-style: `45s`, `23m 10s`, `2h 5m`. */
93export const took = (ms: number) => {
94 const s = Math.round(ms / 1000), h = Math.floor(s / 3600), m = Math.floor((s % 3600) / 60)
95 return h ? `${h}h ${m}m` : m ? `${m}m ${s % 60}s` : `${s}s`
96}
97
98/** The end-of-run toast: `genomeqc done in 23m 10s · 52 succeeded · 49 cached`, or `genomeqc failed after 4m 2s at BUSCO_BUSCO +1 more`. */
99export function endToast(run: Run, pipeline: string, ms?: number): string {
100 if (run.status === 'failed') {
101 const [first, ...rest] = run.procs.filter(p => p.failed > 0).map(p => p.name.split(':').pop()!)
102 return `${pipeline} failed${ms === undefined ? '' : ` after ${took(ms)}`}` +
103 (first ? ` at ${first}${rest.length ? ` +${rest.length} more` : ''}` : '')
104 }
105 const done = run.procs.reduce((n, p) => n + p.completed, 0), cached = run.procs.reduce((n, p) => n + p.cached, 0)
106 return `${pipeline} done${ms === undefined ? '' : ` in ${took(ms)}`} · ${done} succeeded${cached ? ` · ${cached} cached` : ''}`
107}
108
109/** A process's station state from its task counts. */
110export const statusOf = (p: Proc | undefined): Status => {
111 if (!p) return 'pending'
112 const c = tally([p])
113 if (c.running > 0) return 'running'
114 if (c.failed > 0 && c.done === 0) return 'failed'
115 return c.total > 0 && c.done + c.failed === c.total ? 'done' : 'pending' // failed then retried ok: done
116}
117
118/** A run's command line (the weblog's `commandLine`) turned into a no-task preview that writes the DAG. */
119export function previewArgv(line: string, tmp: string): string[] | undefined {
120 // ponytail: whitespace split, quotes stripped per word; an argument with spaces in it breaks the preview (the live view still works)
121 const argv = line.trim().split(/\s+/).map(a => a.replace(/^'(.*)'$|^"(.*)"$/, '$1$2'))
122 const runAt = argv.indexOf('run')
123 if (runAt < 0) return undefined
124 const globals: string[] = []
125 for (let i = 1; i < runAt; i++) {
126 if (argv[i] === '-log' || argv[i] === '-q' || argv[i] === '-quiet') { if (argv[i] === '-log') i++; continue }
127 if (argv[i] === '-c' && argv[i + 1]?.endsWith(WEBLOG_CONFIG)) { i++; continue } // else the preview would report to our own listener
128 globals.push(argv[i]!)
129 }
130 const keep: string[] = []
131 for (let i = runAt + 1; i < argv.length; i++) {
132 const a = argv[i]!
133 if (a === '-bg' || a === '-preview') continue // -preview: we add our own, and Nextflow refuses two
134 if (a === '-with-weblog') { i++; continue } // else the preview would report to our own listener
135 // Nextflow only takes `last` or a session UUID as -resume's value; any other next word is its own argument.
136 if (a === '-resume') { if (/^(last|[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})$/.test(argv[i + 1] ?? '')) i++; continue }
137 keep.push(a)
138 }
139 return [argv[0]!, '-q', '-log', `${tmp}/preview.log`, ...globals, 'run', ...keep,
140 '-preview', '-c', `${tmp}/dag.config`, '--outdir', `${tmp}/out`]
141}
142
143export const dagConfig = (tmp: string) =>
144 `dag { enabled = true; file = '${tmp}/dag.dot'; overwrite = true }\n` +
145 `trace.enabled = false\nreport.enabled = false\ntimeline.enabled = false\n`
146types/index.d.ts 13 lines1export type Proc = { name: string; submitted: number; completed: number; failed: number; cached: number }
2export type RunStatus = 'waiting' | 'running' | 'done' | 'failed'
3/** A Nextflow run as its weblog events describe it. */
4export type Run = { id: string; name: string; dir: string; command: string; isResume: boolean; status: RunStatus; procs: Proc[]; startedAt?: number }
5/** The pipeline's DAG from `nextflow -preview -with-dag`: process names and process->process edges. */
6export type Dag = { runId: string; nodes: string[]; edges: [string, string][]; error?: string }
7
8declare module 'claude-code' {
9 interface PluginState {
10 'metro-map-mod': { run: Run | null; dag: Dag | null; port: number | null; pan: number }
11 }
12}
13