diff --git a/CHANGELOG.md b/CHANGELOG.md index 7eae47af..eb1c3053 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,19 @@ # aicodeman +## 1.1.11 + +### Patch Changes + +- Ultracode (Workflow-tool) run visualization — much better live tracking. + + While a run is in flight, the watcher previously showed empty agent slots ("agent N", 0 tokens, raw `wf_…` id as the title) because the detailed completion JSON only lands when the run finishes. The live path now enriches in-flight runs directly from the on-disk transcript tree: + - **Real per-agent stats mid-run** — tokens and tool-call counts are parsed from each `agent-.jsonl` transcript (tool counts match the final accounting exactly; token totals land within ~1% of the completion value), with model and a prompt preview. All mtime-cached (transcripts, journal, and script meta) so idle polls do no extra reads. + - **Readable window/run title** — workflow name, summary, and phases are derived from the persisted `workflows/scripts/-.js` instead of showing the raw run id. + - **Agent status colors** — done agents show green, working agents show yellow (this also fixes the run/agent status badges, which referenced undefined `--success`/`--warning` CSS variables and were rendering with no color). + - **Connector line** — the floating-window → session-tab line now uses the session-tab accent blue (was purple). + - **Click a run to open its floating window** — clicking a workflow in the dock panel opens (or focuses) its floating window with the connector line, in addition to the auto-popped windows. + - Agents are ordered by journal launch order; concurrent run-detail fetches are de-duplicated. + ## 1.1.10 ### Patch Changes diff --git a/CLAUDE.md b/CLAUDE.md index da69ac5d..22efb009 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -56,7 +56,7 @@ When user says "COM": CI runs `npm run check:lockfile` on every push/PR, so lockfile drift fails the build even if the `version-packages` script is bypassed. -**Version**: 1.1.10 (must match `package.json`) +**Version**: 1.1.11 (must match `package.json`) ## Project Overview diff --git a/package-lock.json b/package-lock.json index d28a1007..f56847ad 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "aicodeman", - "version": "1.1.10", + "version": "1.1.11", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "aicodeman", - "version": "1.1.10", + "version": "1.1.11", "hasInstallScript": true, "license": "MIT", "workspaces": [ diff --git a/package.json b/package.json index a937d0fe..c935ab48 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "aicodeman", - "version": "1.1.10", + "version": "1.1.11", "description": "Mission control for AI coding agents - run 20 autonomous agents with real-time monitoring and session persistence", "type": "module", "main": "dist/index.js", diff --git a/src/web/public/styles.css b/src/web/public/styles.css index c9754e42..998cab16 100644 --- a/src/web/public/styles.css +++ b/src/web/public/styles.css @@ -8473,14 +8473,14 @@ kbd { flex-shrink: 0; } .ultracode-status.completed { - background: var(--success); + background: var(--green); } .ultracode-status.active { - background: var(--warning); + background: var(--yellow); color: black; } .ultracode-status.failed { - background: var(--error, #b3261e); + background: var(--red, #b3261e); color: white; } @@ -8531,6 +8531,16 @@ kbd { border-radius: 6px; padding: 0.4rem 0.5rem; margin-bottom: 0.35rem; + border-left: 3px solid transparent; +} +/* At-a-glance agent state: yellow while working, green when done. */ +.ultracode-agent-card.uw-state-working { + border-left-color: var(--yellow); + background: color-mix(in srgb, var(--yellow) 9%, var(--bg-input)); +} +.ultracode-agent-card.uw-state-done { + border-left-color: var(--green); + background: color-mix(in srgb, var(--green) 9%, var(--bg-input)); } .ultracode-agent-card--clickable { cursor: pointer; @@ -8539,6 +8549,14 @@ kbd { .ultracode-agent-card--clickable:hover { background: var(--bg-card); } +/* Keep the state tint on hover (equal specificity to the hover rule, so it must + come after it) — a tad stronger so the hover still reads as interactive. */ +.ultracode-agent-card--clickable.uw-state-working:hover { + background: color-mix(in srgb, var(--yellow) 16%, var(--bg-card)); +} +.ultracode-agent-card--clickable.uw-state-done:hover { + background: color-mix(in srgb, var(--green) 16%, var(--bg-card)); +} .ultracode-agent-top { display: flex; align-items: center; @@ -8564,10 +8582,10 @@ kbd { color: white; } .ultracode-agent-state.completed { - background: var(--success); + background: var(--green); } .ultracode-agent-state.active { - background: var(--warning); + background: var(--yellow); color: black; } .ultracode-agent-state.idle { @@ -8707,14 +8725,15 @@ kbd { margin-bottom: 0.4rem; } -/* Ultracode connector line — distinct purple to set it apart from blue subagent lines. */ +/* Ultracode connector line — the session-tab accent blue (dashed to stay distinct + from the solid subagent lines, per the same-blue-as-the-tab request). */ .connection-line.ultracode-connection { - stroke: #a855f7; + stroke: var(--accent); stroke-width: 3; stroke-dasharray: 6 3; filter: drop-shadow(0 0 2px rgba(0, 0, 0, 0.8)) - drop-shadow(0 0 5px rgba(168, 85, 247, 0.85)) - drop-shadow(0 0 10px rgba(168, 85, 247, 0.5)); + drop-shadow(0 0 5px rgba(var(--accent-rgb), 0.85)) + drop-shadow(0 0 10px rgba(var(--accent-rgb), 0.5)); animation: ultracode-conn-pulse 1.4s ease-in-out infinite; } .connection-line.ultracode-connection:hover { diff --git a/src/web/public/ultracode-panel.js b/src/web/public/ultracode-panel.js index 4297b4b4..ca7293f6 100644 --- a/src/web/public/ultracode-panel.js +++ b/src/web/public/ultracode-panel.js @@ -96,6 +96,10 @@ Object.assign(CodemanApp.prototype, { this.activeWorkflowPhaseIndex = null; // reset phase filter on run change this._fetchWorkflowRunDetail(runId); this.renderUltracodeAgentsPanel(); + // Clicking a run also pops its floating window (with connector line to the + // session tab), like the auto-popped one — an explicit open, so it ignores the + // floating-windows auto-pop toggle (ultracode-windows.js). + if (typeof this.openUltracodeWindowForRun === 'function') this.openUltracodeWindowForRun(runId); }, selectWorkflowPhase(phaseIndex) { this._ensureWorkflowState(); @@ -105,6 +109,11 @@ Object.assign(CodemanApp.prototype, { }, async _fetchWorkflowRunDetail(runId) { + // De-dupe concurrent fetches for the same run — selecting a run can trigger both + // a panel refresh and a floating-window open, which would otherwise double-fetch. + if (!this._wfDetailInFlight) this._wfDetailInFlight = new Set(); + if (this._wfDetailInFlight.has(runId)) return; + this._wfDetailInFlight.add(runId); try { const res = await fetch(`/api/workflows/${encodeURIComponent(runId)}`); const env = await res.json(); @@ -117,6 +126,8 @@ Object.assign(CodemanApp.prototype, { } } catch { /* transient — next update retries */ + } finally { + this._wfDetailInFlight.delete(runId); } }, @@ -274,10 +285,12 @@ Object.assign(CodemanApp.prototype, { // to the agent-.jsonl stem already tracked by subagent-watcher). 'start' agents have // no agentId yet, so they stay non-clickable. const clickable = !!a.agentId; + // At-a-glance state tint on the whole card: green when done, yellow while working. + const cardStateCls = state === 'done' ? ' uw-state-done' : state === 'progress' ? ' uw-state-working' : ''; const cardAttrs = clickable - ? ` class="ultracode-agent-card ultracode-agent-card--clickable" role="button" tabindex="0"` + + ? ` class="ultracode-agent-card ultracode-agent-card--clickable${cardStateCls}" role="button" tabindex="0"` + ` title="View transcript" onclick="app.openWorkflowAgentTranscript('${escapeHtml(a.agentId)}')"` - : ` class="ultracode-agent-card"`; + : ` class="ultracode-agent-card${cardStateCls}"`; return ( `` + `
` + diff --git a/src/web/public/ultracode-windows.js b/src/web/public/ultracode-windows.js index 3537e8ac..3901c659 100644 --- a/src/web/public/ultracode-windows.js +++ b/src/web/public/ultracode-windows.js @@ -77,12 +77,14 @@ Object.assign(CodemanApp.prototype, { _syncUltracodeFloatingWindow(run, opts) { this._ensureUltracodeWindowState(); if (!run || !run.runId) return; - if (!this._ultracodeFloatingEnabled()) return; const runId = run.runId; + const existing = this.ultracodeWindows.get(runId); + // Auto-pop is gated on the floating-windows toggle, but an ALREADY-open window + // (e.g. one opened by clicking the run in the dock) keeps refreshing regardless. + if (!existing && !this._ultracodeFloatingEnabled()) return; if (this.ultracodeWindowsClosed.has(runId)) return; // respect explicit dismissal const active = this._isWorkflowRunActive(run); - const existing = this.ultracodeWindows.get(runId); if (active) { // Run is alive — cancel any pending auto-close. @@ -120,6 +122,31 @@ Object.assign(CodemanApp.prototype, { } }, + /** + * Explicitly open (or focus) the floating window for a run — the click-through + * from the dock panel's run list. Unlike auto-pop this ignores the floating-windows + * toggle and clears any prior dismissal (it's a direct user action), then draws the + * connector line from the run's session tab. + */ + openUltracodeWindowForRun(runId) { + this._ensureUltracodeWindowState(); + if (!runId) return; + const run = this.workflowRuns && this.workflowRuns.get(runId); + if (!run) return; + this.ultracodeWindowsClosed.delete(runId); // an explicit open overrides a past dismissal + const existing = this.ultracodeWindows.get(runId); + if (existing) { + // Already open — bring to front, expand if collapsed, refresh. + existing.element.style.zIndex = ++this.ultracodeWindowZIndex; + if (existing.collapsed) this.toggleUltracodeWindowCollapse(runId); + this.renderUltracodeWindowContent(runId); + this._fetchWorkflowRunDetail(runId); + this.updateConnectionLines(); + } else { + this.createUltracodeWindow(run); + } + }, + /** Build and mount a floating window for a run, positioned near its parent tab. */ createUltracodeWindow(run) { this._ensureUltracodeWindowState(); diff --git a/src/workflow-run-watcher.ts b/src/workflow-run-watcher.ts index 98a7a922..ee7c9144 100644 --- a/src/workflow-run-watcher.ts +++ b/src/workflow-run-watcher.ts @@ -53,16 +53,99 @@ const RUN_FILE_PREFIX = 'wf_'; const RUN_FILE_SUFFIX = '.json'; const LIVE_JOURNAL_FILE = 'journal.jsonl'; const LIVE_AGENT_PREFIX = 'agent-'; +const LIVE_TRANSCRIPT_SUFFIX = '.jsonl'; +const SCRIPTS_SUBDIR = 'scripts'; /** Hard caps on the largest per-agent strings so a 28-agent run stays compact. */ const PROMPT_PREVIEW_MAX = 200; const RESULT_PREVIEW_MAX = 240; +/** Cap a live script read so a runaway/huge embedded script can't blow memory. */ +const SCRIPT_READ_MAX_BYTES = 512 * 1024; +/** LRU cap on cached per-agent transcript stats across all live runs. */ +const MAX_CACHED_AGENT_STATS = 2000; function truncate(value: string | undefined, max: number): string | undefined { if (typeof value !== 'string') return undefined; return value.length > max ? `${value.slice(0, max)}…` : value; } +/** First text from a transcript message's `content` (string or content-block array). */ +function extractMessageText(content: unknown): string | undefined { + if (typeof content === 'string') return content.trim() || undefined; + if (Array.isArray(content)) { + for (const block of content) { + if (block && typeof block === 'object' && (block as { type?: string }).type === 'text') { + const t = (block as { text?: string }).text; + if (typeof t === 'string' && t.trim()) return t.trim(); + } + } + } + return undefined; +} + +/** + * Isolate the `meta = { … }` object literal from a workflow script so that later + * declarations (e.g. JSON-Schema objects with their own `description:`/`title:` + * fields, or prose containing `phases: [`) can't be mistaken for meta. Brace-matches + * naively — a `{`/`}` inside a string value can still confuse it, in which case we + * return undefined and the caller falls back to scanning the whole body. Best-effort. + */ +function extractMetaBlock(body: string): string | undefined { + const m = /(?:^|[^A-Za-z0-9_])(?:export\s+const|const|let|var)\s+meta\s*=\s*\{/.exec(body); + if (!m) return undefined; + const open = body.indexOf('{', m.index); + if (open < 0) return undefined; + let depth = 0; + for (let i = open; i < body.length; i++) { + if (body[i] === '{') depth++; + else if (body[i] === '}' && --depth === 0) return body.slice(open, i + 1); + } + return undefined; // unterminated (truncated script) — caller falls back to whole body +} + +/** + * Best-effort extraction of a single/double-quoted `key: '...'` value from a workflow + * `meta` literal (the script is JS we must not eval). The key is anchored to a key + * position (`{`/`,`/whitespace before it) so e.g. `description` can't match the tail + * of `..._description`; pass the scoped meta block so later same-named schema fields + * can't win. Returns undefined on no match. + */ +function extractMetaString(body: string, key: string): string | undefined { + const re = new RegExp(`(?:^|[{,\\s])${key}\\s*:\\s*(['"])((?:\\\\.|(?!\\1).)*)\\1`); + const m = re.exec(body); + return m ? m[2].replace(/\\(['"\\`])/g, '$1') : undefined; +} + +/** Best-effort `meta.phases: [{title, detail}, …]` extraction (pass the scoped meta block). */ +function extractMetaPhases(body: string): WorkflowRunPhase[] { + const keyMatch = /(?:^|[^A-Za-z0-9_])phases\s*:\s*\[/.exec(body); + if (!keyMatch) return []; + const open = body.indexOf('[', keyMatch.index); + if (open < 0) return []; + // Brace-match to the array's closing ] so nested arrays don't end it early. + let depth = 0; + let end = -1; + for (let i = open; i < body.length; i++) { + if (body[i] === '[') depth++; + else if (body[i] === ']' && --depth === 0) { + end = i; + break; + } + } + if (end < 0) return []; + const block = body.slice(open, end + 1); + const phases: WorkflowRunPhase[] = []; + const re = /title\s*:\s*(['"])((?:\\.|(?!\1).)*)\1(?:\s*,\s*detail\s*:\s*(['"])((?:\\.|(?!\3).)*)\3)?/g; + let m: RegExpExecArray | null; + while ((m = re.exec(block)) !== null && phases.length < 20) { + phases.push({ + title: m[2].replace(/\\(['"\\`])/g, '$1'), + detail: (m[4] || '').replace(/\\(['"\\`])/g, '$1'), + }); + } + return phases; +} + /** Drop the heavy `agents[]` for list/snapshot use. */ export function summarizeRun(info: WorkflowRunInfo): WorkflowRunSummary { const { agents: _agents, ...summary } = info; @@ -85,6 +168,34 @@ interface DiscoveredLiveRun { runId: string; } +/** Per-agent stats parsed from one `agent-.jsonl` transcript (mtime-cached). */ +interface AgentTranscriptStats { + /** Tokens for the agent's LAST usage-bearing message (in+out+cache) ≈ the completion-JSON `tokens`. */ + tokens: number; + /** Count of `tool_use` blocks across the transcript (matches the completion-JSON `toolCalls`). */ + toolCalls: number; + /** Most-recent model id seen. */ + model: string; + /** Name of the most-recent `tool_use` block. */ + lastToolName?: string; + /** First user-message text, truncated — a hint of what the agent was asked. */ + promptPreview?: string; +} + +/** Workflow meta derived live from `workflows/scripts/-.js`. */ +interface LiveWorkflowMeta { + workflowName?: string; + summary?: string; + phases: WorkflowRunPhase[]; +} + +/** A live agent file on disk (transcript and/or meta), keyed by its `` stem. */ +interface LiveAgentFile { + agentId: string; + transcriptPath?: string; + transcriptMtime?: number; +} + /** A `workflowProgress[]` entry as it appears on disk (loosely typed for defensive parsing). */ interface RawProgressEntry { type?: string; @@ -123,6 +234,17 @@ export class WorkflowRunWatcher extends EventEmitter { private liveDirMtimes = new Map(); /** runId -> absolute live transcript-dir path (for mtime cleanup on removal). */ private runIdToLiveDir = new Map(); + /** transcript abs path -> { mtimeMs, stats }; re-parse a transcript only when its mtime moves. */ + private agentStatCache = new LRUMap({ + maxSize: MAX_CACHED_AGENT_STATS, + }); + /** runId -> derived script meta (immutable per run; re-derived only until a name is found). */ + private liveMetaCache = new Map(); + /** journal abs path -> { mtimeMs, parsed }; re-parse the journal only when it grows. */ + private journalCache = new Map< + string, + { mtimeMs: number; parsed: { startedOrder: string[]; doneIds: Set } } + >(); /** watched-dir absolute path -> chokidar watcher (workflows/ + subagents/workflows/). */ private dirWatchers = new Map(); @@ -160,6 +282,9 @@ export class WorkflowRunWatcher extends EventEmitter { this.runIdToPath.clear(); this.liveDirMtimes.clear(); this.runIdToLiveDir.clear(); + this.agentStatCache.clear(); + this.liveMetaCache.clear(); + this.journalCache.clear(); } /** All cached runs (no recency filter), most-recently-active first. */ @@ -243,7 +368,18 @@ export class WorkflowRunWatcher extends EventEmitter { if (path) this.fileMtimes.delete(path); this.runIdToPath.delete(runId); const liveDir = this.runIdToLiveDir.get(runId); - if (liveDir) this.liveDirMtimes.delete(liveDir); + if (liveDir) { + this.liveDirMtimes.delete(liveDir); + // Drop cached transcript stats + journal parse under this run's live dir. + const prefix = liveDir.endsWith('/') ? liveDir : liveDir + '/'; + for (const key of Array.from(this.agentStatCache.keys())) { + if (key.startsWith(prefix)) this.agentStatCache.delete(key); + } + for (const key of Array.from(this.journalCache.keys())) { + if (key.startsWith(prefix)) this.journalCache.delete(key); + } + } + this.liveMetaCache.delete(runId); this.runIdToLiveDir.delete(runId); this.emit('run_removed', { runId }); } @@ -346,14 +482,58 @@ export class WorkflowRunWatcher extends EventEmitter { } /** - * Re-synthesize an in-flight run from its transcript dir when its newest member - * mtime moved (skip otherwise so we don't re-emit run_updated on idle polls). + * Re-synthesize an in-flight run from its transcript dir. Cheap pass first: stat + * the dir members and skip entirely when the newest mtime is unchanged, so an idle + * poll never reads a single (large) transcript. Only on a real change do we parse — + * and even then per-transcript stats come from an mtime-keyed cache, so only the + * transcripts that actually grew are re-read. */ private async maybeParseLive(live: DiscoveredLiveRun): Promise { - const info = await this.parseLiveDir(live); + let entries: string[]; + try { + entries = await readdir(live.dirPath); + } catch { + return; // vanished between discover and read + } + + // Cheap pass: newest member mtime + the agent-file set (no transcript reads yet). + let newestMtime = 0; + let journalPath: string | null = null; + let journalMtime = 0; + const agentFiles = new Map(); + for (const name of entries) { + if (name !== LIVE_JOURNAL_FILE && !name.startsWith(LIVE_AGENT_PREFIX)) continue; + const full = join(live.dirPath, name); + let m = 0; + try { + m = (await stat(full)).mtimeMs; + } catch { + continue; // entry vanished — ignore + } + if (m > newestMtime) newestMtime = m; + if (name === LIVE_JOURNAL_FILE) { + journalPath = full; + journalMtime = m; + continue; + } + // agent-.jsonl (transcript) or agent-.meta.json (queued-slot marker) + const stem = name.slice(LIVE_AGENT_PREFIX.length).replace(/\.(meta\.json|jsonl)$/, ''); + if (!stem) continue; + const slot = agentFiles.get(stem) || { agentId: stem }; + if (name.endsWith(LIVE_TRANSCRIPT_SUFFIX)) { + slot.transcriptPath = full; + slot.transcriptMtime = m; + } + agentFiles.set(stem, slot); + } + if (agentFiles.size === 0) return; // nothing to show yet + + // Skip the expensive parse when nothing changed since the last synthesis. + if (this.liveDirMtimes.get(live.dirPath) === newestMtime) return; + this.liveDirMtimes.set(live.dirPath, newestMtime); + + const info = await this.parseLiveDir(live, agentFiles, journalPath, journalMtime, newestMtime); if (!info) return; - if (this.liveDirMtimes.get(live.dirPath) === info.lastActivityAt) return; - this.liveDirMtimes.set(live.dirPath, info.lastActivityAt); const existed = this.runs.has(info.runId); this.runs.set(info.runId, info); @@ -362,56 +542,87 @@ export class WorkflowRunWatcher extends EventEmitter { } /** - * Build a minimal ACTIVE WorkflowRunInfo from `subagents/workflows/wf_/`. - * The transcript tree carries no phases/tokens — those arrive with the - * completion wf_*.json — so we expose: the agent slots (keyed by agentId, so the - * card→transcript click still works), each marked done/running from journal - * `result` lines, and lastActivityAt from the newest agent/journal mtime. + * Build an ACTIVE WorkflowRunInfo from a live transcript dir, ENRICHED so the + * floating window / panel show real data mid-run rather than empty slots: + * - per-agent tokens + tool-calls + model + last tool, parsed from each + * `agent-.jsonl` (mtime-cached — only changed transcripts are re-read); + * - per-agent state: 'done' once the journal logs a `result` (→ green badge), + * 'progress' once it logs `started` or a transcript exists (→ yellow), else 'start'; + * - run name/summary/phases from the live `workflows/scripts/-.js` + * (the completion wf_.json, which carries these, only lands at the end); + * - run totals = sums of the per-agent stats. */ - private async parseLiveDir(live: DiscoveredLiveRun): Promise { - let entries: string[]; - try { - entries = await readdir(live.dirPath); - } catch { - return null; // vanished between discover and read - } + private async parseLiveDir( + live: DiscoveredLiveRun, + agentFiles: Map, + journalPath: string | null, + journalMtime: number, + newestMtime: number + ): Promise { + const { startedOrder, doneIds } = journalPath + ? await this.readJournal(journalPath, journalMtime) + : { startedOrder: [] as string[], doneIds: new Set() }; + const startIndex = new Map(); + startedOrder.forEach((id, i) => startIndex.set(id, i)); - const agentIds = new Set(); - let newestMtime = 0; - for (const name of entries) { - if (name.startsWith(LIVE_AGENT_PREFIX)) { - const stem = name.slice(LIVE_AGENT_PREFIX.length).replace(/\.(meta\.json|jsonl)$/, ''); - if (stem) agentIds.add(stem); - } - if (name === LIVE_JOURNAL_FILE || name.startsWith(LIVE_AGENT_PREFIX)) { - try { - const m = (await stat(join(live.dirPath, name))).mtimeMs; - if (m > newestMtime) newestMtime = m; - } catch { - // entry vanished — ignore - } - } - } - if (agentIds.size === 0) return null; // nothing to show yet + // Order agents by journal launch order; not-yet-started ones trail (sorted by id). + const ids = Array.from(agentFiles.keys()).sort((a, b) => { + const ia = startIndex.has(a) ? (startIndex.get(a) as number) : Number.MAX_SAFE_INTEGER; + const ib = startIndex.has(b) ? (startIndex.get(b) as number) : Number.MAX_SAFE_INTEGER; + if (ia !== ib) return ia - ib; + return a < b ? -1 : a > b ? 1 : 0; + }); - const doneIds = await this.readJournalDoneAgents(join(live.dirPath, LIVE_JOURNAL_FILE)); - const agents: WorkflowAgentInfo[] = Array.from(agentIds) - .sort() - .map((id, i) => ({ + let totalTokens = 0; + let totalToolCalls = 0; + let defaultModel = ''; + const agents: WorkflowAgentInfo[] = []; + for (let i = 0; i < ids.length; i++) { + const id = ids[i]; + const file = agentFiles.get(id) as LiveAgentFile; + const stats = file.transcriptPath + ? await this.agentTranscriptStats(file.transcriptPath, file.transcriptMtime || 0) + : null; + const state = doneIds.has(id) ? 'done' : startIndex.has(id) || file.transcriptPath ? 'progress' : 'start'; + if (stats) { + totalTokens += stats.tokens; + totalToolCalls += stats.toolCalls; + if (stats.model && !defaultModel) defaultModel = stats.model; + } + agents.push({ index: i + 1, label: `agent ${i + 1}`, phaseIndex: 1, phaseTitle: '', - model: '', - state: doneIds.has(id) ? 'done' : 'progress', + model: stats?.model || '', + state, agentId: id, - })); + tokens: stats ? stats.tokens : undefined, + toolCalls: stats ? stats.toolCalls : undefined, + lastToolName: stats?.lastToolName, + promptPreview: stats?.promptPreview, + }); + } + + // Script meta is immutable for the life of a run, so derive it at most once. + // Re-derive only while we still lack a name (the script can land a poll or two + // after the first transcripts, which seed the run before the script exists). + let meta = this.liveMetaCache.get(live.runId); + if (!meta || !meta.workflowName) { + meta = await this.deriveLiveWorkflowMeta(live); + this.liveMetaCache.set(live.runId, meta); + } return { runId: live.runId, + workflowName: meta.workflowName, + summary: meta.summary, status: 'running', agentCount: agents.length, - phases: [], + totalTokens, + totalToolCalls, + defaultModel: defaultModel || undefined, + phases: meta.phases, agents, sessionUuid: live.sessionUuid, projectHash: live.projectHash, @@ -419,25 +630,145 @@ export class WorkflowRunWatcher extends EventEmitter { }; } - /** Agent ids that already emitted a `result` event in the run journal. */ - private async readJournalDoneAgents(journalPath: string): Promise> { - const done = new Set(); + /** Parse one agent transcript for token/tool stats; cached on the file's mtime. */ + private async agentTranscriptStats(path: string, mtimeMs: number): Promise { + // peek (not get) on the hit-check so a cache HIT doesn't churn the LRU — every + // active agent is looked up each poll, and get() would delete+reinsert each time. + const cached = this.agentStatCache.peek(path); + if (cached && cached.mtimeMs === mtimeMs) return cached.stats; + let text: string; + try { + text = await readFile(path, 'utf-8'); + } catch { + return null; // vanished mid-poll + } + let tokens = 0; + let toolCalls = 0; + let model = ''; + let lastToolName: string | undefined; + let promptPreview: string | undefined; + for (const line of text.split('\n')) { + if (!line) continue; + let entry: { type?: string; message?: unknown }; + try { + entry = JSON.parse(line) as { type?: string; message?: unknown }; + } catch { + continue; // tolerate a partially-written trailing line + } + const msg = entry.message as + | { + model?: string; + usage?: { + input_tokens?: number; + output_tokens?: number; + cache_read_input_tokens?: number; + cache_creation_input_tokens?: number; + }; + content?: Array<{ type?: string; name?: string; text?: string }> | string; + } + | undefined; + if (!promptPreview && entry.type === 'user' && msg) { + promptPreview = truncate(extractMessageText(msg.content), PROMPT_PREVIEW_MAX); + } + if (!msg || typeof msg !== 'object') continue; + if (typeof msg.model === 'string' && msg.model) model = msg.model; + const u = msg.usage; + if (u) { + const total = + (u.input_tokens || 0) + + (u.output_tokens || 0) + + (u.cache_read_input_tokens || 0) + + (u.cache_creation_input_tokens || 0); + // Anthropic usage already reflects cumulative context, so the agent's running + // token total ≈ its LATEST usage-bearing message — last non-zero wins. + if (total > 0) tokens = total; + } + if (Array.isArray(msg.content)) { + for (const block of msg.content) { + if (block && block.type === 'tool_use' && typeof block.name === 'string') { + toolCalls++; + lastToolName = block.name; + } + } + } + } + const stats: AgentTranscriptStats = { tokens, toolCalls, model, lastToolName, promptPreview }; + this.agentStatCache.set(path, { mtimeMs, stats }); + return stats; + } + + /** + * Read the run journal: agent ids in `started` order plus the set that already + * logged a terminal `result`. (`started` lines appear in launch order.) + */ + private async readJournal( + journalPath: string, + mtimeMs: number + ): Promise<{ startedOrder: string[]; doneIds: Set }> { + const cached = this.journalCache.get(journalPath); + if (cached && cached.mtimeMs === mtimeMs) return cached.parsed; + const startedOrder: string[] = []; + const seenStarted = new Set(); + const doneIds = new Set(); let text: string; try { text = await readFile(journalPath, 'utf-8'); } catch { - return done; // journal not written yet — all agents still in progress + return { startedOrder, doneIds }; // journal not written yet } for (const line of text.split('\n')) { if (!line) continue; try { const ev = JSON.parse(line) as { type?: string; agentId?: string }; - if (ev && ev.type === 'result' && typeof ev.agentId === 'string') done.add(ev.agentId); + if (!ev || typeof ev.agentId !== 'string') continue; + if (ev.type === 'started' && !seenStarted.has(ev.agentId)) { + seenStarted.add(ev.agentId); + startedOrder.push(ev.agentId); + } else if (ev.type === 'result') { + doneIds.add(ev.agentId); + } } catch { // tolerate a partially-written trailing line } } - return done; + const parsed = { startedOrder, doneIds }; + this.journalCache.set(journalPath, { mtimeMs, parsed }); + return parsed; + } + + /** + * Derive name/summary/phases for a live run from its persisted script file + * `…/workflows/scripts/-.js`. The filename yields a reliable name; + * `meta.description`/`meta.phases` are best-effort regex extractions (the script is + * JS we must not eval), so any parse miss simply leaves those optional fields absent. + */ + private async deriveLiveWorkflowMeta(live: DiscoveredLiveRun): Promise { + const empty: LiveWorkflowMeta = { phases: [] }; + const scriptsDir = join(this.projectsDir, live.projectHash, live.sessionUuid, WORKFLOWS_SUBDIR, SCRIPTS_SUBDIR); + let names: string[]; + try { + names = await readdir(scriptsDir); + } catch { + return empty; + } + const suffix = `-${live.runId}.js`; + const fileName = names.find((n) => n.endsWith(suffix)) || names.find((n) => n.includes(live.runId)); + if (!fileName) return empty; + const nameFromFile = fileName.endsWith(suffix) ? fileName.slice(0, -suffix.length) : fileName.replace(/\.js$/, ''); + let body = ''; + try { + body = (await readFile(join(scriptsDir, fileName), 'utf-8')).slice(0, SCRIPT_READ_MAX_BYTES); + } catch { + return { workflowName: nameFromFile || undefined, phases: [] }; + } + // Scope name/summary/phases parsing to the `meta` literal so later script + // declarations (schemas, prose) can't be mistaken for it. + const metaBlock = extractMetaBlock(body) || body; + return { + workflowName: nameFromFile || extractMetaString(metaBlock, 'name') || undefined, + summary: extractMetaString(metaBlock, 'description') || undefined, + phases: extractMetaPhases(metaBlock), + }; } /** diff --git a/test/workflow-run-watcher.test.ts b/test/workflow-run-watcher.test.ts index 39f152be..73fa042b 100644 --- a/test/workflow-run-watcher.test.ts +++ b/test/workflow-run-watcher.test.ts @@ -317,3 +317,99 @@ describe('WorkflowRunWatcher — in-flight (live) runs', () => { expect(watcher.getAllRuns()).toHaveLength(1); }); }); + +/** + * Live ENRICHMENT: while a run is in-flight the watcher now parses each agent + * transcript for real tokens/tool-calls/model, maps state→colour from the journal, + * and derives the run name/summary/phases from the persisted script — so the + * floating window/panel show real data mid-run instead of "0 tok / agent N". + */ +describe('WorkflowRunWatcher — live enrichment (tokens/state/name)', () => { + const RUN = 'wf_enrich01-abc'; + let projectsDir: string; + let watcher: WorkflowRunWatcher; + + beforeEach(async () => { + projectsDir = await mkdtemp(join(tmpdir(), 'wfw-enrich-')); + const sessionDir = join(projectsDir, PROJECT_HASH, SESSION_UUID); + const liveDir = join(sessionDir, 'subagents', 'workflows', RUN); + await mkdir(liveDir, { recursive: true }); + const scriptsDir = join(sessionDir, 'workflows', 'scripts'); + await mkdir(scriptsDir, { recursive: true }); + + // Agent aaa: a user prompt + two assistant turns with usage + two tool_use blocks. + await writeFile( + join(liveDir, 'agent-aaa.jsonl'), + [ + '{"type":"user","message":{"role":"user","content":"Audit the docs"}}', + '{"type":"assistant","message":{"model":"claude-opus-4-8","usage":{"input_tokens":1000,"output_tokens":40},"content":[{"type":"tool_use","name":"Bash"}]}}', + '{"type":"assistant","message":{"model":"claude-opus-4-8","usage":{"input_tokens":2000,"cache_read_input_tokens":500,"output_tokens":80},"content":[{"type":"tool_use","name":"Edit"},{"type":"text","text":"done"}]}}', + ].join('\n') + '\n', + 'utf-8' + ); + await writeFile(join(liveDir, 'agent-aaa.meta.json'), JSON.stringify({ agentType: 'workflow-subagent' }), 'utf-8'); + // Agent bbb: started, no result yet, minimal transcript (no usage). + await writeFile(join(liveDir, 'agent-bbb.jsonl'), '{"type":"assistant","message":{"content":[]}}\n', 'utf-8'); + // bbb starts BEFORE aaa, so journal launch order ['bbb','aaa'] differs from the + // old alphabetical sort ['aaa','bbb'] — the ordering assertion below is adversarial. + await writeFile( + join(liveDir, 'journal.jsonl'), + '{"type":"started","agentId":"bbb"}\n{"type":"started","agentId":"aaa"}\n{"type":"result","agentId":"aaa","result":{}}\n', + 'utf-8' + ); + // Persisted script — name from filename, summary/phases from the meta literal. + await writeFile( + join(scriptsDir, `my-cool-workflow-${RUN}.js`), + "export const meta = {\n name: 'my-cool-workflow',\n description: 'Audit and update the docs',\n phases: [ { title: 'Plan', detail: 'plan it' }, { title: 'Do', detail: 'do it' } ],\n}\n", + 'utf-8' + ); + watcher = new WorkflowRunWatcher(projectsDir); + }); + + afterEach(async () => { + watcher.stop(); + await rm(projectsDir, { recursive: true, force: true }); + }); + + function firstRun(): Promise { + return new Promise((resolve, reject) => { + const t = setTimeout(() => reject(new Error('timed out')), 5000); + watcher.once('run_discovered', (info: WorkflowRunInfo) => { + clearTimeout(t); + resolve(info); + }); + watcher.start(); + }); + } + + it('derives workflowName/summary/phases from the live script file', async () => { + const info = await firstRun(); + expect(info.workflowName).toBe('my-cool-workflow'); + expect(info.summary).toBe('Audit and update the docs'); + expect(info.phases.map((p) => p.title)).toEqual(['Plan', 'Do']); + }); + + it('parses per-agent tokens (last usage-bearing message) + tool-call counts from the transcript', async () => { + const info = await firstRun(); + const aaa = info.agents.find((a) => a.agentId === 'aaa')!; + expect(aaa.tokens).toBe(2580); // last turn: 2000 in + 500 cache + 80 out + expect(aaa.toolCalls).toBe(2); // Bash + Edit + expect(aaa.model).toBe('claude-opus-4-8'); + expect(aaa.promptPreview).toBe('Audit the docs'); + }); + + it('maps state to done (→green) / progress (→yellow) from the journal', async () => { + const info = await firstRun(); + expect(info.agents.find((a) => a.agentId === 'aaa')!.state).toBe('done'); + expect(info.agents.find((a) => a.agentId === 'bbb')!.state).toBe('progress'); + }); + + it('orders agents by journal launch order (NOT alphabetical) and sums run-level totals', async () => { + const info = await firstRun(); + // bbb started first in the journal though it sorts after aaa — launch order wins. + expect(info.agents.map((a) => a.agentId)).toEqual(['bbb', 'aaa']); + expect(info.agents[0].label).toBe('agent 1'); // = bbb, the first-launched + expect(info.totalToolCalls).toBe(2); + expect(info.totalTokens).toBe(2580); + }); +});