Compare commits

...
9 changed files with 566 additions and 66 deletions
+14
View File
@@ -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-<id>.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/<name>-<runId>.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
+1 -1
View File
@@ -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
+2 -2
View File
@@ -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": [
+1 -1
View File
@@ -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",
+28 -9
View File
@@ -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 {
+15 -2
View File
@@ -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-<id>.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 (
`<div${cardAttrs}>` +
`<div class="ultracode-agent-top">` +
+29 -2
View File
@@ -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();
+380 -49
View File
@@ -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-<id>.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/<name>-<runId>.js`. */
interface LiveWorkflowMeta {
workflowName?: string;
summary?: string;
phases: WorkflowRunPhase[];
}
/** A live agent file on disk (transcript and/or meta), keyed by its `<id>` 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<string, number>();
/** runId -> absolute live transcript-dir path (for mtime cleanup on removal). */
private runIdToLiveDir = new Map<string, string>();
/** transcript abs path -> { mtimeMs, stats }; re-parse a transcript only when its mtime moves. */
private agentStatCache = new LRUMap<string, { mtimeMs: number; stats: AgentTranscriptStats }>({
maxSize: MAX_CACHED_AGENT_STATS,
});
/** runId -> derived script meta (immutable per run; re-derived only until a name is found). */
private liveMetaCache = new Map<string, LiveWorkflowMeta>();
/** 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<string> } }
>();
/** watched-dir absolute path -> chokidar watcher (workflows/ + subagents/workflows/). */
private dirWatchers = new Map<string, ChokidarWatcher>();
@@ -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<void> {
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<string, LiveAgentFile>();
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-<stem>.jsonl (transcript) or agent-<stem>.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_<id>/`.
* 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-<id>.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/<name>-<runId>.js`
* (the completion wf_<id>.json, which carries these, only lands at the end);
* - run totals = sums of the per-agent stats.
*/
private async parseLiveDir(live: DiscoveredLiveRun): Promise<WorkflowRunInfo | null> {
let entries: string[];
try {
entries = await readdir(live.dirPath);
} catch {
return null; // vanished between discover and read
}
private async parseLiveDir(
live: DiscoveredLiveRun,
agentFiles: Map<string, LiveAgentFile>,
journalPath: string | null,
journalMtime: number,
newestMtime: number
): Promise<WorkflowRunInfo | null> {
const { startedOrder, doneIds } = journalPath
? await this.readJournal(journalPath, journalMtime)
: { startedOrder: [] as string[], doneIds: new Set<string>() };
const startIndex = new Map<string, number>();
startedOrder.forEach((id, i) => startIndex.set(id, i));
const agentIds = new Set<string>();
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<Set<string>> {
const done = new Set<string>();
/** Parse one agent transcript for token/tool stats; cached on the file's mtime. */
private async agentTranscriptStats(path: string, mtimeMs: number): Promise<AgentTranscriptStats | null> {
// 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<string> }> {
const cached = this.journalCache.get(journalPath);
if (cached && cached.mtimeMs === mtimeMs) return cached.parsed;
const startedOrder: string[] = [];
const seenStarted = new Set<string>();
const doneIds = new Set<string>();
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/<name>-<runId>.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<LiveWorkflowMeta> {
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),
};
}
/**
+96
View File
@@ -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<WorkflowRunInfo> {
return new Promise<WorkflowRunInfo>((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);
});
});