diff --git a/CHANGELOG.md b/CHANGELOG.md index cd94e483..dde4701e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,15 @@ # aicodeman +## 1.1.4 + +### Patch Changes + +- Fix: ultracode floating run windows (and the live dock panel) now appear DURING an in-flight Workflow/ultracode run, not only after it finishes. + + The Workflow runtime writes the run-state file `…/workflows/wf_.json` only at completion (always a terminal status); while a run is live, its only on-disk state is the sibling `…/subagents/workflows/wf_/` transcript tree. `workflow-run-watcher` previously scanned only the completion file, so it never observed a run until it was already terminal — and the floating-window auto-pop is gated on an ACTIVE run, so it never fired for a live run (the feature was effectively dead for in-flight runs). + + The watcher now ALSO scans the `subagents/workflows/wf_/` transcript tree and synthesizes a minimal ACTIVE run (status `running`, agent slots keyed by their `agentId` so the agent-card → live-transcript click still works, `lastActivityAt` from the newest agent/journal mtime, per-agent done/running derived from the run journal's `result` events) when no completion file exists yet. When the run finishes, the real `wf_.json` supersedes the synthesized record (same runId), restoring full phase/token detail and the normal finish → 8s-grace auto-close flow. The watcher stays standalone (it never imports subagent-watcher). Verified end-to-end against a real in-flight run; adds unit coverage for live synthesis, agentId preservation, journal-derived state, empty-dir skipping, and completion-file precedence. + ## 1.1.3 ### Patch Changes diff --git a/CLAUDE.md b/CLAUDE.md index b269f477..f6cdc2ff 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.3 (must match `package.json`) +**Version**: 1.1.4 (must match `package.json`) ## Project Overview diff --git a/package-lock.json b/package-lock.json index ca46b78b..51c88de2 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "aicodeman", - "version": "1.1.3", + "version": "1.1.4", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "aicodeman", - "version": "1.1.3", + "version": "1.1.4", "hasInstallScript": true, "license": "MIT", "workspaces": [ diff --git a/package.json b/package.json index d741ab4e..b4079892 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "aicodeman", - "version": "1.1.3", + "version": "1.1.4", "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/workflow-run-watcher.ts b/src/workflow-run-watcher.ts index 3b9500b1..98a7a922 100644 --- a/src/workflow-run-watcher.ts +++ b/src/workflow-run-watcher.ts @@ -1,19 +1,33 @@ /** * @fileoverview Workflow (ultracode) Run Watcher * - * Watches `~/.claude/projects///workflows/wf_*.json` — - * the run-state JSON the Workflow tool writes for each ultracode run — and emits - * events powering the master-detail "working agents" view (tasks/phases on the - * LEFT, per-agent tokens/tool-calls on the RIGHT). + * Emits events powering the master-detail "working agents" view (tasks/phases on + * the LEFT, per-agent tokens/tool-calls on the RIGHT) AND the floating run + * windows, from TWO disk sources per run: * - * Deliberately STANDALONE: it never imports from or touches subagent-watcher.ts. - * It globs the run-state tree (`.../workflows/wf_*.json`); subagent-watcher globs - * the disjoint, deeper transcript tree (`.../subagents/workflows/wf_/agent-*`). - * Separate singletons, disjoint directories, no shared mutable state. + * 1. COMPLETION artifact — `…/workflows/wf_.json`. The Workflow runtime + * writes this with the FULL run state (phases, per-agent tokens/tool-calls, + * result), but — as of the mid-2026 runtime — only when the run FINISHES + * (always a terminal status). It is the authoritative, detailed record. + * 2. LIVE transcript dir — `…/subagents/workflows/wf_/` (agent-*.jsonl + + * journal.jsonl). This appears WHILE a run is in flight, before any + * `wf_.json` exists. From it we synthesize a minimal ACTIVE run + * (status 'running', agent slots keyed by agentId, lastActivityAt from file + * mtimes) so the floating window pops DURING the run instead of only after. + * + * Precedence: when a completion `wf_.json` exists it ALWAYS supersedes the + * synthesized live record (same runId), so a finished run shows full detail and + * the normal finish→auto-close flow runs. Without source 2 the floating-window + * feature is dead for live runs (the completion file only lands at the end, so + * the watcher would never see a run while it is active). + * + * Still STANDALONE: it never imports from or touches subagent-watcher.ts. It + * independently reads the same `subagents/workflows/` tree subagent-watcher uses, + * but as a separate singleton with no shared mutable state. * * Discovery is dual: a periodic poll (catches new run dirs + removals) plus a - * per-run-dir chokidar watcher (live within-run updates). A per-file mtime skip - * keeps the hot path cheap — the run JSON is rewritten on every agent tick. + * per-dir chokidar watcher (live updates). A per-source mtime skip keeps the hot + * path cheap — the completion JSON / live dir is re-read only when its mtime moves. * * @module workflow-run-watcher */ @@ -34,8 +48,11 @@ import { } from './config/workflow-config.js'; const WORKFLOWS_SUBDIR = 'workflows'; +const SUBAGENTS_SUBDIR = 'subagents'; const RUN_FILE_PREFIX = 'wf_'; const RUN_FILE_SUFFIX = '.json'; +const LIVE_JOURNAL_FILE = 'journal.jsonl'; +const LIVE_AGENT_PREFIX = 'agent-'; /** Hard caps on the largest per-agent strings so a 28-agent run stays compact. */ const PROMPT_PREVIEW_MAX = 200; @@ -60,6 +77,14 @@ interface DiscoveredRun { runId: string; } +/** An in-flight run discovered from its `subagents/workflows/wf_/` transcript dir. */ +interface DiscoveredLiveRun { + dirPath: string; + projectHash: string; + sessionUuid: string; + runId: string; +} + /** A `workflowProgress[]` entry as it appears on disk (loosely typed for defensive parsing). */ interface RawProgressEntry { type?: string; @@ -94,7 +119,11 @@ export class WorkflowRunWatcher extends EventEmitter { private fileMtimes = new Map(); /** runId -> absolute run-file path (for mtime cleanup on removal). */ private runIdToPath = new Map(); - /** workflows-dir absolute path -> chokidar watcher (one per live run dir). */ + /** absolute live transcript-dir path -> newest member mtimeMs (skip unchanged live runs). */ + private liveDirMtimes = new Map(); + /** runId -> absolute live transcript-dir path (for mtime cleanup on removal). */ + private runIdToLiveDir = new Map(); + /** watched-dir absolute path -> chokidar watcher (workflows/ + subagents/workflows/). */ private dirWatchers = new Map(); constructor(projectsDir?: string) { @@ -129,6 +158,8 @@ export class WorkflowRunWatcher extends EventEmitter { this.runs.clear(); this.fileMtimes.clear(); this.runIdToPath.clear(); + this.liveDirMtimes.clear(); + this.runIdToLiveDir.clear(); } /** All cached runs (no recency filter), most-recently-active first. */ @@ -180,42 +211,65 @@ export class WorkflowRunWatcher extends EventEmitter { } private async pollAsync(): Promise { - const { files, dirs } = await this.discover(); + const { files, liveDirs, watchDirs } = await this.discover(); - // Install a live watcher for each run dir; tear down watchers for dirs that vanished. - for (const dir of dirs) this.ensureDirWatcher(dir); + // Install a live watcher for each watched dir; tear down watchers for dirs that vanished. + for (const dir of watchDirs) this.ensureDirWatcher(dir); for (const dir of Array.from(this.dirWatchers.keys())) { - if (!dirs.has(dir)) this.removeDirWatcher(dir); + if (!watchDirs.has(dir)) this.removeDirWatcher(dir); } const seenRunIds = new Set(); + const realRunIds = new Set(); for (const file of files) { seenRunIds.add(file.runId); + realRunIds.add(file.runId); await this.maybeParse(file); } - // Removal by set-diff: a cached run whose file disappeared. + // In-flight runs: synthesize from the transcript tree ONLY while no completion + // wf_*.json exists yet — the real file (full detail + terminal status) supersedes. + for (const live of liveDirs) { + if (realRunIds.has(live.runId)) continue; + seenRunIds.add(live.runId); + await this.maybeParseLive(live); + } + + // Removal by set-diff: a cached run discoverable from neither source. for (const runId of Array.from(this.runs.keys())) { if (!seenRunIds.has(runId)) { this.runs.delete(runId); const path = this.runIdToPath.get(runId); if (path) this.fileMtimes.delete(path); this.runIdToPath.delete(runId); + const liveDir = this.runIdToLiveDir.get(runId); + if (liveDir) this.liveDirMtimes.delete(liveDir); + this.runIdToLiveDir.delete(runId); this.emit('run_removed', { runId }); } } } - /** Walk projects///workflows/ for wf_*.json files. */ - private async discover(): Promise<{ files: DiscoveredRun[]; dirs: Set }> { + /** + * Walk projects/// for both run sources: + * - completion files: `workflows/wf_*.json` + * - in-flight runs: `subagents/workflows/wf_/` (transcript dirs) + * Returns the dirs to chokidar-watch (so a new run/file is caught sub-poll). + */ + private async discover(): Promise<{ + files: DiscoveredRun[]; + liveDirs: DiscoveredLiveRun[]; + watchDirs: Set; + }> { const files: DiscoveredRun[] = []; - const dirs = new Set(); + const liveDirs: DiscoveredLiveRun[] = []; + const watchDirs = new Set(); let projectHashes: string[]; try { projectHashes = await readdir(this.projectsDir); } catch { - return { files, dirs }; + return { files, liveDirs, watchDirs }; } for (const projectHash of projectHashes) { @@ -226,28 +280,50 @@ export class WorkflowRunWatcher extends EventEmitter { continue; } for (const sessionUuid of sessions) { - const workflowsDir = join(this.projectsDir, projectHash, sessionUuid, WORKFLOWS_SUBDIR); - let names: string[]; + const sessionDir = join(this.projectsDir, projectHash, sessionUuid); + + // (1) Completion artifacts: workflows/wf_*.json + const workflowsDir = join(sessionDir, WORKFLOWS_SUBDIR); try { - names = await readdir(workflowsDir); + const names = await readdir(workflowsDir); + let hasRun = false; + for (const name of names) { + if (!name.startsWith(RUN_FILE_PREFIX) || !name.endsWith(RUN_FILE_SUFFIX)) continue; + hasRun = true; + files.push({ + filePath: join(workflowsDir, name), + projectHash, + sessionUuid, + runId: name.slice(0, -RUN_FILE_SUFFIX.length), + }); + } + if (hasRun) watchDirs.add(workflowsDir); } catch { - continue; // no workflows dir for this session — normal + // no workflows dir for this session — normal } - let hasRun = false; - for (const name of names) { - if (!name.startsWith(RUN_FILE_PREFIX) || !name.endsWith(RUN_FILE_SUFFIX)) continue; - hasRun = true; - files.push({ - filePath: join(workflowsDir, name), - projectHash, - sessionUuid, - runId: name.slice(0, -RUN_FILE_SUFFIX.length), - }); + + // (2) In-flight runs: subagents/workflows/wf_*/ + const liveParent = join(sessionDir, SUBAGENTS_SUBDIR, WORKFLOWS_SUBDIR); + try { + const names = await readdir(liveParent); + let hasLive = false; + for (const name of names) { + if (!name.startsWith(RUN_FILE_PREFIX) || name.endsWith(RUN_FILE_SUFFIX)) continue; // wf_ dir, not a file + hasLive = true; + liveDirs.push({ + dirPath: join(liveParent, name), + projectHash, + sessionUuid, + runId: name, + }); + } + if (hasLive) watchDirs.add(liveParent); + } catch { + // no subagents/workflows dir for this session — normal } - if (hasRun) dirs.add(workflowsDir); } } - return { files, dirs }; + return { files, liveDirs, watchDirs }; } private async maybeParse(file: DiscoveredRun): Promise { @@ -269,6 +345,101 @@ export class WorkflowRunWatcher extends EventEmitter { this.emit(existed ? 'run_updated' : 'run_discovered', info); } + /** + * 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). + */ + private async maybeParseLive(live: DiscoveredLiveRun): Promise { + const info = await this.parseLiveDir(live); + 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); + this.runIdToLiveDir.set(info.runId, live.dirPath); + this.emit(existed ? 'run_updated' : 'run_discovered', info); + } + + /** + * 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. + */ + private async parseLiveDir(live: DiscoveredLiveRun): Promise { + let entries: string[]; + try { + entries = await readdir(live.dirPath); + } catch { + return null; // vanished between discover and read + } + + 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 + + const doneIds = await this.readJournalDoneAgents(join(live.dirPath, LIVE_JOURNAL_FILE)); + const agents: WorkflowAgentInfo[] = Array.from(agentIds) + .sort() + .map((id, i) => ({ + index: i + 1, + label: `agent ${i + 1}`, + phaseIndex: 1, + phaseTitle: '', + model: '', + state: doneIds.has(id) ? 'done' : 'progress', + agentId: id, + })); + + return { + runId: live.runId, + status: 'running', + agentCount: agents.length, + phases: [], + agents, + sessionUuid: live.sessionUuid, + projectHash: live.projectHash, + lastActivityAt: newestMtime || 0, + }; + } + + /** Agent ids that already emitted a `result` event in the run journal. */ + private async readJournalDoneAgents(journalPath: string): Promise> { + const done = new Set(); + let text: string; + try { + text = await readFile(journalPath, 'utf-8'); + } catch { + return done; // journal not written yet — all agents still in progress + } + 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); + } catch { + // tolerate a partially-written trailing line + } + } + return done; + } + /** * Parse a wf_.json into WorkflowRunInfo, STRIPPING the heavyweight * `script`/`scriptPath`/`result`/`logs` fields (the embedded script alone is @@ -360,6 +531,9 @@ export class WorkflowRunWatcher extends EventEmitter { watcher.on('add', handler); watcher.on('change', handler); watcher.on('unlink', handler); + // subagents/workflows/ children are wf_/ DIRS — catch their add/remove too. + watcher.on('addDir', handler); + watcher.on('unlinkDir', handler); watcher.on('error', () => { // chokidar surfaced an error for this dir — drop the watcher; poll still covers it. this.removeDirWatcher(workflowsDir); diff --git a/test/workflow-run-watcher.test.ts b/test/workflow-run-watcher.test.ts index deab7f1c..39f152be 100644 --- a/test/workflow-run-watcher.test.ts +++ b/test/workflow-run-watcher.test.ts @@ -210,3 +210,110 @@ describe('WorkflowRunWatcher', () => { expect(summaries[0].agentCount).toBe(3); }); }); + +/** + * In-flight runs: the Workflow runtime writes the completion wf_.json only when + * a run FINISHES, so while it is live the only on-disk state is its + * subagents/workflows/wf_/ transcript dir. The watcher synthesizes a minimal + * ACTIVE run from that dir so the floating window pops DURING the run. + */ +describe('WorkflowRunWatcher — in-flight (live) runs', () => { + const LIVE_RUN_ID = 'wf_live5678-xyz'; + let projectsDir: string; + let liveDir: string; + let watcher: WorkflowRunWatcher; + + beforeEach(async () => { + projectsDir = await mkdtemp(join(tmpdir(), 'wfw-live-')); + liveDir = join(projectsDir, PROJECT_HASH, SESSION_UUID, 'subagents', 'workflows', LIVE_RUN_ID); + await mkdir(liveDir, { recursive: true }); + // Two agents started; one already produced a result (journal `result` line). + await writeFile( + join(liveDir, 'agent-aaa111.meta.json'), + JSON.stringify({ agentType: 'workflow-subagent' }), + 'utf-8' + ); + await writeFile(join(liveDir, 'agent-aaa111.jsonl'), '{"type":"assistant"}\n', 'utf-8'); + await writeFile( + join(liveDir, 'agent-bbb222.meta.json'), + JSON.stringify({ agentType: 'workflow-subagent' }), + 'utf-8' + ); + await writeFile(join(liveDir, 'agent-bbb222.jsonl'), '{"type":"assistant"}\n', 'utf-8'); + await writeFile( + join(liveDir, 'journal.jsonl'), + '{"type":"started","agentId":"aaa111"}\n{"type":"started","agentId":"bbb222"}\n{"type":"result","agentId":"aaa111","result":{}}\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 timer = setTimeout(() => reject(new Error('timed out waiting for run_discovered')), 5000); + watcher.once('run_discovered', (info: WorkflowRunInfo) => { + clearTimeout(timer); + resolve(info); + }); + watcher.start(); + }); + } + + it('synthesizes an ACTIVE run from the transcript dir when no completion file exists', async () => { + const info = await firstRun(); + expect(info.runId).toBe(LIVE_RUN_ID); + expect(info.status).toBe('running'); // active → frontend pops a floating window + expect(info.sessionUuid).toBe(SESSION_UUID); + expect(info.projectHash).toBe(PROJECT_HASH); + expect(info.agents).toHaveLength(2); + expect(info.agentCount).toBe(2); + expect(info.lastActivityAt).toBeGreaterThan(0); + }); + + it('preserves agentId per slot (so the card→transcript click join still works)', async () => { + const info = await firstRun(); + const ids = info.agents.map((a) => a.agentId).sort(); + expect(ids).toEqual(['aaa111', 'bbb222']); + }); + + it('marks an agent done/progress from journal result lines', async () => { + const info = await firstRun(); + expect(info.agents.find((a) => a.agentId === 'aaa111')!.state).toBe('done'); // has a result line + expect(info.agents.find((a) => a.agentId === 'bbb222')!.state).toBe('progress'); // started, no result yet + }); + + it('counts the live run as running in getStats', async () => { + await firstRun(); + expect(watcher.getStats().running).toBe(1); + }); + + it('does NOT surface a live dir that has no agent files yet', async () => { + const empty = join(projectsDir, PROJECT_HASH, SESSION_UUID, 'subagents', 'workflows', 'wf_empty0000-noo'); + await mkdir(empty, { recursive: true }); + await firstRun(); // resolves on the real (populated) live run + // The empty run id must never enter the cache. + expect(watcher.getRun('wf_empty0000-noo')).toBeUndefined(); + expect(watcher.getAllRuns().map((r) => r.runId)).toEqual([LIVE_RUN_ID]); + }); + + it('a completion wf_*.json supersedes the live dir for the same runId (real status wins)', async () => { + const workflowsDir = join(projectsDir, PROJECT_HASH, SESSION_UUID, 'workflows'); + await mkdir(workflowsDir, { recursive: true }); + await writeFile( + join(workflowsDir, `${LIVE_RUN_ID}.json`), + JSON.stringify({ runId: LIVE_RUN_ID, status: 'completed', durationMs: 1234, phases: [], workflowProgress: [] }), + 'utf-8' + ); + const info = await firstRun(); + expect(info.runId).toBe(LIVE_RUN_ID); + expect(info.status).toBe('completed'); // real completion file wins, not synthesized 'running' + expect(info.durationMs).toBe(1234); + // Only one cached entry for the runId — no live/real duplication. + expect(watcher.getAllRuns()).toHaveLength(1); + }); +});