// agents-lib — pure orchestration logic for the spawn_agents tool
// (headless pi child runs). Pure/injectable so node:test covers it
// without ever spawning a real agent.
//
// Children are headless runs of THIS entry (process.argv[1], i.e. the
// running pi entry — same runtime, same global extension/config dir)
// in the parent's workspace, with `--mode json -p` so their output is
// a parseable event stream. Each child writes its own session JSONL
// under <parent's session dir>/subagents/ (subagentSessionDir: kept out
// of the folder /resume and --continue read, so a user only ever sees
// their own conversations there), under a session id the parent
// picks (`--session-id`) and returns in the result (`session_id`): the
// parent session's session uploader reads the child's session file from it and
// records the child as a sub-agent of the turn (session.child events,
// grandchildren included), the way the desktop records task sub-agents.
//
// Children get OMNIRUSH_PARENT_SESSION=<parent session id> in their
// environment: a child process uploads nothing itself.
//
// A child has no wall-clock limit unless the caller sets one: it runs until
// it is done. The only automatic stop is an inactivity watchdog (no output
// at all for CHILD_INACTIVITY_MS, longer while a tool runs). AgentManager
// runs children blocking or in the background and hands finished
// background results back for delivery to the parent.

import { spawn as nodeSpawn } from "node:child_process";
import { randomUUID } from "node:crypto";
import { existsSync } from "node:fs";
import { mkdtemp, writeFile, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import path from "node:path";

import {
  ENV_FALLBACK_EFFORT,
  ENV_FALLBACK_MODEL,
  FALLBACK_MARKER,
  fallbackNote,
  sourceNote,
  piThinking,
  type SubagentFallback,
} from "./subagents-lib";
import { childAuthEnv } from "./auth";
import { sanitizeToolEnvironment } from "./secret-env";
import { yoloActive } from "./yolo-lib";

/** childAuthEnv, never throwing (a spawn must not fail on the auth file). */
function childAuthEnvSafe(): Record<string, string> {
  try {
    return childAuthEnv(process.env);
  } catch {
    return {};
  }
}

/**
 * Inactivity watchdog default: a child that prints nothing (no model
 * stream, no tool event) for this long is treated as hung and stopped.
 * There is no wall-clock limit by default (a child works as long as it
 * needs); OMNIRUSH_AGENT_INACTIVITY_MINUTES overrides the window and 0
 * turns the watchdog off.
 */
export const CHILD_INACTIVITY_MS = 30 * 60_000;
/** While a child's tool call runs (a build, a polling `sleep`), its silence window is this many times longer. */
export const TOOL_INACTIVITY_FACTOR = 4;

/** The watchdog window from the environment (ms; 0 = off). */
export function childInactivityMs(env: NodeJS.ProcessEnv = process.env): number {
  const raw = String(env.OMNIRUSH_AGENT_INACTIVITY_MINUTES ?? "").trim().toLowerCase();
  if (!raw) return CHILD_INACTIVITY_MS;
  if (raw === "off" || raw === "0") return 0;
  const minutes = Number(raw);
  return Number.isFinite(minutes) && minutes > 0 ? Math.round(minutes * 60_000) : CHILD_INACTIVITY_MS;
}
/** Grace between SIGTERM and SIGKILL on timeout. */
export const CHILD_KILL_GRACE_MS = 5_000;
/** Per-child output cap in the structured result (bytes, UTF-8). */
export const CHILD_OUTPUT_CAP_BYTES = 50 * 1024;
/**
 * Finished sub-agents whose result reached the model keep their output (up
 * to CHILD_OUTPUT_CAP_BYTES each) for agents_result; beyond this many per
 * session, the oldest ones' output is let go (a swarm session runs
 * hundreds of children).
 */
export const MAX_KEPT_OUTPUTS = 200;
/** What a released output reads. */
export const RELEASED_OUTPUT = "(output no longer kept in memory: it was delivered to the conversation earlier)";

/**
 * Sub-agent layers below the main session (the desktop app's
 * OMNIRUSH_SUBAGENT_DEPTH): layers 1 and 2 may delegate, layer 3 may not.
 * The session capture records exactly these layers (session-sync.ts).
 */
export const MAX_SUBAGENT_DEPTH = 3;
/** The layer a process runs at, handed down to every child (+1 per layer). */
export const ENV_AGENT_DEPTH = "OMNIRUSH_AGENT_DEPTH";

/**
 * This process's sub-agent layer: 0 for a main session. A child of a parent
 * that predates the depth variable counts as layer 1.
 */
export function agentDepth(env: NodeJS.ProcessEnv = process.env): number {
  const text = String(env[ENV_AGENT_DEPTH] ?? "").trim();
  const raw = /^\d+$/.test(text) ? Number(text) : NaN;
  if (Number.isSafeInteger(raw)) return raw;
  return String(env.OMNIRUSH_PARENT_SESSION ?? "").trim() ? 1 : 0;
}

/** Whether a process at `depth` may start sub-agents. */
export function canDelegate(depth: number): boolean {
  return depth < MAX_SUBAGENT_DEPTH;
}

/**
 * The main session's spawn_agents guideline. Delegation is opt-in for
 * substantial independent work, rather than a default for every multi-step
 * request: each child repeats the parent's context and tool results.
 */
export const SPAWN_GUIDELINE = "Do the work yourself by default. Start a sub-agent only when the user asks for one or the task has at least two independent parts that each need substantial work. Start one child by default; use two only for genuinely independent tracks. Use three or more only when the user explicitly asks or there are that many independent tracks. Do not create separate search, implementation, test, or review children for the same code path. When the user names a number of sub-agents, start exactly that many. Make each task self-contained and limited to its part; include only the paths, symbols, URLs, and result format it needs. Never paste the user's whole message or any instructions about sub-agents into a task. Children do not delegate unless explicitly asked. For codebase exploration, search paths and symbols first, read targeted slices, and return a compact evidence summary instead of whole files.";

/** Main-session guidance for keeping repeated tool context small and useful. */
export const CONTEXT_GUIDELINE = "Keep context focused: search paths and symbols before reading files, use read offset/limit for large files, avoid rereading unchanged output, and summarize evidence after a tool-heavy phase. Preserve exact lines needed for edits and tests.";

/** A sub-agent's spawn_agents guideline (layer >= 1). */
export const SUBAGENT_SPAWN_GUIDELINE = "You are a sub-agent: do your whole task yourself, even when it is long. Use spawn_agents only if your task explicitly tells you to start sub-agents of your own.";

/** System text for a sub-agent: its layer, and that it does its task itself. */
export function subagentLayerNote(depth: number): string {
  const head = `You are a sub-agent: layer ${depth} of at most ${MAX_SUBAGENT_DEPTH} below the main session. Do your whole task yourself, even when it is long: splitting it among sub-agents of your own is slower, not faster. Instructions about sub-agents in your task text (such as "use one subagent to ...") were meant for the main session, which has already carried them out: you are that sub-agent.`;
  return canDelegate(depth)
    ? `${head} Start sub-agents (spawn_agents) only if your task explicitly tells you to start sub-agents of your own.`
    : `${head} You cannot delegate further (there is no spawn_agents at this layer).`;
}

export const AGENT_ROLES = ["code-searcher", "researcher-web", "general-worker"] as const;
export type AgentRole = (typeof AGENT_ROLES)[number];

/**
 * Resolve the model selected for the current parent session. Pi keeps the
 * effort level separately from the model id, while the CLI accepts the
 * canonical `model:effort` spelling for child argv.
 */
export function modelSelectionFromContext(ctx: any): string | undefined {
  const model = ctx?.model;
  const id = typeof model?.id === "string"
    ? model.id.trim()
    : typeof model?.modelID === "string"
      ? model.modelID.trim()
      : "";
  if (!id) return undefined;
  const effort = typeof ctx?.thinkingLevel === "string" ? ctx.thinkingLevel.trim() : "";
  if (!effort || effort === "off" || id.includes(":")) return id;
  return `${id}:${effort}`;
}

export interface RolePreset {
  role: AgentRole;
  label: string;
  description: string;
  /** Appended to the child's system prompt. */
  systemPrompt: string;
}

/**
 * Role presets. Guidance only — the child keeps every tool (hard tool
 * filters would break extension tools like web_fetch); the system
 * prompt sets expectations instead. Keep these prompts short because they
 * are appended to every child request.
 */
export const ROLE_PRESETS: Record<AgentRole, RolePreset> = {
  "code-searcher": {
    role: "code-searcher",
    label: "Code Searcher",
    description: "Read-only codebase exploration: locate code, trace call paths, summarize findings.",
    systemPrompt: [
      "You are a code-searcher subagent.",
      "Your job is READ-ONLY codebase exploration: find where things live, trace how calls flow, and report precise findings (file paths, symbol names, line references).",
      "Do NOT modify files, run state-changing commands, or start long-running processes.",
      "Start with path and symbol searches, then read only relevant slices; do not dump whole files or scan the repository without a reason.",
      "Report a short structured evidence summary: what you found, where (paths/lines), and any unanswered question — not a play-by-play of searches.",
    ].join(" "),
  },
  "researcher-web": {
    role: "researcher-web",
    label: "Researcher (Web)",
    description: "Web research via web_fetch: read documentation pages, articles and API responses.",
    systemPrompt: [
      "You are a researcher subagent with web access.",
      "Use the web_search tool to find pages (returns titles, URLs and snippets) and web_fetch to read the promising ones (HTTPS; pages stripped to text).",
      "Cross-check claims across more than one page when it matters, and prefer primary sources (official docs, spec pages) over blog summaries.",
      "Fetch only pages that can answer the task and report a compact synthesis with the source URLs you actually used — not a list of everything fetched.",
    ].join(" "),
  },
  "general-worker": {
    role: "general-worker",
    label: "General Worker",
    description: "Full-tool worker for self-contained implementation tasks (edits, commands, verification).",
    systemPrompt: [
      "You are a general-worker subagent with the full toolset.",
      "Complete the delegated task end to end: implement, verify (run the relevant tests/commands), and keep changes tight and reviewable.",
      "Stay inside the delegated task — do not refactor beyond it, do not touch unrelated files.",
      "Read task-scoped files first, keep tool output targeted, and avoid rereading unchanged files or dumping large generated files.",
      "Report what you changed (paths), how you verified it, and anything you deliberately left undone.",
    ].join(" "),
  },
};

/**
 * Resolve the child entry: the core's entry as the launcher started it
 * (OMNIRUSH_CORE_ENTRY: through <package>/engine, so process listings name
 * omnirush; the runtime reports the resolved path in argv[1]), else this
 * running script, with the current runtime.
 */
export function childInvocation(args: string[], env: NodeJS.ProcessEnv = process.env): {
  command: string;
  args: string[];
} {
  const launched = (env.OMNIRUSH_CORE_ENTRY || "").trim();
  if (launched && existsSync(launched)) return { command: process.execPath, args: [launched, ...args] };
  const currentScript = process.argv[1];
  const isBunVirtualScript = currentScript?.startsWith("/$bunfs/root/");
  if (currentScript && !isBunVirtualScript && existsSync(currentScript)) {
    return { command: process.execPath, args: [currentScript, ...args] };
  }
  const execName = path.basename(process.execPath).toLowerCase();
  const isGenericRuntime = /^(node|bun)(\.exe)?$/.test(execName);
  if (!isGenericRuntime) {
    return { command: process.execPath, args };
  }
  // Packaged bun builds: fall back to the omnirush on PATH (it forwards core flags).
  return { command: "omnirush", args };
}

/** The folder, inside a session dir, that holds its sub-agents' sessions. */
export const SUBAGENT_SESSIONS_DIR = "subagents";

/**
 * Where a child's session file goes: <parent's session dir>/subagents/.
 * The core's session picker (/resume) and --continue only read the
 * *.jsonl files directly in a session dir, so children stored one level
 * down never show up as sessions to resume. A sub-agent's own children
 * (nested) share its folder: a parent that is itself a sub-agent
 * (OMNIRUSH_PARENT_SESSION set) already runs there. Undefined when the
 * parent has no session dir (--no-session): the child keeps the core's
 * default, and the launcher tidies such stragglers (src/sessions.js).
 */
export function subagentSessionDir(parentSessionDir: unknown, env: NodeJS.ProcessEnv = process.env): string | undefined {
  if (typeof parentSessionDir !== "string" || !parentSessionDir.trim()) return undefined;
  if (String(env.OMNIRUSH_PARENT_SESSION ?? "").trim()) return parentSessionDir;
  return path.join(parentSessionDir, SUBAGENT_SESSIONS_DIR);
}

/**
 * Build the child's argv for one task: JSON print mode with the role's
 * system prompt appended via a temp file path (the caller writes and
 * cleans up that file — see withRolePromptFile). `model` selects a
 * cross-model child (pi's --provider/--model flags); the provider is
 * always our gateway extension.
 */
export function buildChildArgs(
  task: string,
  promptFilePath: string | null,
  model?: string,
  sessionId?: string,
  effort?: string | null,
  approve = false,
  sessionDir?: string,
): string[] {
  const args: string[] = ["--mode", "json", "-p"];
  if (approve) {
    // Yolo mode: a headless child has no trust prompt, so without this it
    // would ignore the project's settings the parent runs with.
    args.push("--approve");
  }
  if (sessionId) {
    // The child's session id, picked by the parent so its session file can
    // be found and captured as this session's sub-agent.
    args.push("--session-id", sessionId);
  }
  if (sessionDir) {
    // Out of the user's resume list (subagentSessionDir).
    args.push("--session-dir", sessionDir);
  }
  if (promptFilePath) {
    args.push("--append-system-prompt", promptFilePath);
  }
  if (model && model.trim()) {
    args.push("--provider", "omnirush", "--model", model.trim());
  }
  if (effort && effort.trim()) {
    // The effort the sub-agent runs on (the picked one, or the main agent's
    // mapped to the levels this model offers).
    args.push("--thinking", piThinking(effort.trim()));
  }
  args.push(`Task: ${task}`);
  return args;
}

/**
 * Write a role's system prompt to a private temp file and run `body`
 * with its path; the file is cleaned up afterwards either way.
 */
export async function withRolePromptFile<T>(
  role: AgentRole,
  body: (promptFilePath: string) => Promise<T>,
): Promise<T> {
  let dir: string | null = null;
  try {
    dir = await mkdtemp(path.join(tmpdir(), "omnirush-agent-"));
    const filePath = path.join(dir, `prompt-${role.replace(/[^\w.-]+/g, "_")}.md`);
    await writeFile(filePath, ROLE_PRESETS[role].systemPrompt, { encoding: "utf8", mode: 0o600 });
    return await body(filePath);
  } finally {
    if (dir) await rm(dir, { recursive: true, force: true }).catch(() => undefined);
  }
}

export interface ChildTask {
  role: AgentRole;
  task: string;
  /** Cross-model children: gateway model id the child runs on
   *  (e.g. "muse-spark-1.1" workers under an astra parent). Undefined =
   *  inherit the parent's model. */
  model?: string;
  /** The effort the child runs on (gateway spelling); undefined = pi's default. */
  effort?: string;
  /** The picked sub-agent model could not be used: `model` is the main one. */
  fallback?: SubagentFallback;
  /** The main model the child's gateway guard moves to when the gateway refuses `model`. */
  gatewayFallback?: { model: string; effort: string | null };
  /** Extra environment for the child (the sub-agent setting and main model, for nested layers). */
  env?: Record<string, string>;
  /** The model[:effort] the delegating agent asked for (the task's own `model`), when it named one. */
  requestedModel?: string;
  /** That model did not run: the user's /subagents choice won, or the gateway does not serve it. */
  override?: ChildModelOverride;
  /** Where `model` comes from ("subagents", "user_prompt", "task", "parent"). */
  modelSource?: string;
}

/** A task model that did not run, and what ran instead (the result's model_override). */
export interface ChildModelOverride {
  requested: string;
  used: string;
  effort: string | null;
  reason: string;
  note: string;
}

/** A sub-agent that ran on the main model instead of the picked one. */
export interface ChildModelFallback {
  requested: string;
  used: string;
  effort?: string | null;
  reason: string;
  /** "selection": picked before it started; "gateway": the gateway refused the picked model mid-run. */
  kind: "selection" | "gateway";
  note: string;
}

export type ChildStatus = "completed" | "failed" | "timeout" | "stalled" | "cancelled" | "interrupted";

export interface ChildResult {
  role: AgentRole;
  task: string;
  /** The child's pi session id (its session file is `<ts>_<id>.jsonl`). */
  session_id?: string;
  /** The gateway model the child ran on (cross-model children). */
  model?: string;
  /** The effort it ran on. */
  effort?: string;
  /** It ran on the main model instead of the picked one (why, and since when). */
  model_fallback?: ChildModelFallback;
  /** The task's own model did not run (the user's /subagents choice, or not served). */
  model_override?: ChildModelOverride;
  /** Where the model comes from: "subagents", "user_prompt", "task" or "parent". */
  model_source?: string;
  /** The model[:effort] the task asked for, when it named one. */
  requested_model?: string;
  status: ChildStatus;
  /** Process exit code (null when killed by a signal or still unknown). */
  exitCode: number | null;
  /** Final assistant text (or the best available partial output). */
  output: string;
  /** True when `output` was capped (full stdout events were larger). */
  outputTruncated: boolean;
  durationMs: number;
  /** Assistant-turn count seen in the child's event stream. */
  turns: number;
  error?: string;
}

export interface ChildResultInput {
  exitCode: number | null;
  finalText: string;
  turns: number;
  status?: ChildStatus;
  error?: string;
  /** The gateway moved the child to the main model mid-run. */
  gatewayFallback?: { requested: string; used: string; effort?: string | null; reason: string };
}

/** The selection or gateway fallback of a child, for its result. */
export function childModelFallback(task: ChildTask, gateway?: ChildResultInput["gatewayFallback"]): ChildModelFallback | undefined {
  if (task.fallback) {
    return {
      requested: task.fallback.requested,
      used: task.fallback.used,
      effort: task.effort ?? null,
      reason: String(task.fallback.reason),
      kind: "selection",
      note: fallbackNote(task.fallback),
    };
  }
  if (gateway) {
    return {
      requested: gateway.requested,
      used: gateway.used,
      effort: gateway.effort ?? null,
      reason: gateway.reason,
      kind: "gateway",
      note: fallbackNote({ requested: gateway.requested, used: gateway.used, reason: gateway.reason }),
    };
  }
  return undefined;
}

/** Final structured result for one child. */
export function buildChildResult(
  task: ChildTask,
  input: ChildResultInput,
  durationMs: number,
  sessionId?: string,
): ChildResult {
  const output = input.finalText;
  let capped = output;
  const truncated = Buffer.byteLength(output, "utf8") > CHILD_OUTPUT_CAP_BYTES;
  if (truncated) {
    capped = output.slice(0, CHILD_OUTPUT_CAP_BYTES);
    while (Buffer.byteLength(capped, "utf8") > CHILD_OUTPUT_CAP_BYTES) capped = capped.slice(0, -1);
  }
  const modelFallback = childModelFallback(task, input.gatewayFallback);
  return {
    role: task.role,
    task: task.task,
    ...(sessionId ? { session_id: sessionId } : {}),
    ...(task.model ? { model: task.model } : {}),
    ...(task.effort ? { effort: task.effort } : {}),
    ...(modelFallback ? { model_fallback: modelFallback } : {}),
    ...(task.override ? { model_override: { ...task.override } } : {}),
    ...(task.modelSource ? { model_source: task.modelSource } : {}),
    ...(task.requestedModel ? { requested_model: task.requestedModel } : {}),
    status: input.status ?? (input.exitCode === 0 ? "completed" : "failed"),
    exitCode: input.exitCode,
    output: capped,
    outputTruncated: truncated,
    durationMs,
    turns: input.turns,
    ...(input.error ? { error: input.error } : {}),
  };
}

// --- child event-stream parsing -------------------------------------------

export interface ParsedChildEvents {
  /** Final assistant text (last assistant message with text parts). */
  finalText: string;
  /** Assistant messages seen (turn count). */
  turns: number;
}

/** Characters of an event line read to tell its type (and role / tool call id) apart. */
export const CHILD_EVENT_HEAD_CHARS = 512;
/**
 * Longest event line kept whole for parsing (characters). Only an
 * assistant `message_end` is ever kept; one longer than this (a huge tool
 * call argument) still counts as a turn, its text is not read.
 */
export const MAX_CHILD_EVENT_LINE_CHARS = 4 * 1024 * 1024;

const EVENT_TYPE_RE = /^\s*\{\s*"type"\s*:\s*"([A-Za-z_]+)"/;
const TOOL_CALL_ID_RE = /"toolCallId"\s*:\s*"((?:[^"\\]|\\.)*)"/;
const ROLE_RE = /^\s*\{\s*"type"\s*:\s*"message_end"\s*,\s*"message"\s*:\s*\{\s*"role"\s*:\s*"([A-Za-z_]+)"/;

/**
 * Streaming reader of a `pi --mode json -p` child's stdout: one JSON event
 * per line. Only what the parent needs is kept — the last assistant text,
 * the assistant turn count and the tool calls in flight — and every other
 * event (streaming deltas, tool results, turn_end / agent_end, which repeat
 * the whole conversation) is skipped as it streams by, without buffering
 * it. A child's stream is several times its conversation (agent_end alone
 * repeats all of it), so keeping it whole cost the parent that much memory
 * per child. Unparseable lines (banners, noise) are ignored.
 */
export class ChildEventScanner {
  finalText = "";
  turns = 0;
  /** Monotonic count of model/tool events that represent real progress. */
  progressSerial = 0;
  readonly toolsRunning = new Set<string>();
  private line = "";
  /** head: deciding from the first characters; keep: buffering the whole line; skip: dropping it. */
  private mode: "head" | "keep" | "skip" = "head";
  private lineType: string | null = null;
  private lineRole: string | null = null;

  /** Feed a decoded chunk of stdout. */
  push(text: string): void {
    let start = 0;
    while (start <= text.length) {
      const newline = text.indexOf("\n", start);
      const end = newline < 0 ? text.length : newline;
      if (end > start) this.feed(text.slice(start, end));
      if (newline < 0) break;
      this.endLine();
      start = newline + 1;
    }
  }

  /** The stream ended: a last line without a newline still counts. */
  end(): void {
    if (this.line || this.mode !== "head") this.endLine();
  }

  private feed(part: string): void {
    if (this.mode === "skip") return;
    this.line += part;
    if (this.mode === "head" && this.line.length >= CHILD_EVENT_HEAD_CHARS) this.decide();
    if (this.mode === "keep" && this.line.length > MAX_CHILD_EVENT_LINE_CHARS) {
      // Too long to keep: an assistant message still counts as a turn.
      if (this.lineType === "message_end" && this.lineRole === "assistant") this.turns += 1;
      this.line = "";
      this.mode = "skip";
    }
  }

  /** Classify the line from its head: keep it, take what it says from the head, or skip it. */
  private decide(): void {
    const head = this.line.slice(0, CHILD_EVENT_HEAD_CHARS);
    const type = EVENT_TYPE_RE.exec(head)?.[1] ?? null;
    this.lineType = type;
    if (type === "message_update" || type === "message_start" || type === "turn_start") this.progressSerial += 1;
    if (type === "message_end") {
      const role = ROLE_RE.exec(head)?.[1] ?? null;
      this.lineRole = role;
      // Only an assistant message carries the final text; a role further
      // into the object than the head (not how pi writes it) is kept too.
      if (role !== null && role !== "assistant") this.skip();
      else this.mode = "keep";
      return;
    }
    if (type === "tool_execution_start" || type === "tool_execution_end") {
      const id = TOOL_CALL_ID_RE.exec(head)?.[1];
      if (id !== undefined) {
        this.tool(type, JSON.parse(`"${id}"`));
        this.skip();
      } else {
        this.mode = "keep";
      }
      return;
    }
    this.skip();
  }

  private skip(): void {
    this.line = "";
    this.mode = "skip";
  }

  private tool(type: string, id: string): void {
    this.progressSerial += 1;
    if (type === "tool_execution_start") this.toolsRunning.add(id);
    else this.toolsRunning.delete(id);
  }

  private endLine(): void {
    const line = this.line;
    const mode = this.mode;
    this.line = "";
    this.mode = "head";
    this.lineType = null;
    this.lineRole = null;
    if (mode === "skip") return;
    const trimmed = line.trim();
    if (!trimmed || !trimmed.includes('"type"')) return;
    // A short line: skip the types nothing is read from without parsing them.
    const type = mode === "head" ? EVENT_TYPE_RE.exec(trimmed)?.[1] : undefined;
    if (type && type !== "message_end" && type !== "message_update" && type !== "message_start" && type !== "turn_start" && type !== "tool_execution_start" && type !== "tool_execution_end") return;
    let event: any;
    try {
      event = JSON.parse(trimmed);
    } catch {
      return;
    }
    if ((event?.type === "tool_execution_start" || event?.type === "tool_execution_end") && typeof event.toolCallId === "string") {
      this.tool(event.type, event.toolCallId);
      return;
    }
    if (event?.type === "message_update" || event?.type === "message_start" || event?.type === "turn_start") {
      this.progressSerial += 1;
      return;
    }
    if (event?.type !== "message_end" || !event.message) return;
    const message = event.message;
    if (message.role !== "assistant") return;
    this.turns += 1;
    this.progressSerial += 1;
    const parts = Array.isArray(message.content) ? message.content : [];
    const text = parts
      .filter((part: any) => part?.type === "text" && typeof part.text === "string")
      .map((part: any) => part.text)
      .join("\n\n");
    if (text.trim()) this.finalText = text;
    else if (typeof message.content === "string" && message.content.trim()) {
      this.finalText = message.content;
    }
  }
}

/**
 * Parse the stdout of a `pi --mode json -p` run: one JSON event per line;
 * `message_end` events carry the messages (see ChildEventScanner, which
 * the runner feeds as the output streams).
 */
export function parseChildEvents(stdout: string): ParsedChildEvents {
  const scanner = new ChildEventScanner();
  scanner.push(stdout);
  scanner.end();
  return { finalText: scanner.finalText, turns: scanner.turns };
}

// --- the runner -------------------------------------------------------------

export interface SpawnChildOptions {
  cwd: string;
  parentSessionId: string;
  /** Where the child's session file goes (subagentSessionDir); the core's default when unset. */
  sessionDir?: string;
  /**
   * Explicit wall-clock limit (the tool's `timeout_minutes`). Undefined =
   * none: a child runs until it is done, stalls (inactivityMs) or is
   * cancelled.
   */
  timeoutMs?: number;
  /**
   * Inactivity watchdog: a child that prints no event for this long is
   * treated as hung and killed (status "stalled"). While one of its tools
   * is running (a long build, a `sleep` poll), the window is
   * TOOL_INACTIVITY_FACTOR times longer. Default: childInactivityMs(env);
   * 0 disables the watchdog.
   */
  inactivityMs?: number;
  /**
   * Kill the child when this fires (cancellation). An abort whose reason is
   * "interrupted" (the session ended under it) reports status
   * "interrupted"; any other abort reports "cancelled".
   */
  signal?: AbortSignal;
  /** Grace between SIGTERM and SIGKILL when killing (default CHILD_KILL_GRACE_MS). */
  killGraceMs?: number;
  spawnImpl?: typeof nodeSpawn;
  now?: () => number;
  /** Progress callback: partial output lines as the child runs. */
  onChildStdout?: (role: AgentRole, chunk: string) => void;
  /** Meaningful model/tool progress from the child, with its event types. */
  onActivity?: (activity: ChildActivity) => void;
  /** The child's session id (a fresh UUID by default). */
  sessionId?: string;
  /** The child process started (its pid), for memory accounting. */
  onSpawn?: (pid: number | undefined) => void;
  /** Extra environment for the child (e.g. the memory tree root). */
  extraEnv?: Record<string, string>;
}

export interface ChildActivity {
  at: number;
  /** Assistant messages finished so far. */
  turns: number;
  /** Tool calls running right now. */
  toolsRunning: number;
  meaningful: true;
}

/**
 * The session uploader reaches the running sub-agents through this global (it must
 * stop them before its last capture; pi runs its session_shutdown handler
 * before the agents extension's).
 */
export const AGENTS_HOOK = Symbol.for("omnirush.agents");
export interface AgentsHook {
  /** Stop every unfinished child (of one session, or all); resolves with them once they exited. */
  interruptAll(parentSessionId?: string): Promise<Array<{ sessionId: string; id: string; role: string }>>;
}

/** Reason passed to AbortController.abort() when the session ends under a child. */
export const INTERRUPTED_REASON = "interrupted";

/**
 * Run ONE child agent to completion. No wall-clock limit unless the caller
 * sets timeoutMs; the only automatic kill is the inactivity watchdog (a
 * child silent for the whole window). The child inherits the parent's
 * environment plus OMNIRUSH_PARENT_SESSION. On a kill: SIGTERM, a short
 * grace, then SIGKILL; whatever the child printed so far is kept as the
 * partial result.
 */
export async function runChildAgent(
  task: ChildTask,
  options: SpawnChildOptions,
): Promise<ChildResult> {
  const now = options.now ?? Date.now;
  const startedAt = now();
  const timeoutMs = typeof options.timeoutMs === "number" && options.timeoutMs > 0
    ? Math.max(1_000, options.timeoutMs)
    : null;
  const inactivityMs = options.inactivityMs ?? childInactivityMs();
  const spawnImpl = options.spawnImpl ?? nodeSpawn;
  const sessionId = options.sessionId ?? randomUUID();

  return withRolePromptFile(task.role, async (promptFile) => {
    const args = buildChildArgs(task.task, promptFile, task.model, sessionId, task.effort, yoloActive(), options.sessionDir);
    const invocation = childInvocation(args);

    return await new Promise<ChildResult>((resolvePromise) => {
      // The event stream is read as it arrives; nothing of it is kept but
      // what the result needs (the scanner), so a long child costs the
      // parent next to nothing.
      const events = new ChildEventScanner();
      let stderr = "";
      let pendingErrLine = "";
      let gatewayFallback: ChildResultInput["gatewayFallback"];
      let settled = false;
      let killedFor: "timeout" | "stalled" | "cancelled" | "interrupted" | null = null;
      const toolsRunning = events.toolsRunning;
      let lastActivity = now();
      let watchdog: ReturnType<typeof setTimeout> | null = null;
      let timer: ReturnType<typeof setTimeout> | null = null;

      const child = spawnImpl(invocation.command, invocation.args, {
        cwd: options.cwd,
        shell: false,
        detached: process.platform !== "win32",
        stdio: ["ignore", "pipe", "pipe"],
        // Shared credentials by location (OMNIRUSH_DIR), never a token
        // frozen at the parent's launch: the child re-reads auth.json for
        // every request and takes part in the refresh lock.
        env: childEnvironment(task, options.parentSessionId, { ...process.env, ...childAuthEnvSafe(), ...(options.extraEnv ?? {}) }),
      });
      options.onSpawn?.(child.pid);

      const finish = (input: ChildResultInput) => {
        if (settled) return;
        settled = true;
        if (timer) clearTimeout(timer);
        if (watchdog) clearTimeout(watchdog);
        if (options.signal && abortHandler) {
          options.signal.removeEventListener("abort", abortHandler);
        }
        resolvePromise(buildChildResult(task, { ...input, ...(gatewayFallback ? { gatewayFallback } : {}) }, now() - startedAt, sessionId));
      };

      const killTree = () => {
        const signalGroup = (signal: NodeJS.Signals) => {
          if (process.platform !== "win32" && typeof child.pid === "number" && child.pid > 0) {
            try {
              process.kill(-child.pid, signal);
              return true;
            } catch {
              /* fall back to the direct child below */
            }
          }
          if (process.platform === "win32" && typeof child.pid === "number" && child.pid > 0) {
            try {
              const args = ["/PID", String(child.pid), "/T"];
              if (signal === "SIGKILL") args.push("/F");
              nodeSpawn("taskkill", args, { windowsHide: true, stdio: "ignore" });
              return true;
            } catch {
              /* fall back to the direct child below */
            }
          }
          return false;
        };
        if (!signalGroup("SIGTERM")) {
          try { child.kill("SIGTERM"); } catch { /* already gone */ }
        }
        setTimeout(() => {
          if (!signalGroup("SIGKILL")) {
            try { child.kill("SIGKILL"); } catch { /* already gone */ }
          }
        }, Math.max(50, options.killGraceMs ?? CHILD_KILL_GRACE_MS)).unref?.();
      };

      const kill = (reason: NonNullable<typeof killedFor>) => {
        if (settled || killedFor) return;
        killedFor = reason;
        killTree();
      };

      // The watchdog re-arms on every sign of life; it fires only after a
      // full window of silence (longer while a tool is running).
      const armWatchdog = () => {
        if (!(inactivityMs > 0) || settled || killedFor) return;
        if (watchdog) clearTimeout(watchdog);
        const window = toolsRunning.size > 0 ? inactivityMs * TOOL_INACTIVITY_FACTOR : inactivityMs;
        watchdog = setTimeout(() => kill("stalled"), window);
        watchdog.unref?.();
      };
      const alive = (meaningful: boolean) => {
        if (!meaningful) return;
        lastActivity = now();
        armWatchdog();
        options.onActivity?.({ at: lastActivity, turns: events.turns, toolsRunning: toolsRunning.size, meaningful: true });
      };
      armWatchdog();

      if (timeoutMs !== null) {
        timer = setTimeout(() => kill("timeout"), timeoutMs);
        timer.unref?.();
      }

      const abortHandler = options.signal
        ? () => kill(options.signal?.reason === INTERRUPTED_REASON ? "interrupted" : "cancelled")
        : null;
      if (options.signal && abortHandler) {
        if (options.signal.aborted) abortHandler();
        else options.signal.addEventListener("abort", abortHandler, { once: true });
      }

      child.stdout?.setEncoding?.("utf8");
      child.stdout?.on("data", (chunk: Buffer | string) => {
        const text = String(chunk);
        const before = events.progressSerial;
        events.push(text);
        alive(events.progressSerial !== before);
        options.onChildStdout?.(task.role, text);
      });
      child.stderr?.setEncoding?.("utf8");
      child.stderr?.on("data", (chunk: Buffer | string) => {
        const text = String(chunk);
        // The child's gateway guard reports a move to the main model here.
        const lines = (pendingErrLine + text).split("\n");
        pendingErrLine = lines.pop() ?? "";
        if (pendingErrLine.length > 64 * 1024) pendingErrLine = pendingErrLine.slice(-64 * 1024);
        let meaningful = false;
        for (const line of lines) {
          const fallback = parseFallbackMarker(line);
          if (fallback) {
            gatewayFallback = fallback;
            meaningful = true;
          }
        }
        stderr += text;
        if (stderr.length > 256 * 1024) stderr = stderr.slice(-64 * 1024);
        alive(meaningful);
      });
      child.on("error", (error: Error) => {
        finish({
          exitCode: null,
          finalText: "",
          turns: 0,
          status: "failed",
          error: `failed to start: ${error.message}`,
        });
      });
      child.on("close", (code: number | null) => {
        events.end();
        const parsed = { finalText: events.finalText, turns: events.turns };
        const partial = parsed.finalText ? " — partial result kept" : "";
        if (killedFor) {
          const error = killedFor === "timeout"
            ? `timed out after ${formatMinutes(timeoutMs ?? 0)}${partial}`
            : killedFor === "stalled"
              ? `stopped: no activity for ${formatMinutes(toolsRunning.size > 0 ? inactivityMs * TOOL_INACTIVITY_FACTOR : inactivityMs)}${partial}`
              : killedFor === "interrupted"
                ? `interrupted: the parent session ended${partial}`
                : `cancelled by the parent session${partial}`;
          finish({ exitCode: code, finalText: parsed.finalText, turns: parsed.turns, status: killedFor, error });
          return;
        }
        finish({
          exitCode: code,
          finalText: parsed.finalText || (code !== 0 ? stderr.slice(-2_000) : ""),
          turns: parsed.turns,
          status: code === 0 ? "completed" : "failed",
          ...(code !== 0 && !parsed.finalText && stderr.trim() ? { error: stderr.trim().slice(-500) } : {}),
        });
      });
    });
  });
}

function formatMinutes(ms: number): string {
  const minutes = ms / 60_000;
  if (minutes >= 1) return `${Math.round(minutes * 10) / 10} min`;
  return `${Math.max(1, Math.round(ms / 1000))} s`;
}

/**
 * A child's environment: the parent's, the parent session id (the child
 * uploads nothing itself), the sub-agent setting and main model handed down
 * to nested layers, and the main model its gateway guard falls back to (only
 * for this child: a nested one gets its own or none).
 */
export function childEnvironment(task: ChildTask, parentSessionId: string, base: NodeJS.ProcessEnv = process.env): NodeJS.ProcessEnv {
  const env: NodeJS.ProcessEnv = { ...base, ...(task.env ?? {}), OMNIRUSH_PARENT_SESSION: parentSessionId };
  // One layer further down than the delegating agent.
  env[ENV_AGENT_DEPTH] = String(agentDepth(base) + 1);
  // Guarded mode: the name the parent's approval prompt shows for this child.
  const oneLine = task.task.replace(/\s+/g, " ").trim();
  env.OMNIRUSH_SUBAGENT_LABEL = `${ROLE_PRESETS[task.role]?.label ?? task.role}: ${oneLine.length > 60 ? `${oneLine.slice(0, 59)}…` : oneLine}`;
  delete env[ENV_FALLBACK_MODEL];
  delete env[ENV_FALLBACK_EFFORT];
  if (task.gatewayFallback?.model) {
    env[ENV_FALLBACK_MODEL] = task.gatewayFallback.model;
    if (task.gatewayFallback.effort) env[ENV_FALLBACK_EFFORT] = task.gatewayFallback.effort;
  }
  return env;
}

/** A child's stderr line reporting its gateway fallback (see sota.ts), or null. */
export function parseFallbackMarker(line: string): ChildResultInput["gatewayFallback"] | null {
  const at = line.indexOf(FALLBACK_MARKER);
  if (at < 0) return null;
  try {
    const parsed = JSON.parse(line.slice(at + FALLBACK_MARKER.length));
    if (typeof parsed?.requested !== "string" || typeof parsed?.used !== "string") return null;
    return {
      requested: parsed.requested,
      used: parsed.used,
      effort: typeof parsed.effort === "string" ? parsed.effort : null,
      reason: typeof parsed.reason === "string" ? parsed.reason : "refused",
    };
  } catch {
    return null;
  }
}

/**
 * Run tasks with a concurrency cap, preserving input order in the
 * results (ports the pi subagent example's mapWithConcurrencyLimit).
 */
export async function mapWithConcurrency<TIn, TOut>(
  items: TIn[],
  concurrency: number,
  fn: (item: TIn, index: number) => Promise<TOut>,
): Promise<TOut[]> {
  if (items.length === 0) return [];
  const limit = Math.max(1, Math.min(concurrency, items.length));
  const results: TOut[] = new Array(items.length);
  let nextIndex = 0;
  const workers = Array.from({ length: limit }, async () => {
    for (;;) {
      const current = nextIndex++;
      if (current >= items.length) return;
      results[current] = await fn(items[current], current);
    }
  });
  await Promise.all(workers);
  return results;
}

/** The model and effort a child really ran on (after a gateway fallback), and why that one. */
export function ranOn(result: Pick<ChildResult, "model" | "effort" | "model_fallback">): { model: string | null; effort: string | null } {
  if (result.model_fallback?.kind === "gateway") return { model: result.model_fallback.used, effort: result.model_fallback.effort ?? null };
  return { model: result.model ?? null, effort: result.effort ?? null };
}

function ranOnText(result: ChildResult): string {
  const { model, effort } = ranOn(result);
  if (!model) return "the parent's model (not an omnirush.ai model)";
  const why = result.model_fallback?.kind === "gateway"
    ? "the gateway refused the chosen model"
    : result.model_fallback
      ? "the main model"
      : sourceNote(result.model_source);
  return `${model}${effort ? ` · ${effort}` : ""}${why ? ` (${why})` : ""}`;
}

/** Structured result text for the parent model (one section per child). */
export function renderChildResults(results: ChildResult[], ids?: readonly string[]): string {
  const succeeded = results.filter((result) => result.status === "completed").length;
  const sections = results.map((result, index) => {
    const minutes = Math.round((result.durationMs / 60_000) * 10) / 10;
    const ran = ranOn(result);
    const label = ran.model ? ` [${ran.model}${ran.effort ? ` · ${ran.effort}` : ""}]` : "";
    const header = `### ${ids?.[index] ? `${ids[index]} ` : ""}${result.role}${label} — ${result.status} (${minutes} min${result.outputTruncated ? ", output capped" : ""})`;
    const meta: string[] = [`task: ${result.task}`];
    meta.push(`ran on: ${ranOnText(result)}`);
    if (result.model_override) meta.push(`note: ${result.model_override.note}`);
    if (result.model_fallback) meta.push(`note: ${result.model_fallback.note}`);
    if (result.error) meta.push(`error: ${result.error}`);
    return `${header}\n${meta.join("\n")}\n\n${result.output || "(no output)"}`;
  });
  return `spawn_agents: ${succeeded}/${results.length} completed\n\n${sections.join("\n\n---\n\n")}`;
}

// --- sub-agent manager: blocking and background dispatch --------------------

export type AgentState = "queued" | "running" | ChildStatus;

/** One sub-agent of a session, running or finished. */
export interface AgentRecord {
  /** Short handle the model uses with the agents_* tools ("a1", "a2", ...). */
  id: string;
  /** The spawn_agents call it came from ("b1", ...). */
  batch: string;
  parentSessionId: string;
  /** The child's pi session id (its trace is captured under it). */
  sessionId: string;
  role: AgentRole;
  task: string;
  model?: string;
  effort?: string;
  /** Dispatched with wait:false: its result comes back as a message. */
  background: boolean;
  /** Null means the batch was intentionally uncapped; otherwise the batch limit. */
  parallelLimit: number | null;
  /** Background delivery: one message per child, or one per batch. */
  notify: "each" | "batch";
  status: AgentState;
  /** Queued because memory is short (see memory-lib.ts): why, and since when. */
  waiting: { reason: "memory"; detail: string; since: number } | null;
  queuedAt: number;
  startedAt: number | null;
  finishedAt: number | null;
  lastActivityAt: number;
  turns: number;
  toolsRunning: number;
  result: ChildResult | null;
  /** Its result reached the model (tool result or delivered message). */
  delivered: boolean;
  controller: AbortController;
  done: Promise<ChildResult>;
}

export type AgentRunner = (
  task: ChildTask,
  options: {
    sessionId: string;
    parentSessionId: string;
    /** Where the child's session file goes (subagentSessionDir). */
    sessionDir?: string;
    signal: AbortSignal;
    onActivity: (activity: ChildActivity) => void;
    /** The call's explicit timeout_minutes, if any. */
    timeoutMs?: number;
    cwd?: string;
    /** The child process started (its pid). */
    onSpawn?: (pid: number | undefined) => void;
  },
) => Promise<ChildResult>;

/**
 * Memory-aware admission (memory-lib.ts MemoryAdmission): whether one more
 * child may start, and the children it counts.
 */
export interface AgentAdmission {
  check(): { ok: true } | { ok: false; reason: string; detail: string };
  started(key: object, pid?: number | null): void;
  spawned(key: object, pid: number | null | undefined): void;
  finished(key: object): void;
}

/** A queued child: waits for its batch's max_parallel slot and for memory. */
interface QueueEntry {
  record: AgentRecord;
  task: ChildTask;
  options: DispatchOptions;
  batch: { running: number; limit: number };
}

export interface DispatchOptions {
  background: boolean;
  /** Children running at once (default: all of them). */
  concurrency?: number;
  notify?: "each" | "batch";
  /** Blocking calls: the tool's signal cancels the children. */
  signal?: AbortSignal;
  /** A child settled (progress updates). */
  onSettled?: (record: AgentRecord) => void;
  /** A child was queued, started, updated, or settled. */
  onChanged?: () => void;
  /** Explicit wall-clock limit per child (none by default). */
  timeoutMs?: number;
  /** Workspace the children run in. */
  cwd?: string;
  /** Where the children's session files go (subagentSessionDir). */
  sessionDir?: string;
}

/**
 * Tracks every sub-agent of this process per parent session: runs them
 * (through an injected runner), answers status/wait/result/cancel, and
 * queues finished background results for delivery back to the parent
 * (`takeDeliverable`, signalled through `onDeliverable`).
 */
export class AgentManager {
  private readonly runner: AgentRunner;
  private readonly now: () => number;
  private readonly onDeliverable: () => void;
  private readonly onChanged: () => void;
  private readonly bySession = new Map<string, AgentRecord[]>();
  private waiters: Array<{ records: AgentRecord[]; resolve: () => void }> = [];
  private changeWaiters: Array<() => void> = [];
  private readonly resolvers = new Map<AgentRecord, (result: ChildResult) => void>();
  private nextAgent = 1;
  private nextBatch = 1;
  private readonly admission: AgentAdmission | null;
  private readonly retryMs: number;
  private readonly onMemoryWait: (info: { critical: boolean; detail: string; queued: number } | null) => void;
  private queue: QueueEntry[] = [];
  private retryTimer: ReturnType<typeof setTimeout> | null = null;
  /** The last memory wait reported (null when nothing waits for memory). */
  private memoryWait: { critical: boolean; detail: string } | null = null;

  constructor(options: {
    run: AgentRunner;
    now?: () => number;
    onDeliverable?: () => void;
    onChanged?: () => void;
    /** Memory-aware admission; none = children start as soon as their batch allows. */
    admission?: AgentAdmission | null;
    /** How often a child waiting for memory checks again (ms). */
    retryMs?: number;
    /** Children started or stopped waiting for memory (null: nothing waits any more). */
    onMemoryWait?: (info: { critical: boolean; detail: string; queued: number } | null) => void;
  }) {
    this.runner = options.run;
    this.now = options.now ?? Date.now;
    this.onDeliverable = options.onDeliverable ?? (() => undefined);
    this.onChanged = options.onChanged ?? (() => undefined);
    this.admission = options.admission ?? null;
    this.retryMs = options.retryMs ?? 1_000;
    this.onMemoryWait = options.onMemoryWait ?? (() => undefined);
  }

  dispatch(
    parentSessionId: string,
    tasks: ChildTask[],
    options: DispatchOptions,
  ): { batch: string; records: AgentRecord[]; done: Promise<ChildResult[]> } {
    const batch = `b${this.nextBatch++}`;
    this.releaseOldOutputs(parentSessionId);
    const list = this.bySession.get(parentSessionId) ?? [];
    this.bySession.set(parentSessionId, list);
    const at = this.now();
    const records = tasks.map((task) => {
      let resolveDone!: (result: ChildResult) => void;
      const done = new Promise<ChildResult>((resolve) => { resolveDone = resolve; });
      const record: AgentRecord = {
        id: `a${this.nextAgent++}`,
        batch,
        parentSessionId,
        sessionId: randomUUID(),
        role: task.role,
        task: task.task,
        ...(task.model ? { model: task.model } : {}),
        ...(task.effort ? { effort: task.effort } : {}),
        background: options.background,
        parallelLimit: options.concurrency === undefined ? null : Math.max(1, Math.floor(options.concurrency)),
        notify: options.notify ?? "batch",
        status: "queued",
        waiting: null,
        queuedAt: at,
        startedAt: null,
        finishedAt: null,
        lastActivityAt: at,
        turns: 0,
        toolsRunning: 0,
        result: null,
        // A blocking call returns its results itself.
        delivered: !options.background,
        controller: new AbortController(),
        done,
      };
      this.resolvers.set(record, resolveDone);
      list.push(record);
      return record;
    });
    this.onChanged();
    options.onChanged?.();
    if (options.signal) {
      const abort = () => { void this.cancel(records); };
      if (options.signal.aborted) abort();
      else options.signal.addEventListener("abort", abort, { once: true });
    }
    // No count cap: every task of the batch may run at once unless the call
    // set max_parallel; memory admission decides when each one starts.
    const slots = { running: 0, limit: Math.max(1, Math.floor(options.concurrency ?? records.length)) };
    records.forEach((record, index) => this.queue.push({ record, task: tasks[index], options, batch: slots }));
    this.pump();
    return { batch, records, done: Promise.all(records.map((record) => record.done)) };
  }

  /**
   * Start queued children, oldest first, while their batch has a free slot
   * and memory admission allows; the rest wait for a child to finish or
   * for the next check (every retryMs while something waits for memory).
   */
  private pump(): void {
    if (this.retryTimer) {
      clearTimeout(this.retryTimer);
      this.retryTimer = null;
    }
    let blocked: { reason: string; detail: string } | null = null;
    const still: QueueEntry[] = [];
    let changed = false;
    for (const entry of this.queue) {
      const { record } = entry;
      if (record.result) {
        // Cancelled while it waited: it never starts.
        entry.options.onSettled?.(record);
        continue;
      }
      if (blocked || entry.batch.running >= entry.batch.limit) {
        if (blocked && entry.batch.running < entry.batch.limit) changed = this.markWaiting(record, blocked.detail) || changed;
        still.push(entry);
        continue;
      }
      const verdict = this.admission ? this.admission.check() : { ok: true as const };
      if (!verdict.ok) {
        blocked = verdict;
        changed = this.markWaiting(record, verdict.detail) || changed;
        still.push(entry);
        continue;
      }
      entry.batch.running += 1;
      void this.start(entry);
    }
    this.queue = still;
    const waiting = still.filter((entry) => entry.record.waiting);
    if (blocked && waiting.length > 0) {
      const info = { critical: blocked.reason === "critical", detail: blocked.detail };
      // Once when children start waiting, and when it turns critical (or back):
      // the figures in the detail change with every check.
      if (!this.memoryWait || this.memoryWait.critical !== info.critical) {
        this.memoryWait = info;
        this.onMemoryWait({ ...info, queued: waiting.length });
      }
      // A ref'd timer: a one-shot run must not end while children wait.
      this.retryTimer = setTimeout(() => this.pump(), this.retryMs);
    } else if (this.memoryWait) {
      this.memoryWait = null;
      this.onMemoryWait(null);
    }
    if (changed) this.onChanged();
  }

  private markWaiting(record: AgentRecord, detail: string): boolean {
    if (record.waiting?.detail === detail) return false;
    record.waiting = { reason: "memory", detail, since: record.waiting?.since ?? this.now() };
    return true;
  }

  private async start(entry: QueueEntry): Promise<void> {
    const { record, task, options } = entry;
    record.status = "running";
    record.waiting = null;
    record.startedAt = this.now();
    record.lastActivityAt = record.startedAt;
    this.admission?.started(record);
    this.onChanged();
    options.onChanged?.();
    let result: ChildResult;
    try {
      result = await this.runner(task, {
        sessionId: record.sessionId,
        parentSessionId: record.parentSessionId,
        signal: record.controller.signal,
        onSpawn: (pid) => this.admission?.spawned(record, pid),
        ...(options.timeoutMs !== undefined ? { timeoutMs: options.timeoutMs } : {}),
        ...(options.cwd ? { cwd: options.cwd } : {}),
        ...(options.sessionDir ? { sessionDir: options.sessionDir } : {}),
        onActivity: (activity) => {
          record.lastActivityAt = activity.at;
          record.turns = activity.turns;
          record.toolsRunning = activity.toolsRunning;
          this.onChanged();
          options.onChanged?.();
        },
      });
    } catch (error) {
      result = buildChildResult(task, {
        exitCode: null,
        finalText: "",
        turns: record.turns,
        status: "failed",
        error: error instanceof Error ? error.message : String(error),
      }, this.now() - (record.startedAt ?? this.now()), record.sessionId);
    }
    this.admission?.finished(record);
    entry.batch.running -= 1;
    this.settle(record, result);
    options.onSettled?.(record);
    // Its slot and its memory are free: the queue moves on.
    this.pump();
  }

  private settle(record: AgentRecord, result: ChildResult): void {
    if (record.result) return;
    record.result = result;
    record.status = result.status;
    record.finishedAt = this.now();
    record.turns = result.turns;
    record.toolsRunning = 0;
    this.onChanged();
    // The dispatch callback is not available here; settle callers also invoke
    // their onSettled hook immediately after start() returns.
    this.resolvers.get(record)?.(result);
    this.resolvers.delete(record);
    // Waiters first: what they return is delivered through their tool result.
    const still: typeof this.waiters = [];
    for (const waiter of this.waiters) {
      if (waiter.records.every((candidate) => candidate.result)) {
        waiter.records.forEach((candidate) => { candidate.delivered = true; });
        waiter.resolve();
      } else {
        still.push(waiter);
      }
    }
    this.waiters = still;
    const changed = this.changeWaiters;
    this.changeWaiters = [];
    changed.forEach((resolve) => resolve());
    if (record.background && !record.delivered) this.onDeliverable();
  }

  /** Memory: delivered results beyond the newest MAX_KEPT_OUTPUTS keep no output. */
  private releaseOldOutputs(parentSessionId: string): void {
    const list = this.bySession.get(parentSessionId) ?? [];
    let kept = 0;
    for (let index = list.length - 1; index >= 0; index--) {
      const result = list[index].result;
      if (!result || !list[index].delivered || result.output === RELEASED_OUTPUT) continue;
      kept += 1;
      if (kept > MAX_KEPT_OUTPUTS) list[index].result = { ...result, output: RELEASED_OUTPUT };
    }
  }

  list(parentSessionId: string): AgentRecord[] {
    return [...(this.bySession.get(parentSessionId) ?? [])];
  }

  /** Records by id (or all of the session's when `ids` is empty); unknown ids are listed apart. */
  resolve(parentSessionId: string, ids: readonly string[] | undefined): { records: AgentRecord[]; unknown: string[] } {
    const all = this.list(parentSessionId);
    if (!ids || ids.length === 0 || ids.includes("all")) return { records: all, unknown: [] };
    const records: AgentRecord[] = [];
    const unknown: string[] = [];
    for (const raw of ids) {
      const id = String(raw).trim();
      const record = all.find((candidate) => candidate.id === id || candidate.sessionId === id);
      if (record && !records.includes(record)) records.push(record);
      else if (!record) unknown.push(id);
    }
    return { records, unknown };
  }

  /** Children not finished yet (queued or running), of one session or of all. */
  active(parentSessionId?: string): AgentRecord[] {
    const lists = parentSessionId ? [this.list(parentSessionId)] : [...this.bySession.values()];
    return lists.flat().filter((record) => !record.result);
  }

  /** Background results nobody has seen yet, of one session or of all. */
  hasPending(parentSessionId?: string): boolean {
    const lists = parentSessionId ? [this.list(parentSessionId)] : [...this.bySession.values()];
    return lists.flat().some((record) => record.background && !record.delivered);
  }

  /**
   * Finished background results ready to go back to the parent, grouped per
   * delivery (a whole batch, or one child with notify "each"); marked
   * delivered.
   */
  takeDeliverable(parentSessionId?: string): Array<{ parentSessionId: string; records: AgentRecord[] }> {
    const out: Array<{ parentSessionId: string; records: AgentRecord[] }> = [];
    const sessions = parentSessionId ? [parentSessionId] : [...this.bySession.keys()];
    for (const session of sessions) {
      const records = this.list(session).filter((record) => record.background);
      const batches = new Map<string, AgentRecord[]>();
      for (const record of records) batches.set(record.batch, [...(batches.get(record.batch) ?? []), record]);
      for (const members of batches.values()) {
        const ready = members.filter((record) => record.result && !record.delivered);
        if (ready.length === 0) continue;
        if (members[0].notify === "batch" && members.some((record) => !record.result)) continue;
        ready.forEach((record) => { record.delivered = true; });
        out.push({ parentSessionId: session, records: ready });
      }
    }
    return out;
  }

  /**
   * Wait until every one of `records` finished, the wait's own limit passed
   * or `signal` fired (the children keep running then). Finished ones are
   * marked delivered: the caller returns them.
   */
  async wait(records: AgentRecord[], options: { signal?: AbortSignal; timeoutMs?: number } = {}): Promise<void> {
    if (!records.every((record) => record.result)) {
      await new Promise<void>((resolve) => {
        let timer: ReturnType<typeof setTimeout> | null = null;
        const waiter = {
          records,
          resolve: () => {
            if (timer) clearTimeout(timer);
            options.signal?.removeEventListener("abort", stop);
            resolve();
          },
        };
        const stop = () => {
          this.waiters = this.waiters.filter((candidate) => candidate !== waiter);
          waiter.resolve();
        };
        this.waiters.push(waiter);
        if (options.timeoutMs && options.timeoutMs > 0) timer = setTimeout(stop, options.timeoutMs);
        if (options.signal?.aborted) stop();
        else options.signal?.addEventListener("abort", stop, { once: true });
      });
    }
    records.filter((record) => record.result).forEach((record) => { record.delivered = true; });
  }

  /** Resolves at the next child that settles. */
  changed(): Promise<void> {
    return new Promise((resolve) => this.changeWaiters.push(resolve));
  }

  /** Stop children (SIGTERM, then SIGKILL) and wait for them to exit. */
  async cancel(records: AgentRecord[], reason: string = "cancelled", options: { delivered?: boolean } = {}): Promise<void> {
    const live = records.filter((record) => !record.result);
    // The caller hands the results over itself: no delivery message.
    if (options.delivered) live.forEach((record) => { record.delivered = true; });
    for (const record of live) {
      record.controller.abort(reason);
      if (record.startedAt === null) {
        // Still waiting for a slot: settled here, it never starts.
        this.settle(record, buildChildResult(record, {
          exitCode: null,
          finalText: "",
          turns: 0,
          status: reason === INTERRUPTED_REASON ? "interrupted" : "cancelled",
          error: "cancelled before it started",
        }, 0, record.sessionId));
      }
    }
    if (live.some((record) => record.startedAt === null)) this.pump();
    await Promise.all(live.map((record) => record.done));
  }

  /** The session ends: every unfinished child is interrupted; returns them once they exited. */
  async interruptAll(parentSessionId?: string): Promise<AgentRecord[]> {
    const live = this.active(parentSessionId);
    // Nobody is left to read them.
    await this.cancel(live, INTERRUPTED_REASON, { delivered: true });
    return live;
  }
}

/** Minutes, rounded for status lines. */
function minutesOf(ms: number): number {
  return Math.round((ms / 60_000) * 10) / 10;
}

/** Compact status used by both the tool stream and the TUI status bar. */
export function agentStatusSummary(records: AgentRecord[]): string {
  if (records.length === 0) return "No sub-agents in this session.";
  const running = records.filter((record) => record.status === "running").length;
  const queued = records.filter((record) => record.status === "queued").length;
  const forMemory = records.filter((record) => record.status === "queued" && record.waiting?.reason === "memory").length;
  const settled = records.filter((record) => Boolean(record.result)).length;
  const limits = [...new Set(records.map((record) => record.parallelLimit))];
  const parallel = limits.length === 1 && limits[0] === null
    ? "unlimited"
    : limits.length === 1
      ? String(limits[0])
      : "per-batch";
  const queuedText = forMemory > 0 ? `${queued} queued: waiting for memory` : `${queued} queued`;
  return `${records.length} sub-agents: ${running} running, ${queuedText}, ${settled} settled (parallel: ${parallel})`;
}

/** One status line per sub-agent (agents_status, /agents). */
export function renderAgentStatus(records: AgentRecord[], now: number = Date.now()): string {
  if (records.length === 0) return "No sub-agents in this session.";
  const lines = records.map((record) => {
    const since = record.startedAt ?? record.queuedAt;
    const elapsed = minutesOf((record.finishedAt ?? now) - since);
    const state = record.status === "queued" && record.waiting ? "queued: waiting for memory" : record.status;
    const bits = [`${record.id}`, `[${state}]`, record.role];
    if (record.model) bits.push(`[${record.model}${record.effort ? ` · ${record.effort}` : ""}]`);
    bits.push(record.background ? "background" : "blocking");
    const detail: string[] = [`${elapsed} min`];
    if (!record.result && record.startedAt) {
      detail.push(`${record.turns} turns`);
      if (record.toolsRunning > 0) detail.push("tool running");
      detail.push(`last activity ${minutesOf(now - record.lastActivityAt)} min ago`);
    }
    if (record.result && !record.delivered) detail.push("result not yet read");
    if (!record.result && record.waiting) detail.push(record.waiting.detail);
    const task = record.task.length > 80 ? `${record.task.slice(0, 79)}…` : record.task;
    return `- ${bits.join(" ")} (${detail.join(", ")}) — ${task}`;
  });
  return `${agentStatusSummary(records)}:\n${lines.join("\n")}`;
}

/** The message that brings finished background agents back to the parent. */
export function renderDelivery(records: AgentRecord[]): string {
  const results = records.map((record) => record.result).filter((result): result is ChildResult => Boolean(result));
  const ids = records.map((record) => record.id).join(", ");
  return [
    `Background sub-agents finished (${ids}). Their results follow; continue your work with them.`,
    "",
    renderChildResults(results, records.map((record) => record.id)),
  ].join("\n");
}
