import * as fs from "node:fs"; import * as path from "node:path"; import type { AgentConfig } from "../agents/agent-config.ts"; import { allAgents, discoverAgents, listDynamicAgents, registerDynamicAgent, unregisterDynamicAgent } from "../agents/discover-agents.ts"; import { loadConfig } from "../config/config.ts"; import { DEFAULT_PATHS } from "../config/defaults.ts"; import { getCrewEnv } from "../config/env-vars.ts"; // Heavy runtime — lazy-loaded to avoid 1.4s import cost at extension registration. // executeTeamRun is only called when a team run actually executes. import type { executeTeamRun as _executeTeamRunFn } from "../runtime/team-runner.ts"; import type { TeamToolParamsValue } from "../schema/team-tool-schema.ts"; import { TEAM_TERMINAL_TASK_STATUSES } from "../state/contracts.ts"; import { withRunLock } from "../state/coordination/locks.ts"; import { replayPendingMailboxMessages } from "../state/coordination/mailbox.ts"; import { isTaskClaimExpired } from "../state/coordination/task-claims.ts"; import { appendEventAsync, appendEventFireAndForget } from "../state/event-log/event-log.ts"; import { ACTIVE_RUN_STALE_MS, activeRunEntries, registerActiveRun, unregisterActiveRun } from "../state/stores/active-run-registry.ts"; import { writeArtifact } from "../state/stores/artifact-store.ts"; import { loadRunManifestById, saveRunManifestAsync, saveRunTasks, updateRunStatus } from "../state/stores/state-store.ts"; import type { ArtifactDescriptor, TeamRunManifest, TeamTaskState } from "../state/types.ts"; import { allTeams, discoverTeams } from "../teams/discover-teams.ts"; import type { TeamConfig } from "../teams/team-config.ts"; import { logInternalError } from "../utils/internal-error.ts"; import { findRepoRoot, projectCrewRoot, userCrewRoot } from "../utils/paths.ts"; import { resolveRealContainedPath } from "../utils/safe-paths.ts"; import { allWorkflows, discoverWorkflows } from "../workflows/discover-workflows.ts"; import { listRecentRuns } from "./run-index.ts"; import type { PiTeamsToolResult } from "./tool-result.ts"; type ExecuteTeamRunFn = typeof _executeTeamRunFn; async function executeTeamRun(...args: Parameters): Promise>> { // LAZY: heavy runtime — defer 1.4s import cost until team run actually executes. const mod = await import("../runtime/team-runner.ts"); return mod.executeTeamRun(...args); } import { directTeamAndWorkflowFromRun } from "../runtime/direct-run.ts"; import { resolveCrewRuntime, runtimeResolutionState } from "../runtime/model/runtime-resolver.ts"; import { resolveParentModel } from "../runtime/model/session-model.ts"; import { parsePiJsonOutput } from "../runtime/output/pi-json-output.ts"; import { effectiveRunConfig } from "./team-tool/config-patch.ts"; import { buildParentContext, formatScoped, result, type TeamContext } from "./team-tool/context.ts"; // Lazy-loaded: run.ts pulls in spawnBackgroundTeamRun, resolveCrewRuntime, etc. // Static import fails silently in some jiti contexts (child-process), leaving handleRun undefined. import type { handleRun as _handleRunFn } from "./team-tool/run.ts"; import { resolveRunDeadline } from "./team-tool/run-deadline.ts"; type HandleRunFn = typeof _handleRunFn; async function handleRun(...args: Parameters): Promise>> { // LAZY: run.ts pulls in spawnBackgroundTeamRun + resolveCrewRuntime; also avoids jiti import race in child-process contexts. const mod = await import("./team-tool/run.ts"); return mod.handleRun(...args); } import { t } from "../i18n.ts"; import { waitForRun } from "../runtime/run-tracker.ts"; import { normalizeSkillOverride } from "../runtime/skill-instructions.ts"; import { formatActionSuggestion } from "./action-suggestions.ts"; import { type CacheControlDeps, invalidateSnapshot } from "./team-tool/cache-control.ts"; // API-5 facade dispatch: domain routers replace the former 54-case switch. import { domainForAction, handleAutomateDomain, handleControlDomain, handleManageDomain, handleRunDomain, handleStatusDomain, } from "./team-tool/dispatch/index.ts"; import { RUN_NOT_FOUND_HINT } from "./team-tool/run-not-found.ts"; export { handleApi } from "./team-tool/api.ts"; export { handleRetry } from "./team-tool/cancel.ts"; export type { TeamContext } from "./team-tool/context.ts"; export { handleDoctor } from "./team-tool/doctor.ts"; export { handleSchedule } from "./team-tool/handle-schedule.ts"; export { handleArtifacts, handleEvents, handleSummary, } from "./team-tool/inspect.ts"; export { handleCleanup, handleExport, handleForget, handleImport, handleImports, handlePrune, handleWorktrees, } from "./team-tool/lifecycle-actions.ts"; export { handleOrchestrate } from "./team-tool/orchestrate.ts"; export { handlePlan } from "./team-tool/plan.ts"; export { handleStatus } from "./team-tool/status.ts"; export type { TeamToolDetails } from "./team-tool-types.ts"; export { handleRun }; export function handleList(params: TeamToolParamsValue, ctx: TeamContext): PiTeamsToolResult { const resource = params.resource; const blocks: string[] = []; if (!resource || resource === "team") { const teams = allTeams(discoverTeams(ctx.cwd)); blocks.push( "Teams:", ...(teams.length ? teams.map((team) => formatScoped(team.name, team.source, team.description)) : ["- (none)"]), ); } if (!resource || resource === "workflow") { const workflows = allWorkflows(discoverWorkflows(ctx.cwd)); blocks.push( "", "Workflows:", ...(workflows.length ? workflows.map((workflow) => formatScoped(workflow.name, workflow.source, workflow.description)) : ["- (none)"]), ); } if (!resource || resource === "agent") { const agents = allAgents(discoverAgents(ctx.cwd)); blocks.push( "", "Agents:", ...(agents.length ? agents.map((agent) => formatScoped(agent.name, agent.source, agent.description)) : ["- (none)"]), ); } if (!resource) { // PERF (2026-08-24): listRecentRuns caps at source — collectRuns slices the // run-directory listing before reading manifests, instead of parsing every // manifest in scope and discarding all but 10. const runs = listRecentRuns(ctx.cwd, 10); blocks.push( "", "Recent runs:", ...(runs.length ? runs.map((run) => `- ${run.runId} [${run.status}] ${run.team}/${run.workflow ?? "none"}: ${run.goal}`) : ["- (none)"]), ); } return result(blocks.join("\n"), { action: "list", status: "ok" }); } export function handleGet(params: TeamToolParamsValue, ctx: TeamContext): PiTeamsToolResult { if (params.team) { const team = allTeams(discoverTeams(ctx.cwd)).find((item) => item.name === params.team); if (!team) return result(`Team '${params.team}' not found.`, { action: "get", status: "error" }, true); const lines = [ `Team: ${team.name} (${team.source})`, `Path: ${team.filePath}`, `Description: ${team.description}`, `Default workflow: ${team.defaultWorkflow ?? "(none)"}`, `Workspace mode: ${team.workspaceMode ?? "single"}`, "Roles:", ...(team.roles.length ? team.roles.map((role) => `- ${role.name} -> ${role.agent}${role.description ? `: ${role.description}` : ""}`) : ["- (none)"]), ]; return result(lines.join("\n"), { action: "get", status: "ok" }); } if (params.workflow) { const workflow = allWorkflows(discoverWorkflows(ctx.cwd)).find((item) => item.name === params.workflow); if (!workflow) return result(`Workflow '${params.workflow}' not found.`, { action: "get", status: "error" }, true); const lines = [ `Workflow: ${workflow.name} (${workflow.source})`, `Path: ${workflow.filePath}`, `Description: ${workflow.description}`, "Steps:", ...(workflow.steps.length ? workflow.steps.map((step) => `- ${step.id} [${step.role}] dependsOn=${step.dependsOn?.join(",") ?? "none"}`) : ["- (none)"]), ]; return result(lines.join("\n"), { action: "get", status: "ok" }); } if (params.agent) { const agent = allAgents(discoverAgents(ctx.cwd)).find((item) => item.name === params.agent); if (!agent) return result(`Agent '${params.agent}' not found.`, { action: "get", status: "error" }, true); const lines = [ `Agent: ${agent.name} (${agent.source})`, `Path: ${agent.filePath}`, `Description: ${agent.description}`, agent.model ? `Model: ${agent.model}` : undefined, agent.skills?.length ? `Skills: ${agent.skills.join(", ")}` : undefined, "", agent.systemPrompt || "(empty system prompt)", ].filter((line): line is string => line !== undefined); return result(lines.join("\n"), { action: "get", status: "ok" }); } return result("Specify team, workflow, or agent for get.", { action: "get", status: "error" }, true); } function artifactKey(artifact: ArtifactDescriptor): string { return `${artifact.kind}:${artifact.path}`; } /** Optional numeric params that accept the empty-string unset marker AND * stringified numbers in their schema (Union with Literal("") + a numeric-string * pattern branch). pi-ai's tool-argument coercion sometimes stringifies a * numeric value (e.g. interval:0 → "0") when the schema Union has a * string-literal branch, so the schema accepts the string form (Value.Check * passes) but handlers expect real numbers (TeamToolParamsValue types them as * number). This converts the string form back to a number — and treats "" as * unset (deleted) — before any domain router reads the param. */ const LOOSE_NUMERIC_PARAM_KEYS = ["interval", "budgetWarning", "budgetAbort", "tokenBudget", "replyDeadline", "budgetTotal"] as const; export function normalizeLooseNumericFields(params: TeamToolParamsValue): TeamToolParamsValue { let mutated = false; const out: Record = { ...params }; for (const key of LOOSE_NUMERIC_PARAM_KEYS) { const v = out[key]; if (v === undefined || v === null) continue; if (v === "") { delete out[key]; mutated = true; continue; } if (typeof v === "string") { const n = Number(v); if (Number.isFinite(n)) { out[key] = n; mutated = true; } else { delete out[key]; mutated = true; } } } // `once` is Union([Boolean, String, Number]); normalize its string form too. if (typeof out.once === "string") { if (out.once === "false") out.once = false; else if (out.once === "true") out.once = true; else { const n = Number(out.once); if (Number.isFinite(n)) out.once = n; } mutated = true; } return mutated ? (out as TeamToolParamsValue) : params; } async function recoverCheckpointedTasks( manifest: TeamRunManifest, tasks: TeamTaskState[], ): Promise<{ manifest: TeamRunManifest; tasks: TeamTaskState[]; recovered: string[] }> { const recovered: string[] = []; let nextManifest = manifest; const nextTasks = tasks.map((task) => { if (task.status !== "running" || !task.checkpoint) return task; if (task.checkpoint.phase === "artifact-written" && task.resultArtifact) { recovered.push(task.id); return { ...task, status: "completed" as const, finishedAt: task.finishedAt ?? task.checkpoint.updatedAt, error: undefined, claim: undefined, }; } if (task.checkpoint.phase === "child-stdout-final") { // transcripts are written with .attempt-${i}.jsonl suffix; find the most recent one const transcriptsDir = path.join(manifest.artifactsRoot, "transcripts"); let transcriptPath: string | undefined; if (fs.existsSync(transcriptsDir)) { const files = fs.readdirSync(transcriptsDir).filter((f) => f.startsWith(`${task.id}.attempt-`) && f.endsWith(".jsonl")); if (files.length > 0) { // Sort by attempt index descending to get the most recent files.sort((a, b) => { const idxA = parseInt(a.match(/\.attempt-(\d+)\./)?.[1] ?? "0", 10); const idxB = parseInt(b.match(/\.attempt-(\d+)\./)?.[1] ?? "0", 10); return idxB - idxA; }); transcriptPath = path.join(transcriptsDir, files[0]); } } if (!transcriptPath) return task; const transcript = fs.readFileSync(transcriptPath, "utf-8"); const parsed = parsePiJsonOutput(transcript); if (!parsed.finalText && !parsed.usage) return task; const resultArtifact = writeArtifact(manifest.artifactsRoot, { kind: "result", relativePath: `results/${task.id}.txt`, content: parsed.finalText ?? "(recovered from completed child transcript)", producer: task.id, }); const transcriptArtifact = writeArtifact(manifest.artifactsRoot, { kind: "log", relativePath: `transcripts/${task.id}.jsonl`, content: transcript, producer: task.id, }); recovered.push(task.id); return { ...task, status: "completed" as const, finishedAt: task.finishedAt ?? task.checkpoint.updatedAt, error: undefined, claim: undefined, resultArtifact, transcriptArtifact, usage: parsed.usage, jsonEvents: parsed.jsonEvents, }; } return task; }); if (recovered.length) { const artifacts = new Map(nextManifest.artifacts.map((artifact) => [artifactKey(artifact), artifact])); for (const task of nextTasks) { if (!recovered.includes(task.id)) continue; for (const artifact of [task.promptArtifact, task.resultArtifact, task.logArtifact, task.transcriptArtifact].filter( Boolean, ) as ArtifactDescriptor[]) artifacts.set(artifactKey(artifact), artifact); } nextManifest = { ...nextManifest, artifacts: [...artifacts.values()], updatedAt: new Date().toISOString(), }; await saveRunManifestAsync(nextManifest); saveRunTasks(nextManifest, nextTasks); } return { manifest: nextManifest, tasks: nextTasks, recovered }; } /** PID liveness probe — same semantics as filterAliveEntries: only ESRCH/ENOENT * mean "process does not exist"; EPERM means alive in another security context. */ function isPidAliveForResume(pid: number): boolean { try { process.kill(pid, 0); return true; } catch (error) { const code = (error as NodeJS.ErrnoException).code; return code !== "ESRCH" && code !== "ENOENT"; } } /** G12 (SDD-3 W-C WI-1): liveness evidence for a run being resumed. Returns a * human-readable refusal message when the run is LIVE, undefined when it is * safe to resume. force:true bypasses OWNERSHIP, never LIVENESS — callers must * refuse regardless of force when this returns a message. * * Three independent signals (any one ⇒ live): * 1. active-run-registry entry (activeRunEntries already applies the full * filter: terminal status, dead async PID, >30-min staleness); * 2. a running task holding a PRESENT, UNEXPIRED worker claim (task-claims * lease). Review fix 2026-10-01 (MAJOR): the claim must EXIST — coalesced * dispatch (run-coalesced-task-group.ts) and the task-graph scheduler mark * tasks running WITHOUT claims, so counting claim-less running tasks as * live permanently refused resume of crashed coalesced runs (with a * misleading "until undefined" message). Runs dispatched by real engines * carry the registry (signal 1) / async-PID (signal 3) liveness instead; * 3. the manifest's detached async PID is alive and the manifest is fresh * (covers goal-loop/DWF background runs whose registry entry was pruned). * Terminal runs (completed/failed/cancelled) never count as live via signal 3 * — a lingering post-finalize PID must not block legitimate re-runs. */ function resumeLivenessRefusal(manifest: TeamRunManifest, tasks: TeamTaskState[]): string | undefined { const runId = manifest.runId; const registryEntry = activeRunEntries().find((entry) => entry.runId === runId); if (registryEntry) { return [ `Run ${runId} is still live (active-run registry, last heartbeat ${registryEntry.updatedAt}) — resume refused to prevent double-dispatch.`, `Current status: ${manifest.status}. Wait for it to finish, or cancel it first (team action=cancel runId=${runId}).`, "force:true bypasses OWNERSHIP checks only — it can never bypass this liveness check.", ].join("\n"); } const claimed = tasks.find( // Review fix 2026-10-01 (MAJOR): require the claim to EXIST — see the // signal-2 note above (claim-less running tasks come from coalesced // dispatch / scheduler markings and must not read as permanently live). (task) => task.status === "running" && task.claim !== undefined && !isTaskClaimExpired(task.claim), ); if (claimed) { return [ `Run ${runId} is still live (task '${claimed.id}' holds an unexpired worker claim until ${claimed.claim?.leasedUntil}) — resume refused to prevent double-dispatch.`, `Current status: ${manifest.status}. Wait for the worker to finish, or cancel the run first (team action=cancel runId=${runId}).`, "force:true bypasses OWNERSHIP checks only — it can never bypass this liveness check.", ].join("\n"); } const pid = manifest.async?.pid; if ( typeof pid === "number" && Number.isInteger(pid) && pid > 0 && manifest.status !== "completed" && manifest.status !== "failed" && manifest.status !== "cancelled" ) { const updatedAt = Date.parse(manifest.updatedAt); // NIT 1 fix: shared registry constant — the resume horizon can never drift // from filterAliveEntries' staleness horizon. const fresh = Number.isFinite(updatedAt) && Date.now() - updatedAt <= ACTIVE_RUN_STALE_MS; if (fresh && isPidAliveForResume(pid)) { return [ `Run ${runId} is still live (background process PID ${pid} is alive) — resume refused to prevent double-dispatch.`, `Current status: ${manifest.status}. Wait for the background run to finish, or cancel it first (team action=cancel runId=${runId}).`, "force:true bypasses OWNERSHIP checks only — it can never bypass this liveness check.", ].join("\n"); } } return undefined; } /** G11 (SDD-3 W-C WI-2): adopt ownership + refresh heartbeat for a special-kind * (goal-loop / dynamic-workflow) resume, under the same per-run lock semantics * as the static path. Mirrors the static path's B1 adoption (battery 2026-08-18 * case b): a force-resumed run must not keep its dead original ownerSessionId, * or a third session's orphan-scan cancels the live resumed run. The adopted * manifest PRESERVES the original runKind — the whole point of the branch. */ async function adoptSpecialKindRunForResume( manifest: TeamRunManifest, runKind: "goal-loop" | "dynamic-workflow", ctx: TeamContext, ): Promise { return await withRunLock(manifest, async () => { const fresh = loadRunManifestById(manifest.cwd, manifest.runId); const base = fresh?.manifest ?? manifest; // W-C2 (SDD-4 WI-6): drop the dispatch-time detached-runner pointer. Resume // arms re-execute IN THIS session's process (runGoalLoop / // runDynamicWorkflow) — no runner is re-spawned — yet goal-wrapped runs carry // `async: { pid }` from their original background spawn. A dead pid on an // adopted "running" manifest is exactly the state that // transitionStaleAsyncUnderLock (status.ts) flips to failed on the next // status poll ("Async process stale: process does not exist"). The drop is // audited on the run.resume_requested event (clearedAsyncPid). const adoptedBase: TeamRunManifest = { ...base }; const staleAsyncPid = adoptedBase.async?.pid; delete adoptedBase.async; const adopted: TeamRunManifest = { ...adoptedBase, // WI-2 (G11): keep the ORIGINAL runKind — an accidental fallback to // "team-run" would route the NEXT resume through the static path. runKind, // Mark running at adoption: registerActiveRun refuses terminal // entries, and the run IS executing from this dispatch until terminal. // allowTerminalExit mirrors the static resume path (finding 8 write // guard — resume is the one legitimate terminal-exit flow). status: "running", summary: `Resuming ${runKind} run.`, updatedAt: new Date().toISOString(), ...(ctx.sessionId ? { ownerSessionId: ctx.sessionId } : {}), }; await saveRunManifestAsync(adopted, { allowTerminalExit: true }); await appendEventAsync(adopted.eventsPath, { type: "run.resume_requested", runId: adopted.runId, data: { runKind, action: "resume", ...(staleAsyncPid !== undefined ? { clearedAsyncPid: staleAsyncPid } : {}) }, }); return adopted; }); } /** G11 (SDD-3 W-C WI-2): resume arm for runKind="dynamic-workflow" — mirrors * run.ts's DWF dispatch (:390-435): re-synthesize the dynamic team, adopt the * run, then hand off to runDynamicWorkflow (which hydrates the dwf-checkpoint * state, dwf-runner.ts round-18 P2-3). The static team lookup CANNOT resolve * the synthetic `dwf-*` team, which is why this branch exists. */ async function resumeDynamicWorkflowRun( params: TeamToolParamsValue, ctx: TeamContext, manifest: TeamRunManifest, ): Promise { const dwfWorkflow = allWorkflows(discoverWorkflows(ctx.cwd)).find((candidate) => candidate.name === manifest.workflow); if (dwfWorkflow?.runtime !== "dynamic" || !dwfWorkflow?.dynamicScript) { return result( `Workflow '${manifest.workflow ?? ""}' is not a dynamic workflow (runKind=dynamic-workflow); cannot resume run ${manifest.runId}. Fix or restore the workflow file, or re-dispatch with action=run.`, { action: "resume", status: "error", runId: manifest.runId }, true, ); } // Re-synthesize the dynamic team (§0c C9 mirror of run.ts) — role resolution only. const dwfTeam: TeamConfig = { name: manifest.team, description: `Dynamic workflow run for ${dwfWorkflow.name}`, source: "dynamic", filePath: "", roles: [{ name: "worker", agent: params.agent ?? "executor" }], workspaceMode: "single", }; const adopted = await adoptSpecialKindRunForResume(manifest, "dynamic-workflow", ctx); registerActiveRun(adopted); // CORE-8 mirror: unified deadline (params.timeoutMs > config.maxRunMinutes > 1h default). const dwfDeadline = resolveRunDeadline(ctx, params); try { // LAZY: defer dynamic import of the DWF runner to its call site (mirrors run.ts:399). const { runDynamicWorkflow } = await import("../runtime/goal-workflow/dynamic-workflow-runner.ts"); const dwfResult = await runDynamicWorkflow({ manifest: adopted, workflow: dwfWorkflow as import("../workflows/workflow-config.ts").DynamicWorkflowConfig, team: dwfTeam, signal: dwfDeadline.signal, modelOverride: params.model, tokenBudget: params.tokenBudget ?? (dwfWorkflow as import("../workflows/workflow-config.ts").DynamicWorkflowConfig).maxTokenBudget, }); await saveRunManifestAsync(dwfResult.manifest); return result( [ `Resumed dynamic-workflow run ${dwfResult.manifest.runId}.`, `Status: ${dwfResult.manifest.status}`, dwfResult.manifest.summary ? `Result: ${dwfResult.manifest.summary}` : undefined, ] .filter((line): line is string => line !== undefined) .join("\n"), { action: "resume", status: dwfResult.manifest.status === "failed" ? "error" : "ok", runId: dwfResult.manifest.runId, artifactsRoot: dwfResult.manifest.artifactsRoot, }, dwfResult.manifest.status === "failed", ); } catch (runnerError) { // Round-11 runtime-fix mirror (run.ts): persist the failure instead of // leaving the manifest at its pre-resume status forever. const failureReason = runnerError instanceof Error ? runnerError.message : String(runnerError); const failedManifest = { ...adopted, status: "failed" as const, summary: `Dynamic workflow '${dwfWorkflow.name}' resume failed: ${failureReason}`.slice(0, 2000), updatedAt: new Date().toISOString(), }; await saveRunManifestAsync(failedManifest); return result( `Dynamic workflow '${dwfWorkflow.name}' resume failed: ${failureReason}`, { action: "resume", status: "error", runId: failedManifest.runId, artifactsRoot: failedManifest.artifactsRoot }, true, ); } finally { unregisterActiveRun(adopted.runId); clearTimeout(dwfDeadline.timer); // RC-02 } } /** G11 (SDD-3 W-C WI-2): resume arm for runKind="goal-loop" — mirrors the * background-runner.ts goal-loop short-circuit (:492): hydrate GoalLoopState, * continue the loop via runGoalLoop, map the goal outcome to a run status. */ async function resumeGoalLoopRun(params: TeamToolParamsValue, ctx: TeamContext, manifest: TeamRunManifest): Promise { // Check the goal state BEFORE adoption: a missing GoalLoopState is not // resumable at all, so the run must be left untouched (no ownership flip, // no status change) for a clearer retry story. // LAZY: defer heavy goal-loop imports to the call site (mirrors background-runner.ts:498-506). const { GoalStore } = await import("../runtime/goal-workflow/goal-state-store.ts"); const store = new GoalStore(manifest.cwd); const goalState = store.load(manifest.runId); if (!goalState) { return result( `runKind="goal-loop" but GoalLoopState '${manifest.runId}' not found (cwd=${manifest.cwd}); cannot resume. The goal state file may have been pruned — start a new goal run instead.`, { action: "resume", status: "error", runId: manifest.runId }, true, ); } const adopted = await adoptSpecialKindRunForResume(manifest, "goal-loop", ctx); // LAZY: the runner import stays after the state check for fast failure. const { runGoalLoop } = await import("../runtime/goal-workflow/goal-loop-runner.ts"); registerActiveRun(adopted); // CORE-8 mirror (review MINOR 4): honor params.timeoutMs like the DWF arm — // previously this arm resolved the deadline with {} and silently ignored the // caller's timeout override. const goalDeadline = resolveRunDeadline(ctx, params); try { const goalResult = await runGoalLoop({ goalState, manifest: adopted, signal: goalDeadline.signal, deps: { discoverAgents: (cwd: string) => allAgents(discoverAgents(cwd)), }, }); // Fix P1-1 + round-6 #5 mirror (background-runner.ts:525-536): persist the // terminal status reflecting the goal's actual outcome, not a blanket // 'completed'. const goalStatusToRunStatus: Record = { achieved: "completed", max_turns: "completed", budget_exceeded: "completed", blocked: "blocked", cancelled: "cancelled", paused: "blocked", running: "running", }; const runStatus = goalStatusToRunStatus[goalResult.goalState.state] ?? "completed"; const finalManifest: TeamRunManifest = { ...goalResult.manifest, status: runStatus, updatedAt: new Date().toISOString(), }; await saveRunManifestAsync(finalManifest); return result( [ `Resumed goal-loop run ${finalManifest.runId}.`, `Status: ${finalManifest.status} (goal state: ${goalResult.goalState.state})`, ].join("\n"), { action: "resume", status: runStatus === "failed" ? "error" : "ok", runId: finalManifest.runId, artifactsRoot: finalManifest.artifactsRoot, }, runStatus === "failed", ); } catch (runnerError) { // Round-11 runtime-fix mirror (DWF arm + review MINOR 1): persist the // failure instead of leaving the adopted manifest stuck at // status:"running" with a live registry entry until the stale horizon. const failureReason = runnerError instanceof Error ? runnerError.message : String(runnerError); const failedManifest = { ...adopted, status: "failed" as const, summary: `Goal-loop resume failed: ${failureReason}`.slice(0, 2000), updatedAt: new Date().toISOString(), }; await saveRunManifestAsync(failedManifest); return result( `Goal-loop run ${failedManifest.runId} resume failed: ${failureReason}`, { action: "resume", status: "error", runId: failedManifest.runId, artifactsRoot: failedManifest.artifactsRoot }, true, ); } finally { unregisterActiveRun(adopted.runId); clearTimeout(goalDeadline.timer); // RC-02 } } export async function handleResume(params: TeamToolParamsValue, ctx: TeamContext): Promise { if (!params.runId) return result("Resume requires runId.", { action: "resume", status: "error" }, true); const runCwd = locateRunCwd(params.runId, ctx.cwd); if (!runCwd) return result(`Run '${params.runId}' not found.${RUN_NOT_FOUND_HINT}`, { action: "resume", 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: "resume", status: "error" }, true); // G12 (SDD-3 W-C WI-1): liveness-first gate — BEFORE the ownership check and // BEFORE any reset/dispatch. executeTeamRun runs OUTSIDE the resume lock // (LOCK-2 below), so a resume accepted against a live run double-dispatches // workers (duplicate tokens + duplicate side effects). force:true bypasses // OWNERSHIP only — never LIVENESS — so the refusal applies even when forced. const liveness = resumeLivenessRefusal(loaded.manifest, loaded.tasks); if (liveness) return result(liveness, { action: "resume", status: "error", runId: loaded.manifest.runId }, true); // R1: foreign-ownership check — mirrors handleRetry/handleCancel. Without it, // another session can resume (and re-execute) a run it doesn't own, racing // the owning session. const foreignRun = typeof loaded.manifest.ownerSessionId === "string" && loaded.manifest.ownerSessionId !== ctx.sessionId; if (foreignRun && !params.force) { return result( `Run ${loaded.manifest.runId} belongs to another session. Use force: true to override.`, { action: "resume", status: "error", runId: loaded.manifest.runId }, true, ); } // G4 (SDD-3 W-C WI-1): force-on-foreign is an authorization override — liveness // was already REFUSED above, so force only bypassed OWNERSHIP here. Record it // as a security event (registered in TEAM_EVENT_TYPES) so cross-session forced // resumes are auditable; never silent. if (foreignRun && params.force) { await appendEventAsync(loaded.manifest.eventsPath, { type: "run.resume_forced_foreign", runId: loaded.manifest.runId, message: `Foreign run ${loaded.manifest.runId} force-resumed by session ${ctx.sessionId} (owner ${loaded.manifest.ownerSessionId}).`, data: { action: "resume", forced: true, ownerSessionId: loaded.manifest.ownerSessionId, forcingSessionId: ctx.sessionId, runStatus: loaded.manifest.status, }, }); } if (!loaded.manifest.workflow) return result(`Run '${params.runId}' has no workflow to resume.`, { action: "resume", status: "error" }, true); // G11 (SDD-3 W-C WI-2): goal-loop and dynamic-workflow runs use SYNTHETIC // team names (goal-*/dwf-*) that never appear in discoverTeams — the static // team lookup below fails with "Team not found" and the run is unresumable // even though both engines have resume support (DWF checkpoint hydration, // goal-loop state). Branch on the ORIGINAL runKind (mirrors run.ts:390-435 // dispatch and background-runner.ts:492 short-circuit); the resumed run // PRESERVES its runKind (see adoptSpecialKindRunForResume). const resumeRunKind = loaded.manifest.runKind ?? "team-run"; if (resumeRunKind === "dynamic-workflow") { return await resumeDynamicWorkflowRun(params, ctx, loaded.manifest); } if (resumeRunKind === "goal-loop") { return await resumeGoalLoopRun(params, ctx, loaded.manifest); } const agents = allAgents(discoverAgents(ctx.cwd)); const direct = directTeamAndWorkflowFromRun(loaded.manifest, loaded.tasks, agents); const team = direct?.team ?? allTeams(discoverTeams(ctx.cwd)).find((candidate) => candidate.name === loaded.manifest.team); if (!team) return result(`Team '${loaded.manifest.team}' not found.`, { action: "resume", status: "error" }, true); const workflow = direct?.workflow ?? allWorkflows(discoverWorkflows(ctx.cwd)).find((candidate) => candidate.name === loaded.manifest.workflow); if (!workflow) return result(`Workflow '${loaded.manifest.workflow}' not found.`, { action: "resume", status: "error" }, true); // LOCK-2 (Round 2): Lock held only for recovery + reset (the read-modify-write // critical section). executeTeamRun runs OUTSIDE the lock — it uses // team-runner's own per-operation withRunLock calls (mergeUnitResult, // handleFailedTask merge, finalizeRun), mirroring handleRun (which never wraps // executeTeamRun in an outer lock). Holding the resume lock across the // minutes-long executeTeamRun would let any process steal the lock after // DEFAULT_LOCKS.staleMs (30s) — see defaults.ts. const decision = await withRunLock(loaded.manifest, async () => { // R2: re-read inside the lock so recovery + resetTasks reflect committed // state, not the pre-lock snapshot. Between the pre-lock load and lock // acquisition a task may have been cancelled (stale-reconciler) or // completed (background runner); using the stale snapshot would reset // running→queued and re-execute it (double execution: duplicate tokens + // duplicate side effects). Sibling handlers (cancelOrphanedRuns, // reconcileAllStaleRuns) re-read inside the lock for the same reason. const fresh = loadRunManifestById(runCwd, loaded.manifest.runId); const lockedManifest = fresh?.manifest ?? loaded.manifest; const lockedTasks = fresh?.tasks ?? loaded.tasks; // F1 (2026-10-01 security review): the PRE-lock liveness gate raced // concurrent dispatch — between the pre-lock load and this lock acquisition // another session may have adopted and REGISTERED the run. Re-check // liveness against the LOCKED state before any recovery/reset mutation; // force:true still never bypasses liveness (G12 principle). const lockedLiveness = resumeLivenessRefusal(lockedManifest, lockedTasks); if (lockedLiveness) { return { kind: "blocked" as const, payload: result(lockedLiveness, { action: "resume", status: "error", runId: lockedManifest.runId }, true), }; } const loadedConfig = loadConfig(ctx.cwd); const recovered = await recoverCheckpointedTasks(lockedManifest, lockedTasks); const resumeManifest = recovered.manifest; // W-C2 (SDD-4 WI-6): drop the stale detached-runner pointer — resume never // re-spawns a background runner (executeTeamRun below runs in THIS // session's process), yet the adopted manifest kept the dispatch-time // `async: { pid }` block pointing at the (now dead) original runner. During // the re-execution window (status=running) any status poll read that dead // pid and transitionStaleAsyncUnderLock flipped the LIVE resumed run to // failed ("Async process stale: process does not exist", cancelling running // tasks) — real-test battery 2026-10-02 finding 1: resume of a completed // async run died 14s in while the sync contrast passed. Invariant restored: // manifest.async is present ONLY while a detached runner owns execution; a // resumed run's liveness is carried by registerActiveRun + ownerSessionId. // The drop is audited on the run.resume_requested event (clearedAsyncPid). const resumeBase: TeamRunManifest = { ...resumeManifest }; const staleAsyncPid = resumeBase.async?.pid; delete resumeBase.async; const executedConfig = { ...effectiveRunConfig(loadedConfig.config, params.config), }; // Preserve original manifest scaffold mode when resume has no explicit mode override // AND workers are not explicitly disabled. If workers are disabled, let // resolveCrewRuntime detect it and return blocked safety. if (!executedConfig.runtime?.mode && resumeManifest.runtimeResolution?.safety === "explicit_dry_run") { const workersDisabled = executedConfig.executeWorkers === false || getCrewEnv("PI_CREW_EXECUTE_WORKERS") === "0" || getCrewEnv("PI_TEAMS_EXECUTE_WORKERS") === "0"; if (!workersDisabled) executedConfig.runtime = { ...executedConfig.runtime, mode: "scaffold", }; } const runtime = await resolveCrewRuntime(executedConfig); const runtimeResolution = runtimeResolutionState(runtime); const runtimeManifest = { ...resumeBase, runtimeResolution, updatedAt: new Date().toISOString(), // B1 battery 2026-08-18 (case b root-cause fix): resume ADOPTS the run — // ownerSessionId moves to the resuming session. Pre-fix, a force-resumed // run kept its DEAD original owner id, so any THIRD session's startup // orphan-scan saw "owner no longer exists" and cancelled the live, // freshly-resumed run (observed live: orphan_cancelled 5m into a parked // ask that the resume itself had just dispatched). Same-session resume // is a no-op (id equality). ...(ctx.sessionId ? { ownerSessionId: ctx.sessionId } : {}), }; await saveRunManifestAsync(runtimeManifest); await appendEventAsync(runtimeManifest.eventsPath, { type: "runtime.resolved", runId: runtimeManifest.runId, message: `Runtime resolved for resume: ${runtime.kind} safety=${runtime.safety}`, data: { runtimeResolution, action: "resume" }, }); if (runtime.safety === "blocked") { const runningManifest = updateRunStatus(runtimeManifest, "running", "Checking worker runtime availability before resume.", { // Resume is the ONE legitimate terminal-exit flow (finding 8 write guard). allowTerminalExit: true, }); const blocked = updateRunStatus( runningManifest, "blocked", runtime.reason ?? "Child worker execution is disabled; refusing to resume with no-op scaffold subagents.", ); await appendEventAsync(blocked.eventsPath, { type: "run.blocked", runId: blocked.runId, message: blocked.summary, data: { runtime, action: "resume" }, }); return { kind: "blocked" as const, payload: result( [ `Blocked resume for pi-crew run ${blocked.runId}: real subagent workers are disabled.`, `Runtime: ${runtime.kind} (requested ${runtime.requestedMode})`, runtime.reason ?? "Child worker execution is disabled.", "", "To resume effective subagents, remove executeWorkers=false / PI_CREW_EXECUTE_WORKERS=0 / PI_TEAMS_EXECUTE_WORKERS=0 or set runtime.mode=child-process.", "Use runtime.mode=scaffold only for explicit dry-run prompt/artifact generation.", ].join("\n"), { action: "resume", status: "error", runId: blocked.runId, artifactsRoot: blocked.artifactsRoot, }, true, ), }; } const resetTasks = recovered.tasks.map((task) => task.status === "failed" || task.status === "cancelled" || task.status === "skipped" || task.status === "running" ? { ...task, status: "queued" as const, error: undefined, startedAt: undefined, finishedAt: undefined, claim: undefined, } : task, ); saveRunTasks(runtimeManifest, resetTasks); const replay = replayPendingMailboxMessages(runtimeManifest); await appendEventAsync(runtimeManifest.eventsPath, { type: "run.resume_requested", runId: runtimeManifest.runId, data: { replayedMailboxMessages: replay.messages.length, recoveredCheckpointTasks: recovered.recovered, // W-C2 (SDD-4 WI-6): audit trail for the stale detached-runner pointer // dropped at adoption (see resumeBase above). ...(staleAsyncPid !== undefined ? { clearedAsyncPid: staleAsyncPid } : {}), }, }); if (recovered.recovered.length) await appendEventAsync(runtimeManifest.eventsPath, { type: "task.checkpoint_recovered", runId: runtimeManifest.runId, message: `Recovered ${recovered.recovered.length} task(s) from artifact-written checkpoints.`, data: { taskIds: recovered.recovered }, }); if (replay.messages.length) await appendEventAsync(runtimeManifest.eventsPath, { type: "mailbox.replayed", runId: runtimeManifest.runId, message: `Replayed ${replay.messages.length} pending inbox message(s).`, data: { messageIds: replay.messages.map((message) => message.id), taskIds: replay.messages.map((message) => message.taskId).filter(Boolean), }, }); const executeWorkers = runtime.kind !== "scaffold"; const resumeSkillOverride = normalizeSkillOverride(params.skill) ?? runtimeManifest.skillOverride; // F1 (2026-10-01 security review): adopt running + REGISTER in the // cross-process active-run registry INSIDE the lock, mirroring run.ts // dispatch (:388) and adoptSpecialKindRunForResume. The static path // previously never registered, so during foreground executeTeamRun (which // runs OUTSIDE the lock) a concurrent resume saw NO cross-process liveness // signal until the first task claim appeared — the double-dispatch window // (even with force:true). Terminal finalize unregisters (updateRunStatus → // unregisterActiveRun); a resume that crashes mid-execution ages out via // the registry's staleness horizon. registerActiveRun throwing here (a // raced external cancel flipped the disk manifest terminal) is intentional // fail-safe: the other writer won the race, so we must not dispatch. const executingManifest = updateRunStatus( runtimeManifest, "running", `Resuming run (adopted by session ${ctx.sessionId ?? "unknown"}).`, // Finding 8 write guard: resume is the legitimate terminal-exit flow. { allowTerminalExit: true }, ); registerActiveRun(executingManifest); return { kind: "execute" as const, runtimeManifest: executingManifest, resetTasks, executeWorkers, resumeSkillOverride, runtime, executedConfig, }; }); // Lock is now RELEASED. executeTeamRun runs without holding the resume // lock — it uses team-runner's own per-operation withRunLock calls // (mergeUnitResult, handleFailedTask merge, finalizeRun), mirroring // handleRun (which never wraps executeTeamRun in an outer lock). if (decision.kind === "blocked") return decision.payload; const executed = await executeTeamRun({ manifest: decision.runtimeManifest, tasks: decision.resetTasks, team, workflow, agents, executeWorkers: decision.executeWorkers, limits: decision.executedConfig.limits, runtime: decision.runtime, runtimeConfig: decision.executedConfig.runtime, parentContext: buildParentContext(ctx), parentModel: resolveParentModel(ctx.model), modelRegistry: ctx.modelRegistry, modelOverride: params.model, skillOverride: decision.resumeSkillOverride, signal: ctx.signal, reliability: decision.executedConfig.reliability, metricRegistry: ctx.metricRegistry, workspaceId: ctx.sessionId ?? ctx.cwd, // Finding 8: resume is the legitimate terminal-exit flow. isResume: true, }); return result( [ `Resumed run ${executed.manifest.runId}.`, `Status: ${executed.manifest.status}`, `Tasks: ${executed.tasks.length}`, `Artifacts: ${executed.manifest.artifactsRoot}`, ].join("\n"), { action: "resume", status: executed.manifest.status === "failed" ? "error" : "ok", runId: executed.manifest.runId, artifactsRoot: executed.manifest.artifactsRoot, }, executed.manifest.status === "failed", ); } export function handleSteer(params: TeamToolParamsValue, ctx: TeamContext): PiTeamsToolResult { const { runId, taskId, message } = params; if (!runId || !taskId || !message) { return result("steer requires runId, taskId, and message", { action: "steer", status: "error" }, true); } const runCwd = locateRunCwd(runId, ctx.cwd); if (!runCwd) return result(`Run '${runId}' not found`, { action: "steer", status: "error" }, true); const loaded = loadRunManifestById(runCwd, runId); // NOTE: no withRunLock - best-effort only; concurrent writes may cause inconsistency if (!loaded) return result(`Run '${runId}' not found`, { action: "steer", status: "error" }, true); const task = loaded.tasks.find((t) => t.id === taskId); if (!task) return result(`Task '${taskId}' not found`, { action: "steer", status: "error" }, true); if (!task.pendingSteers) task.pendingSteers = []; // T-S1: do not allow steering a task that has already reached a terminal status. if (TEAM_TERMINAL_TASK_STATUSES.has(task.status)) { return result(`Task '${taskId}' is ${task.status}; cannot steer.`, { action: "steer", status: "error" }, true); } // HIGH-04: Cap pendingSteers array to prevent unbounded memory growth const MAX_PENDING_STEERS = 100; if (task.pendingSteers.length >= MAX_PENDING_STEERS) { // Log warning before dropping the oldest message appendEventFireAndForget(loaded.manifest.eventsPath, { type: "task.steer_dropped", runId, taskId, data: { droppedMessage: task.pendingSteers[0], reason: "pendingSteers cap exceeded", queueDepth: task.pendingSteers.length, }, }); task.pendingSteers = task.pendingSteers.slice(-(MAX_PENDING_STEERS - 1)); } task.pendingSteers.push(message); saveRunTasks(loaded.manifest, loaded.tasks); // Real-time steer delivery: write to steering file so child can read immediately try { const steeringDir = `${loaded.manifest.artifactsRoot}/steering`; fs.mkdirSync(steeringDir, { recursive: true }); // AUDIT-08 defense-in-depth: validate the steering-file path is contained // within steeringDir. taskId is currently sanitized via createTaskId, but this // guards against future changes to task-id generation (e.g. if it ever // accepted user input). const safeSteeringPath = resolveRealContainedPath(steeringDir, `${taskId}.jsonl`); fs.appendFileSync(safeSteeringPath, JSON.stringify({ type: "steer", message, ts: new Date().toISOString() }) + "\n"); } catch { // Best-effort: file write failure doesn't block the steer from pending array } // H1 (2026-08-10): handleSteer is a SYNC function — fire-and-forget async. void appendEventAsync(loaded.manifest.eventsPath, { type: "task.steer_queued", runId, taskId, data: { message }, }).catch((error) => logInternalError("steer.event", error instanceof Error ? error : new Error(String(error)), `runId=${runId}`)); return result(`Steer queued for task '${taskId}'. It will be delivered when the task's session is ready.`, { action: "steer", status: "ok", }); } export function cacheControlDepsFromContext(ctx: TeamContext): CacheControlDeps | undefined { if (!ctx.getRunSnapshotCache) return undefined; return { getRunSnapshotCache: ctx.getRunSnapshotCache }; } export function handleInvalidate(params: TeamToolParamsValue, ctx: TeamContext): PiTeamsToolResult { const runId = params.runId; if (!runId) return result("Invalidate requires runId.", { action: "invalidate", status: "error" }, true); const runCwd = locateRunCwd(runId, ctx.cwd); if (!runCwd) return result(`Run '${runId}' not found.`, { action: "invalidate", status: "error" }, true); const deps = cacheControlDepsFromContext(ctx); if (!deps) return result("Cache invalidation not available (no snapshot cache).", { action: "invalidate", status: "error" }, true); invalidateSnapshot(runId, runCwd, deps); return result(`Cache invalidated for run ${runId}.`, { action: "invalidate", status: "ok", runId, }); } /** * Locate the CWD where a run's state is stored. * Tries ctx.cwd first, then scans immediate child directories for .crew/state/runs/. * * Defensive bounds (prevent hang on large dirs like /tmp in CI): * - Skips entries that are well-known system/ephemeral dirs (e.g. .npm, node_modules, .git) * - Caps the scan at MAX_SCAN_ENTRIES to avoid pathological scans * - Skips hidden entries (starting with `.`) unless they look like run directories * (e.g. .crew, .pi, .tmp-crew-runs) */ const MAX_SCAN_ENTRIES = 1000; const SKIP_SCAN_DIRS = new Set(["node_modules", ".git", ".npm", ".cache", ".local", "proc", "sys", "dev", "Library", "Applications"]); // PERF (2026-08-24): a stale/typo'd runId from a looping LLM caller paid the // full 1000-entry directory sweep on EVERY attempt. Resolution results (hits // AND misses) are cached briefly; TTL bounds staleness for runs created in a // sibling cwd after a cached miss. const runCwdCache = new Map(); const RUN_CWD_TTL_MS = 30_000; const RUN_CWD_CACHE_MAX = 128; export function locateRunCwd(runId: string, baseCwd: string): string | undefined { const key = `${baseCwd}\0${runId}`; const cached = runCwdCache.get(key); if (cached && cached.expiresAt > Date.now()) return cached.cwd; const cwd = locateRunCwdUncached(runId, baseCwd); // original body, renamed if (runCwdCache.size >= RUN_CWD_CACHE_MAX) { const oldest = runCwdCache.keys().next().value; if (oldest !== undefined) runCwdCache.delete(oldest); } runCwdCache.set(key, { cwd, expiresAt: Date.now() + RUN_CWD_TTL_MS }); return cwd; } /** * Locate the CWD where a run's state is stored. * Tries ctx.cwd first, then scans immediate child directories for .crew/state/runs/. * * Defensive bounds (prevent hang on large dirs like /tmp in CI): * - Skips entries that are well-known system/ephemeral dirs (e.g. .npm, node_modules, .git) * - Caps the scan at MAX_SCAN_ENTRIES to avoid pathological scans * - Skips hidden entries (starting with `.`) unless they look like run directories * (e.g. .crew, .pi, .tmp-crew-runs) */ /** * F-L1 companion (2026-09-23 correction): a hit counts when the resolved * stateRoot sits under EITHER root that `list` unions (run-index * `scopedRunRoots` = userCrewRoot + projectCrewRoot), because this function * answers "which cwd can read this run's state" — not "was it created here". * * History: the first F-L1 fix (5d3a3d0f) restricted hits to the candidate * cwd's PRIMARY root. That kept `locateRunCwd` from claiming a run created in * a sibling, but it also made every by-ID handler (status/events/summary/ * artifacts/worktrees/api/plans/respond/cancel) reject runs that `list` shows * from the user root: `list` unions both roots, `loadRunManifestById` resolves * both roots, yet `locateRunCwd` returned undefined → "Run not found" for * runs the dashboard had just listed. Measured live: 10/10 user-root runs * listed by `action='list'` failed `action='status'`. * * Sibling isolation is preserved WITHOUT the primary-only restriction: a run * living under a sibling's PROJECT root (…/sibling/.crew/state/runs/) is * not under the candidate cwd's project root nor the user root, so it still * does not match (pinned by locate-run-cwd.test.ts "sibling directory"). */ function runUnderListedRoot(cwd: string, runId: string): boolean { const loaded = loadRunManifestById(cwd, runId); if (!loaded) return false; const roots = findRepoRoot(cwd) ? [projectCrewRoot(cwd), userCrewRoot()] : [userCrewRoot()]; return roots.some((root) => loaded.manifest.stateRoot === path.join(root, DEFAULT_PATHS.state.runsSubdir, runId)); } /** * Scan-path variant: a CHILD directory may claim the run only when the run's * state actually lives under that child's OWN crew root (the nested-child * case this scan exists for). Without this, a shared user-root run would be * "claimed" by whichever child happened to be scanned first. */ function runUnderOwnRoot(cwd: string, runId: string): boolean { const loaded = loadRunManifestById(cwd, runId); if (!loaded) return false; return loaded.manifest.stateRoot === path.join(projectCrewRoot(cwd), DEFAULT_PATHS.state.runsSubdir, runId); } export function locateRunCwdUncached(runId: string, baseCwd: string): string | undefined { // Fast path: run resolves under one of the roots `list` shows (see helper note). if (runUnderListedRoot(baseCwd, runId)) { return baseCwd; } // Scan immediate child directories, but with defensive bounds. try { const entries = fs.readdirSync(baseCwd, { withFileTypes: true }); const boundedEntries = entries.length > MAX_SCAN_ENTRIES ? entries.slice(0, MAX_SCAN_ENTRIES) : entries; for (const entry of boundedEntries) { if (!entry.isDirectory()) continue; if (SKIP_SCAN_DIRS.has(entry.name)) continue; // Skip hidden entries except well-known run-storage prefixes if (entry.name.startsWith(".")) { if (!entry.name.startsWith(".crew") && !entry.name.startsWith(".pi") && !entry.name.startsWith(".tmp-crew")) continue; } const candidate = path.join(baseCwd, entry.name); if (runUnderOwnRoot(candidate, runId)) { return candidate; } } } catch { /* ignore unreadable dirs */ } return undefined; } export async function handleWait(params: TeamToolParamsValue, ctx: TeamContext): Promise { const { runId } = params; if (!runId) return result("wait requires runId.", { action: "wait", status: "error" }, true); const timeoutMs = Math.min( Math.max( typeof params.config?.timeoutMs === "number" && Number.isFinite(params.config.timeoutMs) ? params.config.timeoutMs : 300_000, 1_000, // minimum 1 s ), 3_600_000, // maximum 1 h ); const pollIntervalMs = Math.max( Math.min( typeof params.config?.pollIntervalMs === "number" && Number.isFinite(params.config.pollIntervalMs) ? params.config.pollIntervalMs : 2000, 60_000, // maximum 60 s ), 500, // minimum 500 ms ); // Resolve the run's CWD: try ctx.cwd first, then scan child dirs with .crew/ const runCwd = locateRunCwd(runId, ctx.cwd); if (!runCwd) { return result(`Run '${runId}' not found in '${ctx.cwd}' or its subdirectories.`, { action: "wait", status: "error", runId }, true); } try { const { manifest, tasks } = await waitForRun(runId, runCwd, { timeoutMs, pollIntervalMs, }); const taskSummary = tasks.map((t) => ` ${t.id}: ${t.status}`).join("\n"); return result( [`Run ${runId} finished: ${manifest.status}`, `Summary: ${manifest.summary ?? "(none)"}`, `Tasks:`, taskSummary].join("\n"), { action: "wait", status: manifest.status === "failed" ? "error" : "ok", runId: manifest.runId, }, manifest.status === "failed", ); } catch (err) { const msg = err instanceof Error ? err.message : String(err); return result(`wait failed: ${msg}`, { action: "wait", status: "error", runId }, true); } } export async function handleTeamTool(params: TeamToolParamsValue, ctx: TeamContext): Promise { // API-5 fix: normalize action into params so the domain routers (which read // params.action) see the resolved default, not the raw undefined. Without this, // a missing action defaulted to "list" at the facade but the router read // params.action=undefined → "Unhandled status-domain action: undefined". params = { ...params, action: params.action ?? "list" }; // Coerce stringified-numeric params (pi-ai coercion artifact) back to real // numbers before domain routers read them. See normalizeLooseNumericFields. params = normalizeLooseNumericFields(params); const action = params.action ?? "list"; const domain = domainForAction(action); switch (domain) { case "run": return handleRunDomain(params, ctx); case "status": return handleStatusDomain(params, ctx); case "control": return handleControlDomain(params, ctx); case "manage": return handleManageDomain(params, ctx); case "automate": return handleAutomateDomain(params, ctx); default: return result( t("team.unknownAction", { action: String(action) }) + formatActionSuggestion(String(action)), { action: "unknown", status: "error" }, true, ); } } /** * Module-scoped RPC registry for access to pi-crew's team orchestrator. * * EXT-9: Previously this used a `globalThis[Symbol.for("pi-crew:registry")]` * singleton — fragile (cross-realm, no lifecycle, peer extensions could * read/overwrite it). It is now a module-level variable: exactly one instance * per extension load (one per pi session), invisible to peer extensions. * Cross-extension consumers should use the `pi.events` RPC channel * (`registerPiCrewRpc`) instead of poking globalThis. */ interface CrewRegistry { version: 2; getRecord: (runId: string) => TeamRunManifest | undefined; listRuns: () => Array<{ runId: string; status: string; goal: string }>; appendEvent: (runId: string, event: Record) => void; waitForAll: (runId: string) => Promise; hasRunning: (runId: string) => boolean; /** Register a dynamic agent at runtime. Invalidates the discovery cache. */ registerAgent: (config: AgentConfig) => void; /** Unregister a previously registered dynamic agent. Invalidates the discovery cache. */ unregisterAgent: (name: string) => void; /** List all currently registered dynamic agents. */ listDynamicAgents: () => AgentConfig[]; } // ─── Dynamic Agent Registry (Phase 3b) ─────────────────────────────────── // The dynamic agent store lives in discover-agents.ts and is merged into // discovery results with highest priority. The CrewRegistry interface exposes // registerAgent/unregisterAgent/listDynamicAgents for cross-extension access. // Module-scoped singleton instance — one per extension load (EXT-9). let crewRegistryInstance: CrewRegistry | undefined; export function registerCrewGlobalRegistry(registry: CrewRegistry): void { crewRegistryInstance = registry; } /** @internal — exported for lifecycle tests. */ export function getCrewGlobalRegistry(): CrewRegistry | undefined { return crewRegistryInstance; } /** Manifest cache shape needed to construct the global registry's read-side. */ interface ManifestCacheForRegistry { get: (runId: string) => TeamRunManifest | undefined; list: (limit: number) => TeamRunManifest[]; } /** * Build and install the global CrewRegistry singleton in a single atomic step. * * EXT-7 (Round 3): The previous design called `installCrewGlobalRegistry()` to * install stubs, then patched the manifest-backed methods asynchronously inside * `register.ts`. That left a window where cross-extension consumers could observe * the stub object on `globalThis[Symbol.for("pi-crew:registry")]`. By taking the * real dependencies up-front, we install the registry once with no stub phase — * callers see either no registry (pre-init) or the fully-real registry. */ export function installCrewGlobalRegistry(deps?: { manifestCache: ManifestCacheForRegistry; cwdProvider: () => string }): void { const manifestCache = deps?.manifestCache; const cwdProvider = deps?.cwdProvider ?? ((): string => process.cwd()); const registry: CrewRegistry = { version: 2, getRecord: (runId: string) => manifestCache?.get(runId), listRuns: () => manifestCache ? manifestCache.list(100).map((m) => ({ runId: m.runId, status: m.status, goal: m.goal })) : ([] as Array<{ runId: string; status: string; goal: string }>), appendEvent: (runId: string, event: Record) => { if (!manifestCache) return; const manifest = manifestCache.get(runId); if (manifest) { // LAZY: event-log is already loaded at module top, so use the // pre-resolved appendEventFireAndForget instead of re-importing. appendEventFireAndForget(manifest.eventsPath, event as Parameters[1]); } }, waitForAll: async (runId: string) => { if (!manifestCache) return; // LAZY: state-store is already loaded at module top; use the pre-resolved loadRunManifestById. const check = (): boolean => { const loaded = loadRunManifestById(cwdProvider(), runId); if (!loaded) return true; return !loaded.tasks.some((t: { status: string }) => t.status === "running" || t.status === "queued"); }; while (!check()) await new Promise((resolve) => setTimeout(resolve, 500)); }, hasRunning: (runId: string) => { if (!manifestCache) return false; const manifest = manifestCache.get(runId); if (!manifest) return false; // LAZY: state-store is already loaded at module top; use the pre-resolved loadRunManifestById. const loaded = loadRunManifestById(cwdProvider(), runId); if (!loaded) return false; return loaded.tasks.some((t: { status: string }) => t.status === "running" || t.status === "queued"); }, registerAgent: registerDynamicAgent, unregisterAgent: unregisterDynamicAgent, listDynamicAgents, }; registerCrewGlobalRegistry(registry); } /** Remove the CrewRegistry singleton. Call during session cleanup. */ export function uninstallCrewGlobalRegistry(): void { crewRegistryInstance = undefined; }