mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-01 04:59:41 +02:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
98b2124d7e |
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
Generated
+2
-2
@@ -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
@@ -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",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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">` +
|
||||
|
||||
@@ -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
@@ -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),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user