feat: Claude Code Monitor — lanes, pipelines and a merged workspace
Internal SmartGift build of a Claude Code monitoring dashboard. Lanes: a durable unit of parallel agent work, one per working directory, tracked across session restarts. Managed lanes are git worktrees the dashboard provisions and can reset or remove behind a three-check destroy guard and a counted preflight; adopted lanes are directories you already own and are never destroyable. Pipelines: a lane moves through pipeline stages. A stage the agent declares with evidence renders green; a stage inferred from the tool-event stream renders dashed amber and never counts as done. Detection is forward-only within a 30-minute window, and never writes the declared stage. Workspace: one page at /run with a lane grid, the selected lane's pipeline, and a full Claude console behind a disclosure.
This commit is contained in:
@@ -0,0 +1,807 @@
|
||||
/**
|
||||
* Workflow-tool run ingestion.
|
||||
*
|
||||
* The Claude Code "Workflow" tool (and self-paced /loop) spawn fleets of inner
|
||||
* sub-agents that emit NO hooks — so hook-based ingestion can never see them.
|
||||
* Everything lives on disk under the launching session's transcript folder:
|
||||
*
|
||||
* <projects>/<enc-cwd>/<sessionId>/
|
||||
* workflows/
|
||||
* scripts/<name>-wf_<runId>.js ← written at LAUNCH
|
||||
* wf_<runId>.json ← run journal, written at COMPLETION
|
||||
* subagents/workflows/<runId>/
|
||||
* agent-<agentId>.jsonl ← one transcript per inner agent
|
||||
* agent-<agentId>.meta.json
|
||||
*
|
||||
* The run journal is the source of truth for a completed run: identity,
|
||||
* lifecycle, aggregates (agentCount/totalTokens/totalToolCalls), phases[], and
|
||||
* workflowProgress[] — a MIXED log of `type:"workflow_phase"` markers and
|
||||
* `type:"workflow_agent"` entries. Each workflow_agent entry carries agentId,
|
||||
* state ("done"/"error"/…), label, phaseTitle, tokens, toolCalls, durationMs,
|
||||
* etc., and its agentId is the EXACT agent-<agentId>.jsonl basename in the
|
||||
* per-run nested dir above. Because the journal is terminal-only, a running
|
||||
* workflow is detected from its launch script and replaced by the journal
|
||||
* record on completion (idempotent upsert by run_id).
|
||||
*
|
||||
* Inner agents are linked into the existing agents table via the same
|
||||
* `${sessionId}-jsonl-<agentId>` id scheme that importSubagentFromJsonl uses,
|
||||
* so ingestion CONVERGES with any prior subagent import (no duplicate rows).
|
||||
* Per-agent token/tool/duration metrics come from the journal's progress[]
|
||||
* JSON — this module never writes token_usage, so it cannot double-count.
|
||||
*
|
||||
* All functions are fail-safe: a malformed/partial journal throws only locally
|
||||
* and is skipped; ingestion never blocks or breaks hook handling.
|
||||
* @author Nguyễn Ngọc Trí Vĩ <vinnt@smartgift.vn>
|
||||
*/
|
||||
|
||||
const fs = require("fs");
|
||||
const path = require("path");
|
||||
|
||||
// Lazy-required to avoid a require cycle (import-history → db → … ) and to keep
|
||||
// startup cheap; mirrors how server/index.js lazy-requires import helpers.
|
||||
function importHistory() {
|
||||
return require("../../scripts/import-history");
|
||||
}
|
||||
|
||||
let claudeHome = null;
|
||||
function getClaudeHomeLib() {
|
||||
if (!claudeHome) claudeHome = require("./claude-home");
|
||||
return claudeHome;
|
||||
}
|
||||
|
||||
/**
|
||||
* Canonical run id derived from a journal/script filename. Both
|
||||
* `wf_<runId>.json` and `<name>-wf_<runId>.js` reduce to the same `wf_<runId>`
|
||||
* token so a launch-detected "running" row and its later journal reconcile on
|
||||
* the same key.
|
||||
*/
|
||||
function extractRunId(filename) {
|
||||
const base = path.basename(filename).replace(/\.(json|js)$/i, "");
|
||||
const m = base.match(/wf_[A-Za-z0-9_-]+$/);
|
||||
return m ? m[0] : base;
|
||||
}
|
||||
|
||||
/** Workflow name from a launch-script basename: strip the `-wf_<runId>` tail. */
|
||||
function nameFromScript(filename) {
|
||||
const base = path.basename(filename).replace(/\.js$/i, "");
|
||||
return base.replace(/-?wf_[A-Za-z0-9_-]+$/, "") || base;
|
||||
}
|
||||
|
||||
function toIso(value) {
|
||||
if (value == null) return null;
|
||||
if (typeof value === "number") {
|
||||
try {
|
||||
return new Date(value).toISOString();
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
return String(value);
|
||||
}
|
||||
|
||||
/** Map a journal progress `state` to an agents.status value. */
|
||||
function mapState(state) {
|
||||
switch (String(state || "").toLowerCase()) {
|
||||
case "error":
|
||||
case "failed":
|
||||
return "error";
|
||||
case "running":
|
||||
case "working":
|
||||
case "active":
|
||||
case "in_progress":
|
||||
case "queued":
|
||||
return "working";
|
||||
case "done":
|
||||
case "completed":
|
||||
case "success":
|
||||
return "completed";
|
||||
default:
|
||||
return "completed";
|
||||
}
|
||||
}
|
||||
|
||||
// Token fields carried on a parsed-subagent bucket (camelCase, matching
|
||||
// writeSessionTokens). Used to fold inner-agent usage into the session's cost.
|
||||
const TOKEN_FIELDS = [
|
||||
"input",
|
||||
"output",
|
||||
"cacheRead",
|
||||
"cacheWrite",
|
||||
"cacheWrite1h",
|
||||
"webSearch",
|
||||
"webFetch",
|
||||
"codeExec",
|
||||
];
|
||||
|
||||
/**
|
||||
* Merge a parsed agent's tokensByModel into a session-level accumulator, keyed
|
||||
* by (model, speed, geo) with the service_tier forced to "workflow". This
|
||||
* namespaces workflow spend into its own token_usage bucket so it never
|
||||
* collides with — or clobbers — the main-transcript writer's rows, while still
|
||||
* being summed per-model by the cost calculator. Inner agents are sidechain
|
||||
* contexts whose usage is NOT in the parent transcript, so this is additive,
|
||||
* not double-counting (same model as combineSessionTokens for subagents).
|
||||
*/
|
||||
function mergeWorkflowTokens(dst, src) {
|
||||
for (const b of Object.values(src || {})) {
|
||||
if (!b || !b.model) continue;
|
||||
const key = `${b.model}|${b.speed}|${b.geo}|workflow`;
|
||||
if (!dst[key]) {
|
||||
dst[key] = {
|
||||
model: b.model,
|
||||
speed: b.speed,
|
||||
geo: b.geo,
|
||||
tier: "workflow",
|
||||
};
|
||||
for (const f of TOKEN_FIELDS) dst[key][f] = 0;
|
||||
}
|
||||
for (const f of TOKEN_FIELDS) dst[key][f] += b[f] || 0;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a session's transcript JSONL path from a session-like row. Prefers an
|
||||
* explicit transcript_path; otherwise derives it from (id, cwd) via claude-home.
|
||||
*/
|
||||
function resolveTranscriptPath(session) {
|
||||
if (session && session.transcript_path) return session.transcript_path;
|
||||
if (session && session.id && session.cwd) {
|
||||
try {
|
||||
return getClaudeHomeLib().getTranscriptPath(session.id, session.cwd);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Locate a session's workflow artifacts from its transcript JSONL path.
|
||||
* Workflows live at `<dir>/<sessionId>/workflows/` next to
|
||||
* `<dir>/<sessionId>.jsonl`; inner-agent transcripts are resolved per-run via
|
||||
* agentsDirForRun(sessionDir, runId).
|
||||
*
|
||||
* @returns {{ sessionDir: string|null, workflowsDir: string|null,
|
||||
* journals: string[], scripts: string[] }}
|
||||
*/
|
||||
function findSessionWorkflows(transcriptPath) {
|
||||
const empty = {
|
||||
sessionDir: null,
|
||||
workflowsDir: null,
|
||||
journals: [],
|
||||
scripts: [],
|
||||
liveRuns: [],
|
||||
};
|
||||
if (!transcriptPath) return empty;
|
||||
const dir = path.dirname(transcriptPath);
|
||||
const sessionId = path.basename(transcriptPath, ".jsonl");
|
||||
const sessionDir = path.join(dir, sessionId);
|
||||
const workflowsDir = path.join(sessionDir, "workflows");
|
||||
|
||||
const journals = [];
|
||||
const scripts = [];
|
||||
try {
|
||||
if (fs.existsSync(workflowsDir)) {
|
||||
for (const f of fs.readdirSync(workflowsDir)) {
|
||||
if (f.startsWith("wf_") && f.endsWith(".json")) journals.push(path.join(workflowsDir, f));
|
||||
}
|
||||
const scriptsDir = path.join(workflowsDir, "scripts");
|
||||
if (fs.existsSync(scriptsDir)) {
|
||||
for (const f of fs.readdirSync(scriptsDir)) {
|
||||
if (f.endsWith(".js")) scripts.push(path.join(scriptsDir, f));
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
/* non-fatal — partial dir during a live run */
|
||||
}
|
||||
|
||||
// Live per-run dirs: <sessionDir>/subagents/workflows/<runId>/ — present while
|
||||
// a workflow is still running (journal.jsonl + growing agent-*.jsonl), before
|
||||
// the terminal wf_<runId>.json journal is written.
|
||||
const liveRuns = [];
|
||||
try {
|
||||
const base = path.join(sessionDir, "subagents", "workflows");
|
||||
if (fs.existsSync(base)) {
|
||||
for (const d of fs.readdirSync(base, { withFileTypes: true })) {
|
||||
if (d.isDirectory()) liveRuns.push({ runId: d.name, dir: path.join(base, d.name) });
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
/* non-fatal */
|
||||
}
|
||||
|
||||
return { sessionDir, workflowsDir, journals, scripts, liveRuns };
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-run inner-agent transcript directory. The Workflow tool writes each
|
||||
* fleet's agents under `<sessionId>/subagents/workflows/<runId>/agent-*.jsonl`
|
||||
* (NOT the session's top-level subagents/ dir).
|
||||
*/
|
||||
function agentsDirForRun(sessionDir, runId) {
|
||||
return path.join(sessionDir, "subagents", "workflows", runId);
|
||||
}
|
||||
|
||||
/** Read + normalize a run journal file. Returns null on any parse failure. */
|
||||
function parseWorkflowJournal(journalPath) {
|
||||
let raw;
|
||||
try {
|
||||
raw = fs.readFileSync(journalPath, "utf8");
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
let j;
|
||||
try {
|
||||
j = JSON.parse(raw);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
const runId = extractRunId(journalPath) || j.runId || null;
|
||||
if (!runId) return null;
|
||||
|
||||
const startedAt = toIso(j.startTime != null ? j.startTime : j.startedAt);
|
||||
const durationMs = Number.isFinite(j.durationMs) ? j.durationMs : null;
|
||||
let endedAt = toIso(j.endTime != null ? j.endTime : j.endedAt);
|
||||
if (!endedAt && startedAt && durationMs != null) {
|
||||
const t = Date.parse(startedAt);
|
||||
if (!Number.isNaN(t)) endedAt = new Date(t + durationMs).toISOString();
|
||||
}
|
||||
const progress = Array.isArray(j.workflowProgress)
|
||||
? j.workflowProgress
|
||||
: Array.isArray(j.progress)
|
||||
? j.progress
|
||||
: [];
|
||||
|
||||
return {
|
||||
runId,
|
||||
taskId: j.taskId || null,
|
||||
name: j.workflowName || j.name || nameFromScript(journalPath),
|
||||
status: String(j.status || "completed"),
|
||||
defaultModel: j.defaultModel || null,
|
||||
startedAt,
|
||||
endedAt,
|
||||
durationMs,
|
||||
agentCount: Number.isFinite(j.agentCount)
|
||||
? j.agentCount
|
||||
: progress.filter((e) => e && e.type === "workflow_agent").length,
|
||||
totalTokens: Number.isFinite(j.totalTokens) ? j.totalTokens : 0,
|
||||
totalToolCalls: Number.isFinite(j.totalToolCalls) ? j.totalToolCalls : 0,
|
||||
phases: Array.isArray(j.phases) ? j.phases : [],
|
||||
progress,
|
||||
journalPath,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Ingest one parsed journal: upsert the workflow row, then link/create each
|
||||
* inner agent. Returns the upserted workflow row, or null on failure.
|
||||
*/
|
||||
async function ingestWorkflowJournal(dbModule, sessionId, journal, opts = {}) {
|
||||
const { stmts } = dbModule;
|
||||
const mainAgentId = `${sessionId}-main`;
|
||||
const ih = importHistory();
|
||||
// Inner-agent transcripts live in a per-run nested dir, not the session's
|
||||
// top-level subagents/. opts.sessionDir is the session transcript folder.
|
||||
const agentDir = opts.sessionDir ? agentsDirForRun(opts.sessionDir, journal.runId) : null;
|
||||
// Accumulate inner-agent token usage (real input/output/cache split from each
|
||||
// transcript) so the run's spend can be folded into the session's cost.
|
||||
const runTokens = {};
|
||||
|
||||
stmts.upsertWorkflow.run(
|
||||
journal.runId,
|
||||
sessionId,
|
||||
journal.taskId,
|
||||
journal.name,
|
||||
journal.status,
|
||||
journal.defaultModel,
|
||||
journal.startedAt,
|
||||
journal.endedAt,
|
||||
journal.durationMs,
|
||||
journal.agentCount,
|
||||
journal.totalTokens,
|
||||
journal.totalToolCalls,
|
||||
JSON.stringify(journal.phases),
|
||||
JSON.stringify(journal.progress),
|
||||
opts.scriptPath || null,
|
||||
journal.journalPath || null,
|
||||
"journal"
|
||||
);
|
||||
|
||||
// Only `workflow_agent` entries are real agents; `workflow_phase` entries are
|
||||
// phase markers (kept in progress[] for the phase chips, skipped here).
|
||||
const agentEntries = journal.progress.filter(
|
||||
(e) => e && e.type === "workflow_agent" && e.agentId
|
||||
);
|
||||
for (const entry of agentEntries) {
|
||||
const agentId = entry.agentId;
|
||||
const jsonlId = `${sessionId}-jsonl-${agentId}`;
|
||||
const status = mapState(entry.state);
|
||||
const phase = entry.phaseTitle || null;
|
||||
// subagent_type: prefer the label's prefix (e.g. "scout:starship" → "scout")
|
||||
// for nicer grouping; otherwise the generic workflow-subagent type.
|
||||
const subType =
|
||||
(entry.label && entry.label.includes(":") ? entry.label.split(":")[0] : null) ||
|
||||
entry.agentType ||
|
||||
"workflow-subagent";
|
||||
|
||||
// Prefer parsing the real transcript so tool events + metadata land via the
|
||||
// shared importer (idempotent, dedups by tool_use_id). Fall back to a
|
||||
// minimal row built from the journal entry if the file is gone.
|
||||
let parsed = null;
|
||||
if (agentDir) {
|
||||
const subPath = path.join(agentDir, `agent-${agentId}.jsonl`);
|
||||
if (fs.existsSync(subPath)) {
|
||||
try {
|
||||
parsed = await ih.parseSubagentFile(subPath);
|
||||
} catch {
|
||||
parsed = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
try {
|
||||
if (parsed) {
|
||||
ih.importSubagentFromJsonl(dbModule, sessionId, mainAgentId, parsed);
|
||||
mergeWorkflowTokens(runTokens, parsed.tokensByModel);
|
||||
} else if (!stmts.getAgent.get(jsonlId)) {
|
||||
stmts.insertAgent.run(
|
||||
jsonlId,
|
||||
sessionId,
|
||||
entry.label || `Subagent ${String(agentId).slice(0, 8)}`,
|
||||
"subagent",
|
||||
subType,
|
||||
status,
|
||||
entry.label || entry.promptPreview || null,
|
||||
mainAgentId,
|
||||
JSON.stringify({
|
||||
imported: true,
|
||||
source: "workflow",
|
||||
workflow_run_id: journal.runId,
|
||||
model: entry.model || null,
|
||||
tokens: entry.tokens || 0,
|
||||
tool_calls: entry.toolCalls || 0,
|
||||
})
|
||||
);
|
||||
}
|
||||
// Stamp the workflow linkage + journal-authoritative status/phase.
|
||||
stmts.setAgentWorkflow.run(journal.runId, phase, status, jsonlId);
|
||||
} catch {
|
||||
/* one bad agent must not abort the whole run ingest */
|
||||
}
|
||||
}
|
||||
|
||||
return { row: stmts.getWorkflow.get(journal.runId), tokens: runTokens };
|
||||
}
|
||||
|
||||
function shortLabel(s) {
|
||||
if (!s) return null;
|
||||
const first = String(s).split("\n")[0].trim();
|
||||
return first.length > 80 ? first.slice(0, 79) + "…" : first;
|
||||
}
|
||||
|
||||
function safeStringify(v) {
|
||||
if (v == null) return null;
|
||||
if (typeof v === "string") return v;
|
||||
try {
|
||||
return JSON.stringify(v);
|
||||
} catch {
|
||||
return String(v);
|
||||
}
|
||||
}
|
||||
|
||||
function bucketTotal(tokensByModel) {
|
||||
let n = 0;
|
||||
for (const b of Object.values(tokensByModel || {})) {
|
||||
n +=
|
||||
(b.input || 0) +
|
||||
(b.output || 0) +
|
||||
(b.cacheRead || 0) +
|
||||
(b.cacheWrite || 0) +
|
||||
(b.cacheWrite1h || 0);
|
||||
}
|
||||
return n;
|
||||
}
|
||||
|
||||
/**
|
||||
* Live ingest for a RUNNING workflow — before its terminal wf_<runId>.json
|
||||
* exists. Builds progress[] + aggregates in real time from the streaming
|
||||
* `<runDir>/journal.jsonl` (started/result events per agent) plus the growing
|
||||
* `<runDir>/agent-<id>.jsonl` transcripts (real token/tool/duration usage via
|
||||
* parseSubagentFile). Phase/label aren't available live (those come from the
|
||||
* terminal journal), so phaseTitle is null and label falls back to the agent's
|
||||
* prompt. The fast poll re-runs this as the files grow, so tokens/tools/agents
|
||||
* update live. Returns { row, tokens } or null.
|
||||
*/
|
||||
async function ingestLiveWorkflow(dbModule, sessionId, sessionDir, runId, scriptPath) {
|
||||
const { stmts } = dbModule;
|
||||
const mainAgentId = `${sessionId}-main`;
|
||||
const ih = importHistory();
|
||||
const dir = agentsDirForRun(sessionDir, runId);
|
||||
if (!fs.existsSync(dir)) return null;
|
||||
|
||||
// Streaming journal: which agents started / finished (+ their result payload).
|
||||
const started = new Set();
|
||||
const doneResults = new Map();
|
||||
try {
|
||||
const jj = path.join(dir, "journal.jsonl");
|
||||
if (fs.existsSync(jj)) {
|
||||
for (const line of fs.readFileSync(jj, "utf8").split("\n")) {
|
||||
if (!line.trim()) continue;
|
||||
let o;
|
||||
try {
|
||||
o = JSON.parse(line);
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
if (!o || !o.agentId) continue;
|
||||
if (o.type === "started") started.add(o.agentId);
|
||||
else if (o.type === "result") doneResults.set(o.agentId, o.result);
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
|
||||
let agentFiles = [];
|
||||
try {
|
||||
agentFiles = fs.readdirSync(dir).filter((f) => f.startsWith("agent-") && f.endsWith(".jsonl"));
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
if (agentFiles.length === 0 && started.size === 0) return null;
|
||||
|
||||
const progress = [];
|
||||
const runTokens = {};
|
||||
let totalTokens = 0;
|
||||
let totalToolCalls = 0;
|
||||
let earliest = null;
|
||||
let latest = null;
|
||||
let model = null;
|
||||
|
||||
for (const f of agentFiles) {
|
||||
const agentId = f.replace(/^agent-/, "").replace(/\.jsonl$/, "");
|
||||
let parsed = null;
|
||||
try {
|
||||
parsed = await ih.parseSubagentFile(path.join(dir, f));
|
||||
} catch {
|
||||
parsed = null;
|
||||
}
|
||||
const done = doneResults.has(agentId);
|
||||
const state = done ? "done" : "running";
|
||||
const aTok = parsed ? bucketTotal(parsed.tokensByModel) : 0;
|
||||
const tools = parsed && parsed.toolNames ? parsed.toolNames : [];
|
||||
const startedAt = parsed && parsed.startedAt ? parsed.startedAt : null;
|
||||
const endedAt = parsed && parsed.endedAt ? parsed.endedAt : null;
|
||||
const durationMs = startedAt && endedAt ? Date.parse(endedAt) - Date.parse(startedAt) : null;
|
||||
const label = parsed && parsed.task ? shortLabel(parsed.task) : null;
|
||||
if (parsed && parsed.model && !model) model = parsed.model;
|
||||
totalTokens += aTok;
|
||||
totalToolCalls += tools.length;
|
||||
if (startedAt) {
|
||||
const ts = Date.parse(startedAt);
|
||||
if (!earliest || ts < earliest) earliest = ts;
|
||||
}
|
||||
if (endedAt) {
|
||||
const ts = Date.parse(endedAt);
|
||||
if (!latest || ts > latest) latest = ts;
|
||||
}
|
||||
|
||||
progress.push({
|
||||
type: "workflow_agent",
|
||||
agentId,
|
||||
label,
|
||||
phaseTitle: null,
|
||||
model: parsed ? parsed.model : null,
|
||||
state,
|
||||
tokens: aTok,
|
||||
toolCalls: tools.length,
|
||||
durationMs,
|
||||
lastToolName: tools.length ? tools[tools.length - 1] : null,
|
||||
promptPreview: parsed ? parsed.task : null,
|
||||
resultPreview: done ? safeStringify(doneResults.get(agentId)) : null,
|
||||
});
|
||||
|
||||
try {
|
||||
const jsonlId = `${sessionId}-jsonl-${agentId}`;
|
||||
if (parsed) {
|
||||
ih.importSubagentFromJsonl(dbModule, sessionId, mainAgentId, parsed);
|
||||
mergeWorkflowTokens(runTokens, parsed.tokensByModel);
|
||||
} else if (!stmts.getAgent.get(jsonlId)) {
|
||||
stmts.insertAgent.run(
|
||||
jsonlId,
|
||||
sessionId,
|
||||
label || `Subagent ${agentId.slice(0, 8)}`,
|
||||
"subagent",
|
||||
"workflow-subagent",
|
||||
mapState(state),
|
||||
label,
|
||||
mainAgentId,
|
||||
JSON.stringify({ imported: true, source: "workflow-live", workflow_run_id: runId })
|
||||
);
|
||||
}
|
||||
stmts.setAgentWorkflow.run(runId, null, mapState(state), jsonlId);
|
||||
} catch {
|
||||
/* one bad agent must not abort the live ingest */
|
||||
}
|
||||
}
|
||||
|
||||
// Agents that have a `started` event but no transcript file yet (queued).
|
||||
for (const agentId of started) {
|
||||
if (agentFiles.includes(`agent-${agentId}.jsonl`)) continue;
|
||||
progress.push({
|
||||
type: "workflow_agent",
|
||||
agentId,
|
||||
label: null,
|
||||
phaseTitle: null,
|
||||
model: null,
|
||||
state: doneResults.has(agentId) ? "done" : "running",
|
||||
tokens: 0,
|
||||
toolCalls: 0,
|
||||
durationMs: null,
|
||||
lastToolName: null,
|
||||
});
|
||||
}
|
||||
|
||||
let startedAtIso = earliest ? new Date(earliest).toISOString() : null;
|
||||
if (!startedAtIso && scriptPath) {
|
||||
try {
|
||||
startedAtIso = new Date(fs.statSync(scriptPath).mtimeMs).toISOString();
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
}
|
||||
const durationMs = earliest && latest ? latest - earliest : null;
|
||||
|
||||
stmts.upsertWorkflow.run(
|
||||
runId,
|
||||
sessionId,
|
||||
null,
|
||||
scriptPath ? nameFromScript(scriptPath) : runId,
|
||||
"running",
|
||||
model,
|
||||
startedAtIso,
|
||||
null,
|
||||
durationMs,
|
||||
progress.length,
|
||||
totalTokens,
|
||||
totalToolCalls,
|
||||
null,
|
||||
JSON.stringify(progress),
|
||||
scriptPath || null,
|
||||
null,
|
||||
"live"
|
||||
);
|
||||
return { row: stmts.getWorkflow.get(runId), tokens: runTokens };
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect running workflows: a launch script whose journal hasn't landed yet.
|
||||
* Upsert a minimal `running` row so the UI shows it before completion. Skips
|
||||
* runs that already have a completed/error row (the journal won.) Returns the
|
||||
* upserted rows.
|
||||
*/
|
||||
function detectRunningWorkflows(dbModule, sessionId, paths, handledRunIds) {
|
||||
const { stmts } = dbModule;
|
||||
const changed = [];
|
||||
for (const scriptPath of paths.scripts) {
|
||||
const runId = extractRunId(scriptPath);
|
||||
if (!runId || handledRunIds.has(runId)) continue;
|
||||
const existing = stmts.getWorkflow.get(runId);
|
||||
if (existing && existing.status !== "running") continue; // journal already won
|
||||
|
||||
let startedAt = null;
|
||||
let agentCount = 0;
|
||||
try {
|
||||
const st = fs.statSync(scriptPath);
|
||||
startedAt = new Date(st.mtimeMs).toISOString();
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
// Best-effort fleet size: inner-agent transcripts in this run's nested dir.
|
||||
try {
|
||||
const agentDir = paths.sessionDir ? agentsDirForRun(paths.sessionDir, runId) : null;
|
||||
if (agentDir && fs.existsSync(agentDir)) {
|
||||
agentCount = fs
|
||||
.readdirSync(agentDir)
|
||||
.filter((f) => f.startsWith("agent-") && f.endsWith(".jsonl")).length;
|
||||
}
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
|
||||
stmts.upsertWorkflow.run(
|
||||
runId,
|
||||
sessionId,
|
||||
null,
|
||||
nameFromScript(scriptPath),
|
||||
"running",
|
||||
null,
|
||||
startedAt,
|
||||
null,
|
||||
null,
|
||||
agentCount,
|
||||
0,
|
||||
0,
|
||||
null,
|
||||
null,
|
||||
scriptPath,
|
||||
null,
|
||||
"live"
|
||||
);
|
||||
changed.push(stmts.getWorkflow.get(runId));
|
||||
}
|
||||
return changed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Ingest every workflow artifact for one session: completed journals first,
|
||||
* then running detection for journal-less launch scripts.
|
||||
*
|
||||
* @param {object} dbModule - { db, stmts }
|
||||
* @param {{id: string, transcript_path?: string, cwd?: string}} session
|
||||
* @returns {Promise<object[]>} the workflow rows that were inserted/updated
|
||||
*/
|
||||
async function ingestWorkflowsForSession(dbModule, session) {
|
||||
const sessionId = session && session.id;
|
||||
if (!sessionId) return [];
|
||||
const transcriptPath = resolveTranscriptPath(session);
|
||||
if (!transcriptPath) return [];
|
||||
|
||||
const paths = findSessionWorkflows(transcriptPath);
|
||||
if (paths.journals.length === 0 && paths.scripts.length === 0 && paths.liveRuns.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const changed = [];
|
||||
const journalRunIds = new Set();
|
||||
// Session-wide accumulator of inner-agent token usage across all runs, so the
|
||||
// session's cost includes workflow spend. Recomputed in full each call (all
|
||||
// journals are re-parsed) → writeSessionTokens replace semantics make it
|
||||
// idempotent (no double-count across re-ingests).
|
||||
const workflowTokens = {};
|
||||
// Map runId → its launch script (so a journal row records script_path too).
|
||||
const scriptByRun = new Map();
|
||||
for (const s of paths.scripts) scriptByRun.set(extractRunId(s), s);
|
||||
|
||||
for (const journalPath of paths.journals) {
|
||||
try {
|
||||
const journal = parseWorkflowJournal(journalPath);
|
||||
if (!journal) continue;
|
||||
journalRunIds.add(journal.runId);
|
||||
const res = await ingestWorkflowJournal(dbModule, sessionId, journal, {
|
||||
sessionDir: paths.sessionDir,
|
||||
scriptPath: scriptByRun.get(journal.runId) || null,
|
||||
});
|
||||
if (res && res.row) changed.push(res.row);
|
||||
if (res && res.tokens) mergeWorkflowTokens(workflowTokens, res.tokens);
|
||||
} catch {
|
||||
/* skip malformed journal */
|
||||
}
|
||||
}
|
||||
|
||||
// Live runs (no terminal journal yet): build real-time progress + tokens from
|
||||
// the streaming journal.jsonl + growing agent transcripts.
|
||||
const liveHandled = new Set();
|
||||
for (const lr of paths.liveRuns) {
|
||||
if (journalRunIds.has(lr.runId)) continue; // terminal journal is authoritative
|
||||
try {
|
||||
const res = await ingestLiveWorkflow(
|
||||
dbModule,
|
||||
sessionId,
|
||||
paths.sessionDir,
|
||||
lr.runId,
|
||||
scriptByRun.get(lr.runId) || null
|
||||
);
|
||||
if (res && res.row) {
|
||||
changed.push(res.row);
|
||||
liveHandled.add(lr.runId);
|
||||
}
|
||||
if (res && res.tokens) mergeWorkflowTokens(workflowTokens, res.tokens);
|
||||
} catch {
|
||||
/* non-fatal — partial live run */
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
const handled = new Set([...journalRunIds, ...liveHandled]);
|
||||
changed.push(...detectRunningWorkflows(dbModule, sessionId, paths, handled));
|
||||
} catch {
|
||||
/* non-fatal */
|
||||
}
|
||||
|
||||
// Fold the workflow fleet's token usage into the session cost under a
|
||||
// namespaced `workflow` service_tier (isolated from the main-transcript
|
||||
// writer's buckets). getTokensBySession + calculateCost sum it per model.
|
||||
try {
|
||||
if (Object.keys(workflowTokens).length > 0) {
|
||||
importHistory().writeSessionTokens(dbModule, sessionId, workflowTokens);
|
||||
}
|
||||
} catch {
|
||||
/* non-fatal — cost folding must never break ingestion */
|
||||
}
|
||||
|
||||
return changed;
|
||||
}
|
||||
|
||||
/**
|
||||
* One-time backfill: ingest workflow artifacts for every recorded session.
|
||||
* Used by the legacy auto-import on first boot so historical completed
|
||||
* workflows surface. Idempotent and fail-safe per session.
|
||||
*
|
||||
* @returns {Promise<{sessions: number, workflows: number}>}
|
||||
*/
|
||||
async function ingestAllWorkflows(dbModule) {
|
||||
const { db } = dbModule;
|
||||
let rows = [];
|
||||
try {
|
||||
rows = db.prepare("SELECT id, cwd, transcript_path FROM sessions").all();
|
||||
} catch {
|
||||
return { sessions: 0, workflows: 0 };
|
||||
}
|
||||
let sessions = 0;
|
||||
let workflows = 0;
|
||||
for (const row of rows) {
|
||||
try {
|
||||
const changed = await ingestWorkflowsForSession(dbModule, {
|
||||
id: row.id,
|
||||
cwd: row.cwd,
|
||||
transcript_path: row.transcript_path,
|
||||
});
|
||||
if (changed.length > 0) {
|
||||
sessions++;
|
||||
workflows += changed.length;
|
||||
}
|
||||
} catch {
|
||||
/* non-fatal — skip this session */
|
||||
}
|
||||
}
|
||||
return { sessions, workflows };
|
||||
}
|
||||
|
||||
/**
|
||||
* Cheap change-fingerprint for a session's workflow artifacts: the newest mtime
|
||||
* across its journals, launch scripts, and — crucially for real-time — the
|
||||
* streaming files of any RUNNING run (journal.jsonl + agent-*.jsonl), so the
|
||||
* poll re-ingests as a live workflow's tokens/agents grow. Per-file statting is
|
||||
* bounded to runs without a terminal journal; completed runs contribute only
|
||||
* their (stable) terminal-journal mtime. Returns 0 when nothing exists.
|
||||
*/
|
||||
function workflowsMaxMtime(transcriptPath) {
|
||||
const { journals, scripts, liveRuns } = findSessionWorkflows(transcriptPath);
|
||||
let max = 0;
|
||||
const stat = (p) => {
|
||||
try {
|
||||
const m = fs.statSync(p).mtimeMs;
|
||||
if (m > max) max = m;
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
};
|
||||
for (const p of [...journals, ...scripts]) stat(p);
|
||||
const completed = new Set(journals.map(extractRunId));
|
||||
for (const lr of liveRuns) {
|
||||
if (completed.has(lr.runId)) continue; // terminal journal mtime already counted
|
||||
try {
|
||||
for (const f of fs.readdirSync(lr.dir)) {
|
||||
if (f.endsWith(".jsonl")) stat(path.join(lr.dir, f));
|
||||
}
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
}
|
||||
return max;
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
ingestWorkflowsForSession,
|
||||
ingestAllWorkflows,
|
||||
ingestLiveWorkflow,
|
||||
workflowsMaxMtime,
|
||||
findSessionWorkflows,
|
||||
parseWorkflowJournal,
|
||||
ingestWorkflowJournal,
|
||||
detectRunningWorkflows,
|
||||
extractRunId,
|
||||
nameFromScript,
|
||||
mapState,
|
||||
};
|
||||
Reference in New Issue
Block a user