Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 26 additions & 8 deletions client/src/client.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -678,21 +678,34 @@ function claudeStreamProgress(lines) {

const inactiveAgentStatuses = new Set(["complete", "done", "completed", "closed", "cancelled", "canceled", "failed", "released", "skipped", "finalized", "killed", "missing"]);

function paneOutcomeForEvent(event) {
export function paneOutcomeForEvent(event) {
const type = String(event?.type || "");
if (type === "user_message") return { status: "", summary: "" };
if (type === "session_started") return { status: "running", summary: "" };
if (type === "assistant_message") return { status: "complete", summary: compactOutcomeSummary(event) };
if (type === "assistant_message") {
const summary = compactOutcomeSummary(event);
const status = /^(?:blocker\s*:|.*\bblocked\b)/i.test(summary) ? "blocked" : "complete";
return { status, summary };
}
if (type === "question") return { status: "waiting", summary: compactOutcomeSummary(event) };
if (type === "session_interrupted") return { status: "interrupted", summary: compactOutcomeSummary(event) };
if (type === "progress") return { status: "working", summary: compactOutcomeSummary(event) };
if (type === "session_completed") return { status: "complete", summary: undefined };
return null;
}

function compactOutcomeSummary(event) {
const entries = String(event?.payload?.text || event?.payload?.report || "")
.split(/\r?\n/)
export function finalAgentMessageText(value) {
const sources = String(value || "").split(/\r?\n/);
const finalMessageHeading = sources.findIndex((source) => /^\s*#{1,6}\s+final agent message\s*$/i.test(source));
if (finalMessageHeading < 0) return String(value || "").trim();
const finalSection = sources.slice(finalMessageHeading + 1);
const traceHeading = finalSection.findIndex((source) => /^\s*#{1,6}\s+trace references\s*$/i.test(source));
return (traceHeading >= 0 ? finalSection.slice(0, traceHeading) : finalSection).join("\n").trim();
}

export function compactOutcomeSummary(event) {
const sources = finalAgentMessageText(event?.payload?.text || event?.payload?.report || "").split(/\r?\n/);
const entries = sources
.map((source) => {
const tableRow = /^\s*\|/.test(source) && /\|\s*$/.test(source);
let text = source
Expand All @@ -708,6 +721,8 @@ function compactOutcomeSummary(event) {
.filter(({ text }) => text);
const meaningful = entries.filter(({ text }) =>
!/^(?:result|summary|outcome|answer|final answer)$/i.test(text)
&& !/^(?:status|workflow|session|task)\s*:/i.test(text)
&& !/^trace references$/i.test(text)
&& !/^(?:-+)(?:\s+—\s+-+)*$/.test(text));
const latestIndex = meaningful.findIndex(({ text }) => /\b(?:most recently|latest)\b/i.test(text));
if (latestIndex >= 0) {
Expand All @@ -719,7 +734,7 @@ function compactOutcomeSummary(event) {
}
const lines = meaningful.map(({ text }) => text);
const preferred = lines.find((line) => /\b(?:most recently|latest)\b/i.test(line))
|| lines.find((line) => /^(?:result|answer|outcome)\s*:/i.test(line))
|| lines.find((line) => /^(?:result|answer|outcome|blocker|finding|conclusion)\s*:/i.test(line))
|| lines.find((line) => /\b(?:found|fixed|created|updated|merged|deployed|completed)\b/i.test(line))
|| lines[0]
|| entries[0]?.text
Expand Down Expand Up @@ -774,7 +789,7 @@ export function renderAgentPane(agents, {

function agentStatusGlyph(status) {
const value = String(status || "").toLowerCase();
if (new Set(["failed", "killed", "cancelled", "canceled", "delivery-blocked", "interrupted"]).has(value)) return "×";
if (new Set(["blocked", "failed", "killed", "cancelled", "canceled", "delivery-blocked", "interrupted"]).has(value)) return "×";
if (inactiveAgentStatuses.has(value)) return "✓";
if (new Set(["starting", "queued", "connecting", "restoring", "waiting"]).has(value)) return "◌";
if (new Set(["running", "working", "in-progress", "planning"]).has(value)) return "●";
Expand Down Expand Up @@ -907,7 +922,10 @@ function selectThread(threads, selector) {
}

function renderInteractiveEvent(stdout, event) {
const text = String(event.payload?.text || event.payload?.report || "").trim();
const rawText = String(event.payload?.text || event.payload?.report || "").trim();
const text = new Set(["assistant_message", "question", "session_interrupted"]).has(event.type)
? finalAgentMessageText(rawText)
: rawText;
if (event.type === "user_message") stdout.write(`\nyou> ${text}\n`);
else if (event.type === "assistant_message") stdout.write(`\nassistant> ${text}\n`);
else if (event.type === "question") stdout.write(`\nassistant? ${text}\n`);
Expand Down
62 changes: 57 additions & 5 deletions client/test/client.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,16 @@ import { mkdtemp, readFile, stat, writeFile } from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import test from "node:test";
import { ControlClient, main, renderAgentPane, terminalDelta, terminalProgressView } from "../src/client.mjs";
import {
compactOutcomeSummary,
ControlClient,
finalAgentMessageText,
main,
paneOutcomeForEvent,
renderAgentPane,
terminalDelta,
terminalProgressView,
} from "../src/client.mjs";

function writer() {
return { output: "", write(value) { this.output += String(value); } };
Expand Down Expand Up @@ -537,7 +546,36 @@ test("completed outcome summary wraps within the terminal width", () => {
assert.ok(lines.every((line) => line.length <= 52));
});

test("asynchronous status redraw restores the active input prompt", async () => {
test("interrupted orchestrator pane shows the bounded blocker instead of a generic failure", () => {
const lines = renderAgentPane([], {
columns: 64,
maxRows: 6,
thread: { id: "thread-blocked", state: "interrupted" },
outcomeStatus: "interrupted",
taskSummary: "Grafana read blocked: runbook requests 1.0.0 but prod-mcp certifies 1.1.0.",
});
assert.deepEqual(lines.slice(0, 3), [
"× orchestrator · interrupted",
" ↳ Grafana read blocked: runbook requests 1.0.0 but prod-mcp",
" certifies 1.1.0.",
]);
});

test("structured PR review reports expose the blocker and hide the runtime envelope", () => {
const report = "# session-pr-review\n\nStatus: completed\nWorkflow: run-pr-review\n\n## Final agent message\n# PR #68 Review — Blocked\n\n**Blocker:** GitHub read access does not expose the PR diff, changed files, or CI checks.\n\n## Trace references\n- agents/ops-01/events.jsonl";
assert.equal(finalAgentMessageText(report), "# PR #68 Review — Blocked\n\n**Blocker:** GitHub read access does not expose the PR diff, changed files, or CI checks.");
assert.equal(compactOutcomeSummary({ payload: { text: report } }), "Blocker: GitHub read access does not expose the PR diff, changed files, or CI checks.");
assert.deepEqual(paneOutcomeForEvent({ type: "assistant_message", payload: { text: report } }), {
status: "blocked",
summary: "Blocker: GitHub read access does not expose the PR diff, changed files, or CI checks.",
});
assert.equal(renderAgentPane([], {
outcomeStatus: "blocked",
taskSummary: "Blocker: GitHub read access does not expose the PR diff.",
})[0], "× orchestrator · blocked");
});

test("latest-open-PR interaction ends with an informative summary, completed agent graph, and active prompt", async () => {
const sessionFile = await sessionFixture();
const output = ttyWriter();
let questionCount = 0;
Expand All @@ -559,10 +597,20 @@ test("asynchronous status redraw restores the active input prompt", async () =>
sequence: 1,
type: "assistant_message",
payload: {
text: "# Open PRs — movement-network/aptos-core\n\nChecked current open pull requests.\n\n## Latest opened PR\n\n| PR | Title |\n| --- | --- |\n| **#421** | fix: remove global waypoint signature-verification bypass |",
text: "# thread-latest-open-pr\n\nStatus: completed\nWorkflow: run-latest-open-pr\n\n## Final agent message\nLatest open PR: **#421** fix: remove global waypoint signature-verification bypass\n- Author: contributor\n- URL: https://github.com/movement-network/aptos-core/pull/421\n\n## Trace references\n- agents/ops-01/attempt-0001/events.jsonl",
},
},
})));
threadSocket.emit("message", Buffer.from(JSON.stringify({
type: "agents",
agents: [{
name: "ops-01",
status: "done",
role: "ops",
workingOn: "Found latest open PR #421",
}],
available: true,
})));
setImmediate(() => resolve("/quit"));
});
});
Expand Down Expand Up @@ -595,8 +643,12 @@ test("asynchronous status redraw restores the active input prompt", async () =>
});

assert.match(output.output, /● orchestrator · running/);
assert.match(output.output, /assistant> # Open PRs/);
assert.match(output.output, /assistant> Latest open PR: \*\*#421\*\*/);
assert.match(output.output, /✓ orchestrator · complete/);
assert.match(output.output, /↳ Latest opened PR: #421 — fix: remove global waypoint/);
assert.match(output.output, /↳ Latest open PR: #421 — fix: remove global waypoint/);
assert.match(output.output, /✓ ops-01 · ops · done/);
assert.match(output.output, /↳ Found latest open PR #421/);
assert.doesNotMatch(output.output, /↳ Status: completed/);
assert.doesNotMatch(output.output, /assistant> # thread-latest-open-pr|Status: completed|Workflow: run-latest-open-pr|Trace references/);
assert.ok(prompts.some((prompt) => prompt.label === "› " && prompt.preserveCursor === true));
});
60 changes: 34 additions & 26 deletions control-server/src/server.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import { deliverWorkerReport, reportDeliveryTimeoutMs } from "./worker-report-de
import { issueWorkerToken as createWorkerToken, verifyWorkerAuthorization } from "./worker-token.mjs";
import { readSubagentSnapshot } from "./subagent-status.mjs";
import { renderThreadTask } from "./thread-execution-context.mjs";
import { fetchWorkerSubagents } from "./worker-subagent-client.mjs";
import { fetchWorkerSubagents, fetchWorkerSubagentsWithReconciliation } from "./worker-subagent-client.mjs";
import { configuredRepository, parseRepositoryCatalog } from "./repository-catalog.mjs";
import { visibleLegacySessionIds } from "./session-visibility.mjs";
import {
Expand All @@ -23,13 +23,14 @@ import {
normalizeWorkerReport,
ownsThreadProjection,
responseTypeForMessage,
scopedThreadTranscript,
selectFinalMessage,
sessionControlInvocation,
sessionLaunchInvocation,
shouldAutomaticallyResume,
submitLocalFollowup,
validResourceId,
workerReportInterruptedEvent,
workerReportPublicEvent,
} from "./session-runtime.mjs";

const here = path.dirname(fileURLToPath(import.meta.url));
Expand Down Expand Up @@ -442,10 +443,10 @@ function readLocalWorkerReport(id) {
} catch { return null; }
}

async function deliverCompletedWorkerReport(id) {
async function deliverWorkerOutcomeReport(id) {
if (!workerMode || !workerReportGatewayUrl || !workerReportTokenFile) return;
const report = readLocalWorkerReport(id);
if (!report) throw new Error(`completed session ${id} has no normalized report`);
if (!report) throw new Error(`session ${id} has no normalized outcome report`);
const token = fs.readFileSync(workerReportTokenFile, "utf8").trim();
await deliverWorkerReport({
gatewayUrl: workerReportGatewayUrl,
Expand Down Expand Up @@ -513,12 +514,17 @@ async function gatewaySubagentSnapshot(id) {
if (!record) return unavailableSubagentSnapshot(id);
if (!record.podIP) record = await reconcileGatewaySession(id);
if (!record?.podIP) return unavailableSubagentSnapshot(id);
try {
const agents = await fetchWorkerSubagents({
const result = await fetchWorkerSubagentsWithReconciliation({
record,
fetchSnapshot: (hostname) => fetchWorkerSubagents({
sessionId: id,
hostname: record.podIP,
hostname,
token: issueWorkerToken(id),
});
}),
reconcile: () => reconcileGatewaySession(id),
});
if (Array.isArray(result.agents)) {
const agents = result.agents;
const snapshot = {
sessionId: id,
agents,
Expand All @@ -530,14 +536,14 @@ async function gatewaySubagentSnapshot(id) {
gatewaySubagentSnapshots.set(id, snapshot);
gatewaySubagentSnapshotErrors.delete(id);
return snapshot;
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
if (gatewaySubagentSnapshotErrors.get(id) !== message) {
gatewaySubagentSnapshotErrors.set(id, message);
console.warn("session worker subagent snapshot unavailable", { sessionId: id, error: message });
}
return unavailableSubagentSnapshot(id, error);
}
if (!result.error) return unavailableSubagentSnapshot(id);
const message = result.error instanceof Error ? result.error.message : String(result.error);
if (gatewaySubagentSnapshotErrors.get(id) !== message) {
gatewaySubagentSnapshotErrors.set(id, message);
console.warn("session worker subagent snapshot unavailable", { sessionId: id, error: message });
}
return unavailableSubagentSnapshot(id, result.error);
}

async function reconcileGatewaySession(id) {
Expand Down Expand Up @@ -727,30 +733,28 @@ async function projectSessionToThread(id, status, reportReader = readGatewayRepo
const session = sessions.find((candidate) => candidate.id === id);
if (!session) return;
if (session.inboxAckSequence !== session.inboxHeadSequence) return;
const publicEvent = workerReportPublicEvent(id, report);
await threadStore.appendFencedSessionEvent({
threadId: record.threadId,
sessionId: id,
generation: record.leaseGeneration,
eventId: `final-${id}`,
type: report.responseType,
payload: {
text: report.responseType === "question" && report.message ? report.message : report.report,
transcript: scopedThreadTranscript(id, report.transcript),
},
...publicEvent,
});
await threadStore.markSessionFinishing({ threadId: record.threadId, sessionId: id, generation: record.leaseGeneration });
const finalized = await threadStore.finalizeSession({ threadId: record.threadId, sessionId: id, generation: record.leaseGeneration });
record.threadProjectedAt = new Date().toISOString();
await saveRegistry();
if (finalized.activatedSession) await launchActivatedThreadSession(record, finalized.activatedSession);
} else if (status === "failed" || status === "paused") {
const fallback = status === "paused" ? "Execution session paused" : "Execution session failed";
const publicEvent = workerReportInterruptedEvent(id, reportReader(id), fallback);
await threadStore.appendFencedSessionEvent({
threadId: record.threadId,
sessionId: id,
generation: record.leaseGeneration,
eventId: `interrupted-${id}`,
type: "session_interrupted",
payload: { text: status === "paused" ? "Execution session paused" : "Execution session failed" },
...publicEvent,
});
const finalized = await threadStore.finalizeSession({ threadId: record.threadId, sessionId: id, generation: record.leaseGeneration, status: "interrupted" });
record.threadProjectedAt = new Date().toISOString();
Expand Down Expand Up @@ -1295,7 +1299,7 @@ for (const record of gatewayMode ? [] : Object.values(registry.sessions)) {
writeTraceSummary(record.id, "completed");
saveRegistry();
if (workerMode) {
deliverCompletedWorkerReport(record.id)
deliverWorkerOutcomeReport(record.id)
.catch((error) => console.error(`worker report delivery failed for ${record.id}`, error))
.finally(() => setTimeout(() => process.exit(0), completionGraceMs));
}
Expand All @@ -1314,7 +1318,7 @@ const retirementTimer = setInterval(() => {
if (record.status === "running" && workflowPhase(record.id) === "complete") {
retireSession(record.id, "completed", "workflow-supervisor").then(async () => {
if (workerMode) {
try { await deliverCompletedWorkerReport(record.id); }
try { await deliverWorkerOutcomeReport(record.id); }
catch (error) { console.error(`worker report delivery failed for ${record.id}`, error); }
setTimeout(() => process.exit(0), completionGraceMs);
}
Expand All @@ -1333,8 +1337,12 @@ const retirementTimer = setInterval(() => {
console.error(`automatic resume failed for ${record.id}`, error);
}
}
retireSession(record.id, "failed", "process-exit").then(() => {
if (workerMode) setTimeout(() => process.exit(1), 1000);
retireSession(record.id, "failed", "process-exit").then(async () => {
if (workerMode) {
try { await deliverWorkerOutcomeReport(record.id); }
catch (error) { console.error(`worker outcome report delivery failed for ${record.id}`, error); }
setTimeout(() => process.exit(1), 1000);
}
}).catch((error) => console.error(`failed retirement failed for ${record.id}`, error));
continue;
}
Expand Down
20 changes: 20 additions & 0 deletions control-server/src/session-runtime.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -109,3 +109,23 @@ export function scopedThreadTranscript(sessionId, transcript) {
: [];
return { ...transcript, traceReferences };
}

export function workerReportPublicEvent(sessionId, report) {
return {
type: report.responseType,
payload: {
text: selectFinalMessage(report.message, report.report),
transcript: scopedThreadTranscript(sessionId, report.transcript),
},
};
}

export function workerReportInterruptedEvent(sessionId, report, fallback) {
return {
type: "session_interrupted",
payload: {
text: report ? selectFinalMessage(report.message, fallback) : String(fallback || "").trim(),
transcript: report ? scopedThreadTranscript(sessionId, report.transcript) : null,
},
};
}
3 changes: 2 additions & 1 deletion control-server/src/subagent-status.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,8 @@ function lastProgressLine(text) {
if (/^(?:final status:|Multiagent launch mode:)/i.test(line)) return false;
if (/[{,]\s*\\?"(?:type|session_id|uuid|usage|duration_ms)\\?"\s*:/.test(line)) return false;
try { if (typeof JSON.parse(line) === "object") return false; } catch {}
return true;
if (/[{}\[\]`]|\\[nrt"]|"\s*:|\bsignature\b/i.test(line)) return false;
return /^(?:Analyzing|Checking|Collecting|Comparing|Executing|Finding|Found|Inspecting|Investigating|Preparing|Querying|Reading|Reviewing|Running|Summarizing|Tracing|Validating|Waiting|Working)\b/.test(line);
}) || "";
}

Expand Down
29 changes: 29 additions & 0 deletions control-server/src/worker-subagent-client.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -38,3 +38,32 @@ export function fetchWorkerSubagents({
request.end();
});
}

export async function fetchWorkerSubagentsWithReconciliation({
record,
fetchSnapshot,
reconcile,
}) {
const initialPodIP = record?.podIP || null;
try {
return { agents: await fetchSnapshot(initialPodIP), record, error: null };
} catch (initialError) {
let refreshed;
try {
refreshed = await reconcile();
} catch {
return { agents: null, record, error: initialError };
}
if (refreshed?.status !== "running" || !refreshed.podIP) {
return { agents: null, record: refreshed, error: null };
}
if (refreshed.podIP === initialPodIP) {
return { agents: null, record: refreshed, error: initialError };
}
try {
return { agents: await fetchSnapshot(refreshed.podIP), record: refreshed, error: null };
} catch (retryError) {
return { agents: null, record: refreshed, error: retryError };
}
}
}
Loading
Loading