import { loadConfig } from "../../config/config.ts"; import { applyAttentionState, formatActivityAge, resolveCrewControlConfig } from "../../runtime/agent-control.ts"; import { extractCommandTrace } from "../../runtime/command-trace.ts"; import { readCrewAgents } from "../../runtime/crew-agent-records.ts"; import { deadletterStatusLine } from "../../runtime/deadletter.ts"; import { evaluateRunEffectiveness } from "../../runtime/effectiveness.ts"; import { computePhaseProgress } from "../../runtime/phase-progress.ts"; import { checkProcessLiveness, isActiveRunStatus } from "../../runtime/process-status.ts"; import { formatTaskGraphLines, waitingReason } from "../../runtime/task-display.ts"; import { verifyTaskCompletion } from "../../runtime/verification/completion-guard.ts"; import type { TeamToolParamsValue } from "../../schema/team-tool-schema.ts"; import { withRunLockSync } from "../../state/coordination/locks.ts"; import { readDeliveryState, readMailbox } from "../../state/coordination/mailbox.ts"; import { appendEvent, readEventsCursor } from "../../state/event-log/event-log.ts"; import { loadRunManifestById, saveRunTasks, updateRunStatus } from "../../state/stores/state-store.ts"; import { aggregateUsage, formatCost, formatUsage } from "../../state/usage.ts"; import { formatDuration } from "../../ui/format-helpers.ts"; import { goalAchievedStatusLabel } from "../../ui/goal-flag.ts"; import { locateRunCwd } from "../team-tool.ts"; import type { PiTeamsToolResult } from "../tool-result.ts"; import { result, type TeamContext } from "./context.ts"; import { paramRequired } from "./param-error.ts"; import { RUN_NOT_FOUND_HINT } from "./run-not-found.ts"; // H5 (2026-08-10): idempotency guard for the stale-async side-effect below. // Without this, every status poll of a dead async run (dashboard refresh, // status-bar tick, broker RPC) re-runs saveRunTasks + appendEvent — each // costing ~30ms of blocking I/O on what callers treat as a read-only path. // The runId is added the first time we transition it; it's never removed // (a real resurrection goes through a fresh runId, not a status query). // Bounded FIFO to avoid unbounded growth from long-lived sessions. const STALE_ASYNC_MARKED = new Set(); const STALE_ASYNC_MARKED_MAX = 256; function markStaleAsync(runId: string): boolean { if (STALE_ASYNC_MARKED.has(runId)) return false; STALE_ASYNC_MARKED.add(runId); if (STALE_ASYNC_MARKED.size > STALE_ASYNC_MARKED_MAX) { const oldest = STALE_ASYNC_MARKED.keys().next().value; if (oldest !== undefined) STALE_ASYNC_MARKED.delete(oldest); } return true; } /** * R14-1 (Phase 3.4): the ONLY state mutation on the status read path — the * stale-async transition — executed under the run lock with a FRESH re-read * (respond.ts:43 pattern). handleStatus's `loaded` snapshot is best-effort * (no lock): a concurrent writer can commit a terminal status between that * load and this lock, and `updateRunStatus`'s transition guard would validate * against the STALE status (running→failed legal) and clobber the fresh * terminal state. Inside the lock we re-read and short-circuit when the fresh * state no longer warrants the transition (run gone, async entry gone, status * no longer active, process alive). * * Returns the transitioned manifest+tasks so the caller's read view reflects * the write (same-poll visibility, matching pre-fix behavior), or undefined * when the transition was skipped. * * Exported for unit testing (R14-1 regression: a run that completed on disk * between load and lock must NOT be re-flipped to failed). */ export function transitionStaleAsyncUnderLock( loaded: NonNullable>, runCwd: string, runId: string, ): | { manifest: NonNullable>["manifest"]; tasks: NonNullable>["tasks"]; } | undefined { let transitioned: | { manifest: NonNullable>["manifest"]; tasks: NonNullable>["tasks"]; } | undefined; withRunLockSync(loaded.manifest, () => { const fresh = loadRunManifestById(runCwd, runId); // NOTE: inside withRunLockSync - consistent read if (!fresh?.manifest.async || !isActiveRunStatus(fresh.manifest.status)) return; const freshLiveness = checkProcessLiveness(fresh.manifest.async.pid); if (freshLiveness.alive) return; const failed = updateRunStatus(fresh.manifest, "failed", `Async process stale: ${freshLiveness.detail}`); const freshTasks = fresh.tasks.map((task) => task.status === "running" ? { ...task, status: "cancelled" as const, finishedAt: new Date().toISOString(), error: "Async process died; task was not completed.", } : task, ); saveRunTasks(failed, freshTasks); // async.stale (2026-08-10): sync append (byte-identical to pre-extract // api.ts) so callers reading eventsPath immediately after status see // the event — async fire-and-forget would race. // REVIEW FIX (2026-09-10): reverted the M2b buffered conversion — the // same-poll read-back contract below requires same-tick durability. appendEvent(failed.eventsPath, { type: "async.stale", runId: failed.runId, message: freshLiveness.detail, data: { pid: fresh.manifest.async.pid }, }); transitioned = { manifest: failed, tasks: freshTasks }; }); return transitioned; } export function handleStatus(params: TeamToolParamsValue, ctx: TeamContext): PiTeamsToolResult { if (!params.runId) return result( paramRequired("status", "runId", "{ action: 'status', runId: 'team_...' }"), { action: "status", status: "error" }, true, ); const runCwd = locateRunCwd(params.runId, ctx.cwd); if (!runCwd) return result(`Run '${params.runId}' not found.${RUN_NOT_FOUND_HINT}`, { action: "status", status: "error" }, true); const loaded = loadRunManifestById(runCwd, params.runId); // NOTE: no withRunLock - best-effort only; concurrent writes may cause inconsistency if (!loaded) return result(`Run '${params.runId}' not found.${RUN_NOT_FOUND_HINT}`, { action: "status", status: "error" }, true); let { manifest, tasks } = loaded; // DX (Round 16 F3): compact status mode. Default = full (backward compatible). // details=false gives a tight summary (status, goal, counts, failed/attention // errors) for quick checks without 40 lines of dense key=value noise. const fullDetails = params.details !== false; let asyncLivenessLine: string | undefined; if (manifest.async) { const asyncState = manifest.async; const liveness = checkProcessLiveness(asyncState.pid); asyncLivenessLine = `Async: pid=${asyncState.pid ?? "unknown"} alive=${liveness.alive ? "true" : "false"} detail=${liveness.detail} log=${asyncState.logPath} spawnedAt=${asyncState.spawnedAt}`; if (!liveness.alive && isActiveRunStatus(manifest.status)) { // R14-1 (Phase 3.4): ownership gate — a foreign session must NOT // trigger the stale-async transition on a read path; the owning // session's poll handles it. Mirrors handleResume (team-tool.ts R1). // Backward compat: only applies when ownerSessionId is a non-empty // string — legacy runs without an owner are not gated. const foreignRun = typeof manifest.ownerSessionId === "string" && manifest.ownerSessionId !== ctx.sessionId; if (!foreignRun || params.force) { // H5: only transition once per runId. Repeated status polls of a // dead async run previously re-ran saveRunTasks + sync appendEvent // on every call (~30ms blocking each). if (markStaleAsync(manifest.runId)) { const transitioned = transitionStaleAsyncUnderLock(loaded, runCwd, params.runId); if (transitioned) { manifest = transitioned.manifest; tasks = transitioned.tasks; } } } } } const counts = new Map(); for (const task of tasks) counts.set(task.status, (counts.get(task.status) ?? 0) + 1); const deadletterLine = deadletterStatusLine(manifest); const phaseProgress = computePhaseProgress(tasks); // PERF (2026-08-24): intentionally NOT passing `limit` here — readEventsCursor's // limit is a HEAD cap (oldest-first slice for streaming pagination), not a tail // window, so it would hide recent events from the ack-timeout dedupe below and // re-append duplicate ack_timeout events on every poll. The manifest carries no // last-seq anchor for a sinceSeq-based tail either. The reader is already bounded // internally (4 MB / 5000-event tail), and the downstream filters // (ackTimeoutRequestIds, attentionByTask) intentionally operate on that recent window. const { events: allEvents } = readEventsCursor(manifest.eventsPath); const events = allEvents.slice(-8); // P1-8: pre-build the ack-timeout requestId set once (was O(events × messages) // via allEvents.some() inside the mailbox loop). const ackTimeoutRequestIds = new Set( allEvents.filter((event) => event.type === "agent.group_join.ack_timeout").map((event) => String(event.data?.requestId ?? "")), ); const attentionByTask = new Map( allEvents.filter((event) => event.type === "task.attention" && event.taskId).map((event) => [event.taskId!, event]), ); // H5 (2026-08-10): hoist the config load. Previously three separate // loadConfig(ctx.cwd) calls at lines 77, 85, 119 each stat'd up to 4 files; // status is called per dashboard refresh / RPC poll, so this tripled the // stat cost on every tick. loadConfig is mtime-cached internally so the // three calls returned the same object, but each still paid the cache // lookup + stat syscall. One call, reused everywhere. const cfg = loadConfig(ctx.cwd).config; const controlConfig = resolveCrewControlConfig(cfg); const crewAgents = readCrewAgents(manifest).map((agent) => applyAttentionState(manifest, agent, controlConfig)); const artifactLines = manifest.artifacts .slice(-10) .map( (artifact) => `- ${artifact.kind}: ${artifact.path}${artifact.sizeBytes !== undefined ? ` (${artifact.sizeBytes} bytes)` : ""}`, ); const deliveryState = readDeliveryState(manifest); const ackTimeoutMs = cfg.runtime?.groupJoinAckTimeoutMs; const groupJoinLines: string[] = []; for (const message of readMailbox(manifest, "outbox") .filter((m) => m.data?.kind === "group_join") .slice(-5)) { const ack = deliveryState.messages[message.id] === "acknowledged" ? "acknowledged" : "pending"; const ageMs = Date.now() - new Date(message.createdAt).getTime(); const requestId = String(message.data?.requestId ?? "unknown"); const timedOut = ack === "pending" && ackTimeoutMs !== undefined && Number.isFinite(ageMs) && ageMs > ackTimeoutMs; if (timedOut && !ackTimeoutRequestIds.has(requestId)) { // ack_timeout: sync append (byte-identical to pre-extract api.ts). // REVIEW FIX (2026-09-10): reverted the M2b buffered conversion — the // dedupe set below is rebuilt from readEventsCursor in THIS invocation, // so a buffered write would re-emit duplicates on sub-20ms re-polls. appendEvent(manifest.eventsPath, { type: "agent.group_join.ack_timeout", runId: manifest.runId, message: "Group join delivery ack timed out; mailbox delivery remains the fallback.", data: { requestId, messageId: message.id, batchId: message.data?.batchId, partial: message.data?.partial, ageMs, ackTimeoutMs, }, }); } groupJoinLines.push( `- ${String(message.data?.partial) === "true" ? "partial" : "completed"} request=${requestId} message=${message.id} ack=${timedOut ? "timeout" : ack}`, ); } const totalUsage = aggregateUsage(tasks); const completedTasks = tasks.filter((task) => task.status === "completed"); const effectiveness = evaluateRunEffectiveness({ manifest, tasks, executeWorkers: manifest.runtimeResolution?.kind !== "scaffold", runtimeConfig: cfg.runtime, }); const noObservedWorkTasks = effectiveness.noObservedWorkTaskIds .map((id) => tasks.find((task) => task.id === id)) .filter((task): task is (typeof tasks)[number] => task !== undefined); const attentionTasks = effectiveness.needsAttentionTaskIds .map((id) => tasks.find((task) => task.id === id)) .filter((task): task is (typeof tasks)[number] => task !== undefined); const activeAgents = crewAgents.filter((agent) => agent.status === "running"); const completedAgents = crewAgents.filter((agent) => agent.status !== "running"); const waitingTasks = tasks.filter((task) => task.status === "queued" || task.status === "waiting"); const agentLine = (agent: (typeof crewAgents)[number]): string => `- ${agent.id} [${agent.status}] ${agent.role} -> ${agent.agent} runtime=${agent.runtime}${agent.model ? ` model=${agent.model}` : ""}${agent.usage ? ` usage=${formatUsage(agent.usage)}` : ""}${agent.usage?.cost ? ` cost=${formatCost(agent.usage.cost)}` : ""}${agent.progress?.activityState ? ` activityState=${agent.progress.activityState}` : ""}${formatActivityAge(agent) ? ` activity=${formatActivityAge(agent)}` : ""}${agent.progress?.currentTool ? ` tool=${agent.progress.currentTool}` : ""}${agent.toolUses ? ` tools=${agent.toolUses}` : ""}${!agent.usage && agent.progress?.tokens ? ` tokens=${agent.progress.tokens}` : ""}${agent.progress?.turns ? ` turns=${agent.progress.turns}` : ""}${agent.jsonEvents !== undefined ? ` jsonEvents=${agent.jsonEvents}` : ""}${agent.outputPath ? ` output=${agent.outputPath}` : ""}${agent.transcriptPath ? ` transcript=${agent.transcriptPath}` : ""}${agent.statusPath ? ` status=${agent.statusPath}` : ""}${agent.error ? ` error=${agent.error}` : ""}`; // G19 (W-E Phase 1): surface the goal-achievement verdict on terminal runs — // ⚠ false-green warning line; silent while running, when achieved, or when // the run predates the assessment (goalAchieved undefined). const goalLine = goalAchievedStatusLabel(manifest); const lines = [ `Run: ${manifest.runId}`, `Team: ${manifest.team}`, `Workflow: ${manifest.workflow ?? "(none)"}`, `Status: ${manifest.status}`, ...(goalLine ? [goalLine] : []), `Progress: ${phaseProgress.overallPercentage}% (~${formatDuration(phaseProgress.estimatedRemainingMs)} remaining)`, `Workspace mode: ${manifest.workspaceMode}`, ...(manifest.runtimeResolution ? [ `Runtime: ${manifest.runtimeResolution.kind}`, `Runtime safety: ${manifest.runtimeResolution.safety}`, `Runtime requested: ${manifest.runtimeResolution.requestedMode}${manifest.runtimeResolution.reason ? ` (${manifest.runtimeResolution.reason})` : ""}`, ] : []), `Goal: ${manifest.goal}`, `Created: ${manifest.createdAt}`, `Updated: ${manifest.updatedAt}`, `State: ${manifest.stateRoot}`, `Artifacts: ${manifest.artifactsRoot}`, ...(asyncLivenessLine ? [asyncLivenessLine] : []), "Task graph:", ...formatTaskGraphLines(tasks), "Tasks:", ...(tasks.length ? tasks.map( (task) => `- ${task.id} [${task.status}] ${task.role} -> ${task.agent}${task.taskPacket ? ` scope=${task.taskPacket.scope}` : ""}${task.verification ? ` green=${task.verification.observedGreenLevel}/${task.verification.requiredGreenLevel}` : ""}${task.modelAttempts?.length ? ` attempts=${task.modelAttempts.length}` : ""}${task.modelRouting ? ` modelRouting=${task.modelRouting.requested ? `${task.modelRouting.requested}->` : ""}${task.modelRouting.resolved}${task.modelRouting.usedAttempt ? ` attempt=${task.modelRouting.usedAttempt + 1}` : ""}` : ""}${task.agentProgress?.activityState ? ` activityState=${task.agentProgress.activityState}` : ""}${(() => { const t = extractCommandTrace(task.agentProgress?.recentTools); return t.summary ? ` ${t.summary}` : ""; })()}${attentionByTask.get(task.id)?.data?.reason ? ` attention=${String(attentionByTask.get(task.id)?.data?.reason)}` : ""}${task.jsonEvents !== undefined ? ` jsonEvents=${task.jsonEvents}` : ""}${task.usage ? ` usage=${JSON.stringify(task.usage)}` : ""}${task.resultArtifact ? ` result=${task.resultArtifact.path}` : ""}${task.transcriptArtifact ? ` transcript=${task.transcriptArtifact.path}` : ""}${task.worktree ? ` worktree=${task.worktree.path}` : ""}${task.error ? ` error=${task.error}` : ""}${task.failureCause ? ` failureCause=${task.failureCause}` : ""}`, ) : ["- (none)"]), `Task counts: ${[...counts.entries()].map(([status, count]) => `${status}=${count}`).join(", ") || "none"}`, // US-003 AC-3: surface exhausted-retry failures — only when entries exist. ...(deadletterLine ? [deadletterLine] : []), "Effectiveness:", `- observable=${effectiveness.observable}/${Math.max(1, effectiveness.completed)} completed tasks`, `- workerExecution=${effectiveness.workerExecution} guard=${effectiveness.guardMode} severity=${effectiveness.severity}`, `- noObservedWork=${effectiveness.noObservedWorkTaskIds.length ? effectiveness.noObservedWorkTaskIds.join(",") : "none"}`, `- needsAttention=${effectiveness.needsAttentionTaskIds.length ? effectiveness.needsAttentionTaskIds.join(",") : "none"}`, "Completion verification", ...(tasks.filter((t) => t.status === "completed").length ? tasks .filter((t) => t.status === "completed") .map((t) => { const guard = verifyTaskCompletion(t, manifest); return `- ${t.id} green=${guard.greenLevel}/3${guard.warnings.length ? ` warnings=[${guard.warnings.join(", ")}]` : ""}`; }) : ["- (no completed tasks)"]), "Active agents:", ...(activeAgents.length ? activeAgents.map(agentLine) : ["- (none)"]), "Waiting tasks:", ...(waitingTasks.length ? waitingTasks.map((task) => `- ${task.id} [queued] ${task.role} -> ${task.agent} ${waitingReason(task, tasks) ?? "waiting"}`) : ["- (none)"]), "Completed agents:", ...(completedAgents.length ? completedAgents.map(agentLine) : ["- (none)"]), "Policy decisions:", ...(manifest.policyDecisions?.length ? manifest.policyDecisions.map( (item) => `- ${item.action} (${item.reason})${item.taskId ? ` ${item.taskId}` : ""}: ${item.message}`, ) : ["- (none)"]), `Total usage: ${formatUsage(totalUsage)}`, "Group joins:", ...(groupJoinLines.length ? groupJoinLines : ["- (none)"]), "", "Recent artifacts:", ...(artifactLines.length ? artifactLines : ["- (none)"]), "", "Recent events:", ...(events.length ? events.map( (event) => `- ${event.time} ${event.type}${event.taskId ? ` ${event.taskId}` : ""}${event.message ? `: ${event.message}` : ""}`, ) : ["- (none)"]), ]; if (!fullDetails) { return result(buildCompactStatus(manifest, tasks, counts, asyncLivenessLine, phaseProgress).join("\n"), { action: "status", status: "ok", runId: manifest.runId, artifactsRoot: manifest.artifactsRoot, intent: `status ${manifest.runId}: ${manifest.status} (compact)`, }); } return result(lines.join("\n"), { action: "status", status: "ok", runId: manifest.runId, artifactsRoot: manifest.artifactsRoot, intent: `status ${manifest.runId}: ${manifest.status}`, }); } /** * Compact status builder (DX: Round 16 F3). A tight summary for quick checks: * identity, status, goal, task counts, and ONLY failed / attention task * errors — not the 40-line dense dump. Invoked when params.details === false. * * Exported for unit testing. */ export function buildCompactStatus( manifest: { runId: string; team: string; workflow?: string; status: string; goal: string; workspaceMode?: string; goalAchieved?: boolean | "unknown"; goalAchievementNote?: string; }, tasks: Array<{ id: string; status: string; role: string; agent: string; error?: string; }>, counts: Map, asyncLivenessLine?: string, progress?: { overallPercentage: number; estimatedRemainingMs: number }, ): string[] { const failedOrAttention = tasks.filter((t) => t.status === "failed" || t.status === "needs_attention" || t.status === "cancelled"); const goalLine = goalAchievedStatusLabel(manifest); const lines = [ `Run: ${manifest.runId}`, `Team: ${manifest.team}${manifest.workflow ? ` (${manifest.workflow})` : ""}`, `Status: ${manifest.status}`, ...(goalLine ? [goalLine] : []), ...(progress ? [`Progress: ${progress.overallPercentage}% (~${formatDuration(progress.estimatedRemainingMs)} remaining)`] : []), `Goal: ${manifest.goal}`, ...(asyncLivenessLine ? [asyncLivenessLine] : []), `Tasks: ${[...counts.entries()].map(([s, c]) => `${s}=${c}`).join(", ") || "none"}`, ]; if (failedOrAttention.length > 0) { lines.push("Issues:"); for (const t of failedOrAttention) { lines.push(`- ${t.id} [${t.status}] ${t.role}: ${t.error ?? "(no error detail)"}`); } } lines.push("Tip: pass details=true for full output (task graph, agents, effectiveness, events)."); return lines; }