import { randomBytes } from "node:crypto"; import { execaCommandSync, parseCommandString } from "execa"; import { fromWritable } from "from-node-stream"; import { mkdir, open, readFile, stat, unlink, writeFile } from "fs/promises"; import path from "path"; import DIE from "phpdie"; import sflow from "sflow"; import { XtermProxy } from "./xterm-proxy.ts"; import { agentYesHome } from "./agentYesHome.ts"; import { extractSessionId, getSessionForCwd, isCodexFamily, storeSessionForCwd, } from "./resume/codexSessionManager.ts"; import pty, { ptyPackage } from "./pty.ts"; import { removeControlCharacters } from "./removeControlCharacters.ts"; import { stripClaudeSessionPin } from "./sessionEnv.ts"; import { acquireLock, releaseLock, shouldUseLock } from "./runningLock.ts"; import { logger } from "./logger.ts"; import { createFifoStream } from "./beta/fifo.ts"; import { PidStore } from "./pidStore.ts"; import { TitlePublisher, TitleScanner } from "./titleScanner.ts"; import { sendEnter, sendMessage } from "./core/messaging.ts"; import { AUTO_RETRY_GIVE_UP_MS, AUTO_RETRY_MIN_IDLE_MS, AUTO_RETRY_REASON_FALLBACK, autoRetryBackoffMs, buildAutoRetryMessage, classifyAutoRetryReason, formatDurationSecs, shouldFireRetry, } from "./autoRetry.ts"; import { recordInbox } from "./messageLog.ts"; import { initializeLogPaths, setupDebugLogging, saveLogFile, saveDeprecatedLogFile, } from "./core/logging.ts"; import { spawnAgent } from "./core/spawner.ts"; import { resolveProgram } from "./resolveBinary.ts"; import { AgentContext } from "./core/context.ts"; import { createTerminatorStream } from "./core/streamHelpers.ts"; import { globalAgentRegistry } from "./agentRegistry.ts"; import { notifyWebhook } from "./webhookNotifier.ts"; import { readGlobalPids } from "./globalPidIndex.ts"; import * as reaper from "./reaper.ts"; import { derivePermissions } from "./agentPermissions.ts"; export { removeControlCharacters }; export { AgentContext }; // `ay todo` library surface — only symbols listed here reach the published // `dist/index.js` a consuming project (e.g. a thin project-specific shell // around this engine) actually imports; creating the underlying modules // alone is not sufficient. export { openStore, TodoStore, CycleError } from "./todoStore.ts"; export type { TodoRecord, CreateInput, ListFilter, GateRegistration, GateEvidence, } from "./todoStore.ts"; export { LIFECYCLES, DONE_STATE, ORPHANED_STATE, nextStates, canTransition, requiredGate, initialState, statesOf, isKnownKind, } from "./todoLifecycle.ts"; export type { LifecycleKind, LifecycleGraph, LifecycleTransition } from "./todoLifecycle.ts"; export { monitorHint, describeBlock } from "./todoBlock.ts"; export type { TodoBlock, MonitorHint } from "./todoBlock.ts"; export { unblockedTasks, openBlockers, renderTree, renderDigest } from "./todoDigest.ts"; export { reconcileTodos } from "./todoAutomation.ts"; export type { TodoAction, LiveAgent } from "./todoAutomation.ts"; export { listAsks, listAsksForProject, answerAsk, hasTodoStore } from "./askApi.ts"; export type { AskItem, AnswerInput, AnswerResult } from "./askApi.ts"; // Raw PTY logs are an append-only capture; the reader (`readLogForRender`) only // ever replays the trailing ~64 MiB, so older bytes are dead weight on disk (a // runaway capture was seen at ~1 GB). Once the file passes RAW_LOG_TRIGGER_BYTES, // compact it to its last RAW_LOG_KEEP_BYTES (kept above the render window so no // visible output is lost). Mirrors the Rust runtime's LogWriter cap in // rs/src/log_files.rs — keep the two in sync. const RAW_LOG_KEEP_BYTES = 80 * 1024 * 1024; const RAW_LOG_TRIGGER_BYTES = 160 * 1024 * 1024; /** Shrink a raw log to its trailing `keepBytes` **in place** (rewrite tail→front * + truncate), returning the new length. Deliberately NOT temp-file+rename: * rename swaps the inode, freezing any live follower (`ay serve` `/api/tail`, * incl. the agent-yes.com viewer over WebRTC) whose fd is bound to the old inode. * Rewriting the same inode lets serve's `if (size < offset) offset = size` * "truncated/rotated" guard resume the stream. Mirrors rs/src/log_files.rs. */ async function compactRawLogTail(logPath: string, keepBytes = RAW_LOG_KEEP_BYTES): Promise { const fh = await open(logPath, "r+"); try { const { size } = await fh.stat(); if (size <= keepBytes) return size; // Read the tail fully before overwriting the front so it can't be clobbered. const buf = Buffer.allocUnsafe(keepBytes); const { bytesRead } = await fh.read(buf, 0, keepBytes, size - keepBytes); await fh.write(buf, 0, bytesRead, 0); await fh.truncate(bytesRead); return bytesRead; } finally { await fh.close(); } } export type AgentCliConfig = { // cli install?: | string | { powershell?: string; bash?: string; npm?: string; unix?: string; windows?: string }; // hint user for install command if not installed version?: string; // hint user for version command to check if installed binary?: string; // actual binary name if different from cli, e.g. cursor -> cursor-agent // Env vars injected into the spawned agent (e.g. glm points `claude` at Z.AI). // Values support ${VAR} expansion against the launching env; an entry whose // variable is unset is skipped so it can't blank out an inherited value. env?: Record; defaultArgs?: string[]; // function to ensure certain args are present yesArgs?: string[]; // appended when `-y`/--yes is passed: the per-CLI "yolo" flag (claude: --dangerously-skip-permissions; codex: --dangerously-bypass-approvals-and-sandbox) help?: string; // documentation/help URL for the CLI bunx?: boolean; // metadata for bunx-based launches systemPrompt?: string; // flag name for system prompt injection system?: string; // system prompt content to inject // status detect, and actions ready?: RegExp[]; // regex matcher for stdin ready. Set to empty array [] to disable ready check entirely. fatal?: RegExp[]; // array of regex to match for fatal errors working?: RegExp[]; // regex matcher for working status updateAvailable?: RegExp[]; // regex matcher for update available banners exitCommands?: string[]; // commands to exit the cli gracefully promptArg?: (string & {}) | "first-arg" | "last-arg" | "typed"; // argument name to pass the prompt, e.g. --prompt, first-arg for positional arg, or "typed" to type it into an interactive session after ready (shells) // handle special format /** * The CLI turns bracketed-paste mode on (DECSET 2004), so `ay send` can * deliver a body as ONE paste instead of a burst the CLI has to segment by * timing — see ts/bracketedPaste.ts for the failure this prevents. Opt-in * because a CLI without the mode would show the markers as literal text. * Verify before setting it: the agent's `.raw.log` contains `ESC[?2004h`. */ bracketedPaste?: boolean; noEOL?: boolean; // if true, do not split lines by \n when handling inputs, e.g. for codex, which uses cursor-move csi code instead of \n to move lines // auto responds enter?: RegExp[]; // array of regex to match for sending Enter enterExclude?: RegExp[]; // array of regex to exclude from auto-enter (even if enter matches) typingRespond?: { [message: string]: RegExp[] }; // type specified message to a specified pattern autoRetry?: RegExp[]; // recoverable API errors (overload/rate-limit/usage-limit): type "retry" with exponential backoff (up to 8h) instead of exiting needsInput?: RegExp[]; // agent is blocked on an interactive selection menu it did NOT auto-resolve; surfaced as the `needs_input` state by `ay ls`/`ay status` (query-time only — not consumed by the run loop) // crash/resuming-session behaviour restoreArgs?: string[]; // arguments to continue the session when crashed restartWithoutContinueArg?: RegExp[]; // array of regex to match for errors that require restart without continue args resumeCommand?: RegExp; // optional: scraped from the agent's log on `ay restart`. Capture group 1 is the argument string (whitespace-split) to relaunch with — i.e. a resume command the CLI printed, like "--resume ". When absent or unmatched, restart falls back to restoreArgs (e.g. --continue). }; export type AgentYesConfig = { configDir?: string; // directory to store agent-yes config files, e.g. session store logsDir?: string; // directory to store agent-yes log files clis: { [key: string]: AgentCliConfig }; }; // load user config from agent-yes.config.ts if exists export const config = await import("../agent-yes.config.ts").then((mod) => mod.default || mod); export const CLIS_CONFIG = config.clis as Record< keyof Awaited["clis"], AgentCliConfig >; /** * Main function to run agent-cli with automatic yes/no responses * @param options Configuration options * @param options.continueOnCrash - If true, automatically restart agent-cli when it crashes: * 1. Shows message 'agent-cli crashed, restarting..' * 2. Spawns a new 'agent-cli --continue' process * 3. Re-attaches the new process to the shell stdio (pipes new process stdin/stdout) * 4. If it crashes with "No conversation found to continue", exits the process * @param options.exitOnIdle - Exit when agent-cli is idle. Boolean or timeout in milliseconds, recommended 5000 - 60000, default is false * @param options.cliArgs - Additional arguments to pass to the agent-cli CLI * @param options.removeControlCharactersFromStdout - Remove ANSI control characters from stdout. Defaults to !process.stdout.isTTY * @param options.disableLock - Disable the running lock feature that prevents concurrent agents in the same directory/repo * * @example * ```typescript * import agentYes from 'agent-yes'; * await agentYes({ * prompt: 'help me solve all todos in my codebase', * * // optional * cliArgs: ['--verbose'], // additional args to pass to agent-cli * exitOnIdle: 30000, // exit after 30 seconds of idle * robust: true, // auto restart with --continue if claude crashes, default is true * logFile: 'claude-output.log', // save logs to file * disableLock: false, // disable running lock (default is false) * }); * ``` */ export default async function agentYes({ cli, cliArgs = [], skipPermissions = false, prompt, robust = true, cwd, env, exitOnIdle, logFile, removeControlCharactersFromStdout = false, // = !process.stdout.isTTY, verbose = false, queue = false, install = false, resume = false, useSkills = false, useStdinAppend = false, autoYes = true, idleAction, swarmHint = true, }: { cli: keyof typeof CLIS_CONFIG; cliArgs?: string[]; skipPermissions?: boolean; // if true (`-y`/--yes), append the per-CLI yesArgs ("yolo" flag) prompt?: string; robust?: boolean; cwd?: string; env?: Record; exitOnIdle?: number; logFile?: string; removeControlCharactersFromStdout?: boolean; verbose?: boolean; queue?: boolean; install?: boolean; // if true, install the cli tool if not installed, e.g. will run `npm install -g cursor-agent` resume?: boolean; // if true, resume previous session in current cwd if any useSkills?: boolean; // if true, prepend SKILL.md header to the prompt for non-Claude agents useStdinAppend?: boolean; // if true, enable FIFO input stream on Linux, for additional stdin input autoYes?: boolean; // if true, auto-yes is enabled (default), toggle with Ctrl+Y during session idleAction?: string; // if set, type this message when idle instead of exiting swarmHint?: boolean; // if true (default), inject peer discovery hint when other agents are running; --no-swarm-hint to opt out }) { if (!cli) throw new Error(`cli is required`); const conf = CLIS_CONFIG[cli] || DIE(`Unsupported cli tool: ${cli}, current process.argv: ${process.argv.join(" ")}`); // Acquire lock before starting agent (if in git repo or same cwd and lock is not disabled) const workingDir = cwd ?? process.cwd(); if (queue) { if (queue && shouldUseLock(workingDir)) { await acquireLock(workingDir, prompt ?? "Interactive session"); } // Register cleanup handlers for lock release const cleanupLock = async () => { if (queue && shouldUseLock(workingDir)) { await releaseLock().catch(() => null); // Ignore errors during cleanup } }; process.on("exit", () => { if (queue) releaseLock().catch(() => null); }); process.on("SIGINT", async (code) => { await cleanupLock(); process.exit(code); }); process.on("SIGTERM", async (code) => { await cleanupLock(); process.exit(code); }); } // Initialize process registry const pidStore = new PidStore(workingDir); await pidStore.init(); // Track when user sends Ctrl+C to avoid treating intentional exit as crash let userSentCtrlC = false; if (verbose) logger.debug( `[stdin] isTTY: ${process.stdin.isTTY}, setRawMode available: ${!!process.stdin.setRawMode}`, ); process.stdin.setRawMode?.(true); // must be called any stdout/stdin usage if (verbose) logger.debug(`[stdin] Raw mode set, isRaw: ${(process.stdin as any).isRaw}`); // XtermProxy: headless xterm emulator that auto-responds to all terminal // queries (DSR, DA, etc.) so the spawned CLI never blocks waiting for replies. // writeToPty is set after shell spawn (see below). let shellWrite: (data: string) => void = () => {}; const xtermProxy = new XtermProxy({ ...getTerminalDimensions(), writeToPty: (data) => shellWrite(data), }); logger.debug(`Using ${ptyPackage} for pseudo terminal management.`); // Detect if running as sub-agent const isSubAgent = !!process.env.CLAUDE_PPID; if (isSubAgent) logger.info(`[${cli}-yes] Running as sub-agent (CLAUDE_PPID=${process.env.CLAUDE_PPID})`); // Apply CLI specific configurations (moved to CLI_CONFIGURES) const cliConf = (CLIS_CONFIG as Record)[cli] || {}; cliArgs = cliConf.defaultArgs ? [...cliConf.defaultArgs, ...cliArgs] : cliArgs; // `-y`/--yes appends the per-CLI "yolo" args. Each CLI declares its own: // claude → --dangerously-skip-permissions; codex → // --dangerously-bypass-approvals-and-sandbox (codex rejects the claude flag, // and its bwrap sandbox can't init inside an already-sandboxed container). if (skipPermissions && cliConf.yesArgs?.length) { cliArgs = [...cliArgs, ...cliConf.yesArgs]; } // The AGENT_YES_PID we INHERITED is the wrapper pid of the agent that spawned // us — see the ptyEnv block below, which reads the same value before re-stamping // it with our own pid. Resolved here (before any prompt decoration) because the // `` wrapper has to go around the RAW task, with the SKILL.md // header and the peer hint prepended outside it: those are ambient guidance, not // part of what the parent asked for. const inheritedParentPid = ((): number | undefined => { const n = Number((env ?? (process.env as Record)).AGENT_YES_PID); return Number.isInteger(n) && n > 0 ? n : undefined; })(); // The task as the caller wrote it, before any wrapping/prefixing and before // promptArg clears it — used for the registry (so `ay ls` shows the agent's // role, not wrapper boilerplate) and to remind the parent what it asked for // when a report fires. const originalPrompt = prompt; // Wrap a sub-agent's initial prompt in `` — same attribution + // reply route `ay send` puts on every later message, plus the reporting duty. // No-ops for a top-level (human-launched) agent and for an interactive session // with no prompt. See ts/initMsg.ts (mirrored in rs/src/init_msg.rs). if (prompt && inheritedParentPid) { try { const { resolveSpawner } = await import("./parentLink.ts"); const spawner = await resolveSpawner(inheritedParentPid); if (spawner) { const { buildInitMsg } = await import("./initMsg.ts"); const { formatIdentity } = await import("./identity.ts"); // The spawner runs on THIS host (we resolved it from a local wrapper // pid), so the local user/host and its cwd's branch are the right // defaults — the same call rs/src/main.rs makes. const identity = formatIdentity({ cwd: spawner.cwd, pid: spawner.pid }); prompt = buildInitMsg(prompt, spawner, randomBytes(4).toString("hex"), identity); } } catch (error) { // Attribution is never worth failing a spawn over. if (verbose) logger.warn("[init-msg] failed to wrap initial prompt:", { error }); } } // If enabled, read SKILL.md header and prepend to the prompt for non-Claude agents try { const workingDir = cwd ?? process.cwd(); if (useSkills && cli !== "claude") { // Find git root to determine search boundary let gitRoot: string | null = null; try { const result = execaCommandSync("git rev-parse --show-toplevel", { cwd: workingDir, reject: false, }); if (result.exitCode === 0) { gitRoot = result.stdout.trim(); } } catch { // Not a git repo, will only check cwd } // Walk up from cwd to git root (or stop at filesystem root) collecting SKILL.md files const skillHeaders: string[] = []; let currentDir = workingDir; const searchLimit = gitRoot || path.parse(currentDir).root; while (true) { const skillPath = path.resolve(currentDir, "SKILL.md"); const md = await readFile(skillPath, "utf8").catch(() => null); if (md) { // Extract header (content before first level-2 heading `## `) const headerMatch = md.match(/^[\s\S]*?(?=\n##\s)/); const headerRaw = (headerMatch ? headerMatch[0] : md).trim(); if (headerRaw) { skillHeaders.push(headerRaw); if (verbose) logger.info(`[skills] Found SKILL.md in ${currentDir} (${headerRaw.length} chars)`); } } // Stop if we've reached git root or filesystem root if (currentDir === searchLimit) break; const parentDir = path.dirname(currentDir); if (parentDir === currentDir) break; // Reached filesystem root currentDir = parentDir; } if (skillHeaders.length > 0) { // Combine all headers (most specific first) const combined = skillHeaders.join("\n\n---\n\n"); const MAX = 2000; // increased limit for multiple skills const header = combined.length > MAX ? combined.slice(0, MAX) + "…" : combined; const prefix = `Use this repository skill as context:\n\n${header}`; prompt = prompt ? `${prefix}\n\n${prompt}` : prefix; if (verbose) logger.info( `[skills] Injected ${skillHeaders.length} SKILL.md header(s) (${header.length} chars total)`, ); } else { if (verbose) logger.info("[skills] No SKILL.md found in directory hierarchy"); } } } catch (error) { // Non-fatal; continue without skills if (verbose) logger.warn("[skills] Failed to inject SKILL.md header:", { error }); } // Inject peer discovery hint when other agents are running if (swarmHint) { try { const peers = await readGlobalPids({ liveOnly: true }); if (peers.length > 0) { const hint = `${peers.length} peer agent${peers.length > 1 ? "s are" : " is"} running. Before asking the user for input on any domain-specific topic (design, testing, architecture, etc.), check for relevant peers first: \`ay ls --json\` (see \`prompt\` field for their role). Ask one: \`ay send \`. Read reply: \`ay tail \`. Do not use interactive forms or user prompts when a peer can answer.`; if (cli === "claude") { cliArgs = ["--append-system-prompt", hint, ...cliArgs]; } // Prepend to prompt for all CLIs (including claude) so it's read before the task prompt = prompt ? `[${hint}]\n\n${prompt}` : hint; } } catch { // Non-fatal } } // Handle --continue flag for codex session restoration if (resume) { if (isCodexFamily(cli) && resume) { // Try to get stored session for this directory const storedSessionId = await getSessionForCwd(workingDir); if (storedSessionId) { // Replace or add resume args cliArgs = ["resume", storedSessionId, ...cliArgs]; await logger.debug(`resume|using stored session ID: ${storedSessionId}`); } else { throw new Error( `No stored session found for codex in directory: ${workingDir}, please try without resume option.`, ); } } else if (cli === "claude") { // just add --continue flag for claude cliArgs = ["--continue", ...cliArgs]; await logger.debug(`resume|adding --continue flag for claude`); } else { throw new Error( `Resume option is not supported for cli: ${cli}, make a feature request if you want it. https://github.com/snomiao/agent-yes/issues`, ); } } // If possible pass prompt via cli args, its usually faster than stdin if (prompt && cliConf.promptArg) { if (cliConf.promptArg === "first-arg") { cliArgs = [prompt, ...cliArgs]; prompt = undefined; // clear prompt to avoid sending later } else if (cliConf.promptArg === "last-arg") { cliArgs = [...cliArgs, prompt]; prompt = undefined; // clear prompt to avoid sending later } else if (cliConf.promptArg.startsWith("--")) { cliArgs = [cliConf.promptArg, prompt, ...cliArgs]; prompt = undefined; // clear prompt to avoid sending later } else if (cliConf.promptArg === "typed") { // Shell mode (bash/cmd/powershell): don't pass an argv (that would // run-and-exit). Leave `prompt` set so it's typed into the interactive // session at onStart, keeping the shell alive afterwards. } else { logger.warn(`Unknown promptArg format: ${cliConf.promptArg}`); } } // Opportunistic sweep: reap any process group leaked by an agent whose wrapper // died without cleanup, before we start a new one. See ts/reaper.ts. reaper.sweep().catch(() => {}); // Spawn the agent CLI process const ptyEnv = { ...(env ?? (process.env as Record)) }; // The AGENT_YES_PID we INHERITED (the wrapper of the parent agent that launched // this nested `ay`), read above from the same env before we overwrite it with // our own pid here. undefined when started from a human shell → tree root. const parentPid = inheritedParentPid; ptyEnv.AGENT_YES_PID = String(process.pid); // A caller-injected AGENT_YES_AGENT_ID (from `ay serve`'s /api/spawn) is meant // for THIS agent's record only — pidStore.registerProcess reads it from our own // env. Strip it from the wrapped CLI's env so the agent's subagents (a nested // `ay`) don't inherit it and register under the same id (which would make that // id ambiguous). The wrapper's own process.env still carries it for pidStore. delete ptyEnv.AGENT_YES_AGENT_ID; // Strip the parent Claude Code session markers so the wrapped CLI is a CLEAN // top-level session. Without this, an `ay claude` launched from inside another // Claude Code session (or this claude's Bash tool) inherits // CLAUDE_CODE_CHILD_SESSION — the child claude then disables transcript saving // ("⚠ Transcript saving is off …") — and CLAUDE_CODE_SSE_PORT/SESSION_ID make it // attach to the parent's stale session. The daemon /api/spawn path already drops // these via freshAgentEnv; this covers the direct-launch path. Idempotent, and // AGENT_YES_PID is handled above (read for parentPid, then re-stamped). See // ts/sessionEnv.ts (mirrored in rs/src/pty_spawner.rs for the Rust runtime). stripClaudeSessionPin(ptyEnv); // Inject per-CLI env (e.g. glm → Z.AI endpoint). Expand ${VAR} against the // launching env; skip entries whose vars are unset so we never blank out an // inherited value (e.g. ANTHROPIC_AUTH_TOKEN when ZAI_API_KEY isn't exported). // `${VAR:-default}` falls back to `default` when VAR is unset/empty — used for // overridable defaults like the model (glm → z-ai/glm-5.2). if (cliConf?.env) { for (const [key, raw] of Object.entries(cliConf.env)) { let unresolved = false; const value = raw.replace( /\$\{([A-Za-z_][A-Za-z0-9_]*)(?::-([^}]*))?\}/g, (_, name, fallback) => { const v = ptyEnv[name]; if (v !== undefined && v !== "") return v; if (fallback !== undefined) return fallback; unresolved = true; return ""; }, ); if (unresolved) continue; ptyEnv[key] = value; } } // The agent runs in a PTY (a real terminal), so advertise terminal // capabilities. A console/daemon-spawned agent inherits an env with no TERM/ // COLORTERM: neither the daemon (no controlling terminal) nor the recovered // login-shell env (captured without a tty) carries them — those vars are set // by the terminal emulator, not by the shell. Without them the wrapped CLI // renders colorless in the web console. Fill only when absent so a // terminal-launched agent keeps its real values (e.g. xterm-256color, tmux). if (!ptyEnv.TERM) ptyEnv.TERM = "xterm-256color"; if (!ptyEnv.COLORTERM) ptyEnv.COLORTERM = "truecolor"; const ptyOptions = { name: ptyEnv.TERM, ...getTerminalDimensions(), cwd: cwd ?? process.cwd(), env: ptyEnv, }; let shell = spawnAgent({ cli, cliConf, cliArgs, verbose, install, ptyOptions, }); // Wire up the xterm proxy to write back to the PTY shellWrite = (data: string) => shell.write(data); // Capture the child CLI's terminal title (OSC 0/2) into the registry so // `ay whoami` / `ay ls --json` can answer "what is this agent doing". // TitlePublisher change-gates + rate-limits the registry writes. const titleScanner = new TitleScanner(); const titlePublisher = new TitlePublisher((title) => { pidStore.updateTitle(shell.pid, title).catch(() => null); }); // Attach data handler IMMEDIATELY after spawn to avoid losing early PTY output. // node-pty emits 'data' events eagerly — if no listener is attached, events are lost. function onData(data: string) { const currentPid = shell.pid; xtermProxy.write(data); globalAgentRegistry.appendStdout(currentPid, data); const title = titleScanner.feed(data); if (title) titlePublisher.observe(title); else titlePublisher.poll(); } shell.onData(onData); // Register process in pidStore (non-blocking - failures should not prevent agent from running) try { await pidStore.registerProcess({ pid: shell.pid, cli, args: cliArgs, // The task as the user/parent wrote it — NOT the delivered `prompt`, which // by now carries the `` wrapper (and any SKILL.md header / // peer hint). `ay ls` shows this field as the agent's role; a wall of // wrapper boilerplate there would bury the one line that identifies it. prompt: originalPrompt, cwd: workingDir, // We inject our own pid as AGENT_YES_PID into the agent's env above; record // it so a child `ay send` can map that env value back to this agent. wrapperPid: process.pid, // The parent agent's wrapper pid (inherited AGENT_YES_PID), for the tree. parentPid, // Audit trail: what this agent is actually allowed to do. `auto_continue` // is robust AND a CLI that declares restoreArgs — robust alone only // restarts, it resumes nothing. permissions: derivePermissions({ cliArgs, yesArgs: cliConf.yesArgs, robust, autoContinue: Boolean(robust && cliConf.restoreArgs?.length), }), }); } catch (error) { logger.warn(`[pidStore] Failed to register process ${shell.pid}:`, error); } // Defense-in-depth: record (this wrapper, the agent's process group) so a later // sweep reaps the group if we're killed without running onExit cleanup. The PTY // child is a session leader, so its pgid == shell.pid. See ts/reaper.ts. reaper.register(process.pid, shell.pid).catch(() => {}); notifyWebhook("RUNNING", prompt ?? "", workingDir).catch(() => null); // Initialize log paths (independent of registration) const logPaths = await initializeLogPaths(pidStore, shell.pid); await setupDebugLogging(logPaths.debuggingLogsPath); // Track whether the session ever switched to the alternate screen buffer. // render() reconstructs scrollback of the normal buffer only; the alt buffer // keeps no scrollback, so for alt-screen apps the rendered log would be just // the final frame and must NOT replace the raw byte log. Current CLIs (ink- // based claude/codex) stay on the normal buffer, but guard for the future. let usedAltScreen = false; // Create agent context const ctx = new AgentContext({ shell, pidStore, logPaths, cli, cliConf, verbose, robust, autoYes, }); // Register agent in global registry (non-blocking) try { globalAgentRegistry.register(shell.pid, { pid: shell.pid, context: ctx, cwd: workingDir, cli, prompt: originalPrompt, startTime: Date.now(), stdoutBuffer: [], }); } catch (error) { logger.warn(`[agentRegistry] Failed to register agent ${shell.pid}:`, error); } // Show startup mode if not default (i.e., when starting in manual mode) if (!ctx.autoYesEnabled) { process.stderr.write("\x1b[33m[auto-yes: OFF]\x1b[0m Press Ctrl+Y to toggle\n"); } // If ready check is disabled (empty array) or manual mode, mark stdin ready immediately // Manual mode needs immediate stdin so user can respond to trust prompts if ((cliConf.ready && cliConf.ready.length === 0) || !ctx.autoYesEnabled) { ctx.stdinReady.ready(); ctx.stdinFirstReady.ready(); } // force ready after 10s to avoid stuck forever if the ready-word mismatched sleep(10e3).then(() => { if (!ctx.stdinReady.isReady) ctx.stdinReady.ready(); if (!ctx.stdinFirstReady.isReady) ctx.stdinFirstReady.ready(); }); const pendingExitCode = Promise.withResolvers(); shell.onExit(async function onExit({ exitCode }) { const exitedPid = shell.pid; // Capture PID immediately before any shell reassignment // Reap the exited agent's process group. The PTY child is a session/group // leader, so a `yes | cmd` (or any descendant) it leaked shares its pgid even // after it reparents to PID 1 — kill the group so orphans don't spin at ~100% // CPU forever. Targeting the pgid (not ppid==1) is container-safe and never // touches processes outside this agent's session. Runs on the final exit AND // before each robust restart below. if (process.platform !== "win32") { try { process.kill(-exitedPid, "SIGKILL"); } catch { // ESRCH = no surviving group members left to reap; nothing to do. } } // Unregister from agent registry globalAgentRegistry.unregister(exitedPid); ctx.stdinReady.unready(); // start buffer stdin // Exit codes 130 (SIGINT/Ctrl+C) and 143 (SIGTERM) are intentional exits, not crashes // Also check if user sent Ctrl+C recently (within last 2 seconds) const intentionalExit = exitCode === 130 || exitCode === 143 || userSentCtrlC; const agentCrashed = exitCode !== 0 && !intentionalExit; // Handle restart without continue args (e.g., "No conversation found to continue") // logger.debug(``, { shouldRestartWithoutContinue, robust }) if (ctx.shouldRestartWithoutContinue) { // Update status (non-blocking) try { await pidStore.updateStatus(exitedPid, "exited", { exitReason: "restarted", exitCode: exitCode ?? undefined, }); } catch (error) { logger.warn(`[pidStore] Failed to update status for PID ${exitedPid}:`, error); } ctx.shouldRestartWithoutContinue = false; // reset flag ctx.isFatal = false; // reset fatal flag to allow restart // Enforce restart limit with exponential backoff if (ctx.restartCount >= 10) { logger.error(`${cli} reached max restarts (10), giving up.`); return pendingExitCode.resolve(exitCode); } const backoffMs = 1000 * Math.pow(2, ctx.restartCount); logger.info(`Restart ${ctx.restartCount + 1}/10, waiting ${backoffMs}ms before restart...`); await sleep(backoffMs); ctx.restartCount++; // Restart without continue args - use original cliArgs without restoreArgs const cliCommand = cliConf?.binary || cli; let [bin, ...args] = [ ...parseCommandString(cliCommand), ...cliArgs.filter((arg) => !["--continue", "--resume"].includes(arg)), ]; logger.info(`Restarting ${cli} ${JSON.stringify([bin, ...args])}`); const restartPtyOptions = { name: ptyEnv.TERM, ...getTerminalDimensions(), cwd: cwd ?? process.cwd(), env: ptyEnv, }; // Resolve via PATH so a cwd entry named like the CLI can't shadow it (#138). shell = pty.spawn(resolveProgram(bin!, ptyEnv.PATH), args, restartPtyOptions); shellWrite = (data: string) => shell.write(data); // Register process in pidStore (non-blocking) try { await pidStore.registerProcess({ pid: shell.pid, cli, args, prompt: originalPrompt, cwd: workingDir, }); } catch (error) { logger.warn(`[pidStore] Failed to register restarted process ${shell.pid}:`, error); } // Re-register the NEW process group with the reaper — the restart gave us a // fresh pgid; without this the reaper would track the old (now-dead) group // and the live one would leak if we're SIGKILLed. Mirrors the Rust loop. reaper.register(process.pid, shell.pid).catch(() => {}); // Update context with new shell ctx.shell = shell; // Register new agent in registry (non-blocking) try { globalAgentRegistry.register(shell.pid, { pid: shell.pid, context: ctx, cwd: workingDir, cli, prompt: originalPrompt, startTime: Date.now(), stdoutBuffer: [], }); } catch (error) { logger.warn(`[agentRegistry] Failed to register restarted agent ${shell.pid}:`, error); } shell.onData(onData); shell.onExit(onExit); // Re-mark stdin ready for manual mode after restart if ((cliConf.ready && cliConf.ready.length === 0) || !ctx.autoYesEnabled) { ctx.stdinReady.ready(); ctx.stdinFirstReady.ready(); } return; } if (agentCrashed && robust && conf?.restoreArgs) { if (!conf.restoreArgs) { logger.warn( `robust is only supported for ${Object.entries(CLIS_CONFIG) .filter(([_, v]) => v.restoreArgs) .map(([k]) => k) .join(", ")} currently, not ${cli}`, ); return; } if (ctx.isFatal) { // Update status (non-blocking) try { await pidStore.updateStatus(exitedPid, "exited", { exitReason: "fatal", exitCode: exitCode ?? undefined, }); } catch (error) { logger.warn(`[pidStore] Failed to update status for PID ${exitedPid}:`, error); } notifyWebhook("EXIT", `fatal exitCode=${exitCode ?? "?"}`, workingDir).catch(() => null); return pendingExitCode.resolve(exitCode); } // Enforce restart limit with exponential backoff if (ctx.restartCount >= 10) { logger.error(`${cli} reached max restarts (10), giving up.`); notifyWebhook("EXIT", `max-restarts exitCode=${exitCode ?? "?"}`, workingDir).catch( () => null, ); return pendingExitCode.resolve(exitCode); } const backoffMs = 1000 * Math.pow(2, ctx.restartCount); logger.info( `${cli} crashed (exit code: ${exitCode}), restart ${ctx.restartCount + 1}/10 in ${backoffMs}ms...`, ); await sleep(backoffMs); ctx.restartCount++; // Update status (non-blocking) try { await pidStore.updateStatus(exitedPid, "exited", { exitReason: "restarted", exitCode: exitCode ?? undefined, }); } catch (error) { logger.warn(`[pidStore] Failed to update status for PID ${exitedPid}:`, error); } // For codex, try to use stored session ID for this directory let restoreArgs = conf.restoreArgs; if (isCodexFamily(cli)) { const storedSessionId = await getSessionForCwd(workingDir); if (storedSessionId) { // Use specific session ID instead of --last restoreArgs = ["resume", storedSessionId]; logger.debug(`restore|using stored session ID: ${storedSessionId}`); } else { logger.debug(`restore|no stored session, using default restore args`); } } const restorePtyOptions = { name: ptyEnv.TERM, ...getTerminalDimensions(), cwd: cwd ?? process.cwd(), env: ptyEnv, }; // Resolve via PATH so a cwd entry named like the CLI can't shadow it (#138). shell = pty.spawn(resolveProgram(cli, ptyEnv.PATH), restoreArgs, restorePtyOptions); shellWrite = (data: string) => shell.write(data); // Register process in pidStore (non-blocking) try { await pidStore.registerProcess({ pid: shell.pid, cli, args: restoreArgs, prompt: originalPrompt, cwd: workingDir, }); } catch (error) { logger.warn(`[pidStore] Failed to register restored process ${shell.pid}:`, error); } // Re-register the NEW process group with the reaper (fresh pgid after the // restore) so the reaper tracks the live group, not the dead one. See above. reaper.register(process.pid, shell.pid).catch(() => {}); // Update context with new shell ctx.shell = shell; // Register new agent in registry (non-blocking) try { globalAgentRegistry.register(shell.pid, { pid: shell.pid, context: ctx, cwd: workingDir, cli, prompt: originalPrompt, startTime: Date.now(), stdoutBuffer: [], }); } catch (error) { logger.warn(`[agentRegistry] Failed to register restored agent ${shell.pid}:`, error); } shell.onData(onData); shell.onExit(onExit); // Re-mark stdin ready for manual mode after restart if ((cliConf.ready && cliConf.ready.length === 0) || !ctx.autoYesEnabled) { ctx.stdinReady.ready(); ctx.stdinFirstReady.ready(); } return; } const exitReason = agentCrashed ? "crash" : "normal"; // Update status (non-blocking) try { await pidStore.updateStatus(exitedPid, "exited", { exitReason, exitCode: exitCode ?? undefined, }); } catch (error) { logger.warn(`[pidStore] Failed to update status for PID ${exitedPid}:`, error); } notifyWebhook("EXIT", `${exitReason} exitCode=${exitCode ?? "?"}`, workingDir).catch( () => null, ); return pendingExitCode.resolve(exitCode); }); // Record the agent's current PTY size to ~/.agent-yes/ptysize/ so `ay serve` // / the web console can render the existing buffer at the agent's real width // before adapting. Mirrors the Rust runtime (rs/src/pty_spawner.rs). const writeCurrentPtysize = (cols: number, rows: number) => { const dir = path.join(agentYesHome(), "ptysize"); void mkdir(dir, { recursive: true }) .then(() => writeFile(path.join(dir, String(process.pid)), `${cols} ${rows}\n`)) .catch(() => null); }; { const { cols, rows } = getTerminalDimensions(); writeCurrentPtysize(cols, rows); } // when current tty resized, resize both pty and xterm proxy process.stdout.on("resize", () => { const { cols, rows } = getTerminalDimensions(); shell.resize(cols, rows); xtermProxy.resize(cols, rows); writeCurrentPtysize(cols, rows); }); const isStillWorkingQ = () => { const rendered = xtermProxy.tail(24).replace(/\s+/g, " "); return conf.working?.some((rgx) => rgx.test(rendered)); }; // Heartbeat for auto-response on rendered terminal output // This catches patterns that appear via CSI positioning instead of newlines let lastHeartbeatRendered = ""; // Auto-retry backoff state (mirrors rs/src/context.rs). `streak` doubles the // backoff each consecutive failed retry; `startedAt` anchors the 8h give-up // window; `nextAt` is non-null while a retry is scheduled. `autoRetryScreen` // is the latest rendered screen captured by the stdout pipeline below — the // heartbeat timer reads it (its own xtermProxy.tail() is empty for EOL CLIs). let retryStreak = 0; let retryStartedAt: number | null = null; let retryNextAt: number | null = null; let autoRetryScreen = ""; // Paraphrased reason captured when the error banner was matched — folded // into the typed retry message so the agent knows why it is being nudged. let retryReason: string | null = null; const heartbeatInterval = setInterval(async () => { try { const rendered = removeControlCharacters(xtermProxy.tail(12)); // Auto-retry backoff timer — fires the scheduled "retry" using the latest // rendered screen captured by the stdout pipeline (consoleResponder). Runs // every tick (independent of output) so it still fires while the agent sits // idle on an error. Arming/reset lives in the stdout pipeline because this // heartbeat's own xtermProxy.tail() is empty for newline (EOL) CLIs like // claude. Only types "retry" when idle at a prompt (never mid-work). if (retryNextAt !== null) { const now = Date.now(); if (retryStartedAt !== null && now - retryStartedAt >= AUTO_RETRY_GIVE_UP_MS) { logger.warn(`[${cli}-yes] auto-retry: giving up after 8h with no recovery`); retryNextAt = null; retryStartedAt = null; retryStreak = 0; } else if (now >= retryNextAt) { const working = conf.working?.some((rx: RegExp) => rx.test(autoRetryScreen)) ?? false; const readyNow = conf.ready?.some((rx: RegExp) => rx.test(autoRetryScreen)) ?? false; // Also require a few quiet seconds on top of the backoff delay, so a // scheduled retry doesn't collide with a line the user is actively // typing into the prompt (see AUTO_RETRY_MIN_IDLE_MS). const idleMs = ctx.idleWaiter.idleTimeMs(); if (!shouldFireRetry(working, readyNow, idleMs, AUTO_RETRY_MIN_IDLE_MS)) { retryNextAt = now + 500; // busy / not at prompt / still active — re-check shortly } else { retryStreak += 1; const nextBackoffMs = autoRetryBackoffMs(retryStreak); const reason = retryReason ?? AUTO_RETRY_REASON_FALLBACK; const sinceFirstSecs = retryStartedAt === null ? 0 : (now - retryStartedAt) / 1000; const line = buildAutoRetryMessage( retryStreak, reason, sinceFirstSecs, nextBackoffMs / 1000, ); logger.warn( `[${cli}-yes] auto-retry: typing retry nudge (attempt ${retryStreak}, reason: ${reason})`, ); // Write the nudge + Enter atomically (mirrors rs do_send_retry); using // sendMessage would split text/Enter across the fast heartbeat ticks. ctx.messageContext.shell.write(line + "\r"); ctx.idleWaiter.ping(); // Structured trace for `ay msgs` / the console (mirrors the Rust // runtime's record_auto_retry_inbox) — best-effort, never blocks. void recordInbox({ at: now, origin: "wrapper", from: null, to: { pid: process.pid, cli, cwd: workingDir }, kind: "auto-retry", body: `${reason} (attempt ${retryStreak}, next backoff ${formatDurationSecs(nextBackoffMs / 1000)})`, wrapped: false, }); // Self-schedule the next retry with escalated backoff. (Leaving nextAt // null and re-arming from the stdout pipeline would tight-loop while the // error banner stays on screen.) Reset on recovery cancels this. retryNextAt = now + nextBackoffMs; } } } // Skip if output hasn't changed since last heartbeat if (rendered === lastHeartbeatRendered) return; lastHeartbeatRendered = rendered; const lines = rendered.split("\n").filter((line) => line.trim()); for (const line of lines) { // ready matcher: if matched, mark stdin ready if (conf.ready?.some((rx: RegExp) => rx.test(line))) { logger.debug(`heartbeat|ready |${line}`); ctx.stdinReady.ready(); ctx.stdinFirstReady.ready(); } // enter matchers: send Enter when any enter regex matches if (conf.enter?.some((rx: RegExp) => rx.test(line))) { logger.debug(`heartbeat|sendEnter matched|${line}`); await sendEnter(ctx.messageContext, 400); continue; } // typingRespond matcher: if matched, send the specified message const typeingRespondMatched = Object.entries(conf.typingRespond ?? {}).filter( ([_sendString, onThePatterns]) => onThePatterns.some((rx) => rx.test(line)), ); if (typeingRespondMatched.length) { await sflow(typeingRespondMatched) .map( async ([sendString]) => await sendMessage(ctx.messageContext, sendString, { waitForReady: false }), ) .toCount(); continue; } // fatal matchers: set isFatal flag when matched if (conf.fatal?.some((rx: RegExp) => rx.test(line))) { logger.debug(`heartbeat|fatal |${line}`); ctx.isFatal = true; await exitAgent(); break; } // restartWithoutContinueArg matchers: set flag to restart without continue args if (conf.restartWithoutContinueArg?.some((rx: RegExp) => rx.test(line))) { logger.debug(`heartbeat|restart-without-continue|${line}`); ctx.shouldRestartWithoutContinue = true; ctx.isFatal = true; await exitAgent(); break; } // session ID capture for codex if (isCodexFamily(cli)) { const sessionId = extractSessionId(line); if (sessionId) { logger.debug(`heartbeat|session|captured session ID: ${sessionId}`); await storeSessionForCwd(workingDir, sessionId); } } } } catch (error) { // Silently ignore heartbeat errors to avoid disrupting main flow logger.debug(`heartbeat|error: ${error}`); } }, 100); // Run every 100ms — cheap when unchanged (see the rendered === lastHeartbeatRendered // guard above): most ticks bail after one xtermProxy.tail(12) + string compare. A short // interval matters for two things this heartbeat drives: no-EOL ready/fatal detection // (CSI-redrawn output never fires the newline-driven consoleResponder path) and auto-retry // backoff timing precision (AUTO_RETRY_MIN_IDLE_MS gating). Previously 800ms; still coarser // than Rust's HEARTBEAT_INTERVAL_MS=50 (rs/src/context.rs) since Rust's per-tick cost is lower. // Clear heartbeat on exit const cleanupHeartbeat = () => clearInterval(heartbeatInterval); shell.onExit(cleanupHeartbeat); // Report to whoever spawned us when this agent settles idle (done), parks on a // question (stuck), or exits. `` asks the AGENT to do this; this // loop is the guarantee, because an agent that has run out of context, crashed, // or simply forgotten is exactly the case the parent is waiting on. No-ops // entirely for a top-level agent. See ts/parentPingLoop.ts. const { startParentPingLoop } = await import("./parentPingLoop.ts"); const parentPingLoop = startParentPingLoop({ parentPid, selfWrapperPid: process.pid, self: { cli, pid: shell.pid, cwd: workingDir, prompt: originalPrompt ?? null }, patterns: { ready: conf.ready, working: conf.working, needsInput: conf.needsInput }, // 48 lines: enough to carry a menu plus the question above it, matching what // `ay notifyd` sends as evidence. screen: () => removeControlCharacters(xtermProxy.tail(48)).split("\n"), }); if (exitOnIdle) (async () => { while (true) { await ctx.idleWaiter.wait(exitOnIdle); await pidStore.updateStatus(shell.pid, "idle").catch(() => null); if (isStillWorkingQ()) { logger.warn(`[${cli}-yes] ${cli} is idle, but seems still working, not exiting yet`); continue; } if (idleAction) { logger.info(`[${cli}-yes] ${cli} is idle, performing idle action: ${idleAction}`); notifyWebhook("IDLE", `action=${idleAction}`, workingDir).catch(() => null); await sendMessage(ctx.messageContext, idleAction); continue; } logger.info(`[${cli}-yes] ${cli} is idle, exiting...`); notifyWebhook("IDLE", "", workingDir).catch(() => null); await exitAgent(); break; } })(); // Message streaming // Message streaming with stdin and optional FIFO (Linux only) // read stdin stream // CRITICAL FIX: fromReadable() from 'from-node-stream' doesn't work properly with stdin // because it doesn't handle Node.js stream modes correctly. We create a custom ReadableStream // that properly manages stdin's flowing mode and event listeners. const stdinStream = new ReadableStream( { start(controller) { // Set up stdin in flowing mode so 'data' events fire process.stdin.resume(); let closed = false; // Handle data events const dataHandler = (chunk: Buffer) => { try { controller.enqueue(chunk); } catch { // Ignore enqueue errors (stream may be closed) } }; // Handle end/close - both events can fire, so track state const endHandler = () => { if (closed) return; closed = true; try { controller.close(); } catch { // Ignore close errors (already closed) } }; const errorHandler = (err: Error) => { if (closed) return; closed = true; try { controller.error(err); } catch { // Ignore error after close } }; process.stdin.on("data", dataHandler); process.stdin.on("end", endHandler); process.stdin.on("close", endHandler); process.stdin.on("error", errorHandler); }, cancel(_reason) { process.stdin.pause(); }, }, { highWaterMark: 16 }, ); let aborted = false; await sflow(stdinStream) .map((buffer) => { const str = buffer.toString(); // CRITICAL FIX: Handle Ctrl+C directly in map instead of forkTo // The previous implementation used .forkTo() which created a separate stream branch // that wasn't being consumed properly, causing Ctrl+C to never be detected. const CTRL_Z = "\u001A"; const CTRL_C = "\u0003"; // handle CTRL+Z and filter it out (not supported yet) if (!aborted && str === CTRL_Z) { return ""; } // handle CTRL+C when stdin is not ready (agent is loading) if (!aborted && !ctx.stdinReady.isReady && str === CTRL_C) { logger.error("User aborted: SIGINT"); shell.kill("SIGINT"); pendingExitCode.resolve(130); // SIGINT exit code aborted = true; return str; // still pass to agent, but they'll probably be killed } // Track Ctrl+C when stdin is ready (user is interrupting running CLI) if (str === CTRL_C) { userSentCtrlC = true; // Reset flag after 2 seconds in case CLI doesn't exit immediately setTimeout(() => { userSentCtrlC = false; }, 2000); } return str; }) // Detect Ctrl+Y or /auto command to toggle auto-yes mode .map( (() => { let line = ""; const toggleAutoYes = () => { ctx.autoYesEnabled = !ctx.autoYesEnabled; // When switching to manual mode, mark stdin ready so user keystrokes are not blocked if (!ctx.autoYesEnabled) { ctx.stdinReady.ready(); ctx.stdinFirstReady.ready(); } const status = ctx.autoYesEnabled ? "\x1b[32m[auto-yes: ON]\x1b[0m" : "\x1b[33m[auto-yes: OFF]\x1b[0m"; process.stderr.write(`\r${status} (Ctrl+Y to toggle)\n`); }; return (data: string) => { let out = ""; for (const ch of data) { // Ctrl+Y (\x19) toggles auto-yes immediately if (ch === "\x19") { toggleAutoYes(); // Do not forward Ctrl+Y to the PTY continue; } // Handle Enter if (ch === "\r" || ch === "\n") { // Only check for /auto if line is short enough if (line.length <= 20) { const cleanLine = line // oxlint-disable-next-line no-control-regex -- intentional: strip ANSI/control chars .replace(/[\x00-\x1f]|\x1b\[[0-9;]*[A-Za-z]|\[[A-Z]/g, "") .trim(); if (cleanLine === "/auto") { out += "\x15"; // Ctrl+U to clear the /auto text from shell input toggleAutoYes(); line = ""; continue; } } line = ""; out += ch; continue; } // Handle backspace if (ch === "\x7f" || ch === "\b") { if (line.length > 0) line = line.slice(0, -1); out += ch; continue; } // Track only printable ASCII for line, with size limit if (ch >= " " && ch <= "~" && line.length < 50) line += ch; out += ch; } return out; }; })(), ) // Read from IPC stream if available (FIFO on Linux, Named Pipes on Windows) .by(async (s) => { if (!useStdinAppend) return s; const fifoPath = pidStore.getFifoPath(shell.pid); const ipcResult = await createFifoStream(cli, fifoPath); if (!ipcResult) return s; pendingExitCode.promise.finally(async () => await ipcResult[Symbol.asyncDispose]()); process.stderr.write(`\n Append prompts: ${cli}-yes --append-prompt '...'\n\n`); return s.merge(ipcResult.stream); }) .confluenceByConcat() // necessary because .by() above is async // .map((e) => e.replaceAll('\x1a', '')) // remove ctrl+z from user's input, to prevent bug (but this seems bug) // .forEach(e => appendFile('.cache/io.log', "input |" + JSON.stringify(e) + '\n')) // for debugging .onStart(async function promptOnStart() { // send prompt when start logger.debug("Sending prompt message: " + JSON.stringify(prompt)); if (prompt) await sendMessage(ctx.messageContext, prompt); }) // pipe content by shell .by({ writable: new WritableStream({ write: async (data) => { await ctx.stdinReady.wait(); shell.write(data); // Forwarded user input counts as activity too (mirrors the Rust // runtime's ping on stdin forward) — the auto-retry idle gate below // must not fire while the user is actively typing. ctx.idleWaiter.ping(); }, }), readable: xtermProxy.readable, }) .forEach((chunk) => { // Only ping activity if there's visible content (not just ANSI/cursor // control sequences) — mirrors the Rust runtime's handle_output gate. // Without this, periodic control-only chatter (e.g. cursor position // queries) would keep resetting the auto-retry idle clock and the // scheduled retry could defer indefinitely without ever firing. if (removeControlCharacters(chunk).trim()) ctx.idleWaiter.ping(); pidStore.updateStatus(shell.pid, "active").catch(() => null); ctx.nextStdout.ready(); }) .forkTo(async function rawLogger(f) { const rawLogPath = ctx.logPaths.rawLogPath; if (!rawLogPath) return f.run(); // no stream // try stream the raw log for realtime debugging, including control chars, note: it will be a huge file return await mkdir(path.dirname(rawLogPath), { recursive: true }) .then(async () => { logger.debug(`[${cli}-yes] raw logs streaming to ${rawLogPath}`); // Track size (seeded from any pre-existing log) to cap disk growth // without a stat() per chunk. See compactRawLogTail. let rawWritten = await stat(rawLogPath) .then((s) => s.size) .catch(() => 0); return f .forEach(async (chars) => { // Detect alt-screen enter (DECSET 1049/1047/47) so we know whether // the rendered log can safely stand in for this raw log on exit. if (!usedAltScreen && /\[\?(?:1049|1047|47)h/.test(chars)) { usedAltScreen = true; } await writeFile(rawLogPath, chars, { flag: "a" }).catch(() => null); rawWritten += Buffer.byteLength(chars); if (rawWritten > RAW_LOG_TRIGGER_BYTES) { rawWritten = await compactRawLogTail(rawLogPath).catch(() => rawWritten); } }) .run(); }) .catch(() => f.run()); }) // handle cursor position requests and render terminal output .by(function consoleResponder(e) { // TODO: wait for cli ready and send prompt if provided // if (cli === "codex" && !process.stdin.isTTY) shell.write(`\u001b[1;1R`); // send cursor position response when stdin is not tty let lastRendered = ""; return ( e // Terminal query responses (DA, DSR, etc.) are handled automatically // by XtermProxy via @xterm/headless — no ad-hoc interception needed. .forEach(async (line, lineIndex) => { // ============ respond on rendered screen const rendered = xtermProxy.tail(24); // Skip processing if output hasn't changed if (rendered === lastRendered) return; lastRendered = rendered; logger.debug(`stdout|${line}`); // Auto-retry on recoverable API errors (overload / rate-limit / usage- // limit): arm/reset the backoff on the whole rendered screen (the error // banner and the ready prompt are on different lines, so this can't be a // per-line check). The firing happens on the heartbeat timer, which // reads `autoRetryScreen`. Done here, before the `fatal` check below, so // these recoverable errors retry instead of exiting. if (conf.autoRetry?.length) { autoRetryScreen = rendered; const errVisible = conf.autoRetry.some((rx: RegExp) => rx.test(rendered)); const readyVisible = conf.ready?.some((rx: RegExp) => rx.test(rendered)) ?? false; if (errVisible && readyVisible) { // Remember WHY (paraphrased — see classifyAutoRetryReason) so // the typed message can explain itself. Refresh on every match: // the banner may change across attempts (e.g. overload → 5xx). retryReason = classifyAutoRetryReason(rendered); if (retryNextAt === null) { if (retryStartedAt === null) retryStartedAt = Date.now(); const delayMs = autoRetryBackoffMs(retryStreak); retryNextAt = Date.now() + delayMs; logger.warn( `[${cli}-yes] auto-retry armed: recoverable error detected, retrying in ${ delayMs / 1000 }s (attempt ${retryStreak + 1})`, ); } } else if (readyVisible && !errVisible && retryStartedAt !== null) { logger.debug(`[${cli}-yes] auto-retry: recovered, resetting backoff`); retryStreak = 0; retryStartedAt = null; retryNextAt = null; } } // ready matcher: if matched, mark stdin ready if (conf.ready?.some((rx: RegExp) => line.match(rx))) { logger.debug(`ready |${line}`); ctx.stdinReady.ready(); ctx.stdinFirstReady.ready(); } // enter matchers: send Enter when any enter regex matches if (conf.enter?.some((rx: RegExp) => line.match(rx))) { logger.debug(`sendEnter matched|${line}`); return await sendEnter(ctx.messageContext, 400); // wait for idle for a short while and then send Enter } // typingRespond matcher: if matched, send the specified message const typeingRespondMatched = Object.entries(conf.typingRespond ?? {}).filter( ([_sendString, onThePatterns]) => onThePatterns.some((rx) => line.match(rx)), ); const typingResponded = typeingRespondMatched.length && (await sflow(typeingRespondMatched) .map( async ([sendString]) => await sendMessage(ctx.messageContext, sendString, { waitForReady: false }), ) .toCount()); if (typingResponded) return; // fatal matchers: set isFatal flag when matched if (conf.fatal?.some((rx: RegExp) => line.match(rx))) { logger.debug(`fatal |${line}`); ctx.isFatal = true; await exitAgent(); } // restartWithoutContinueArg matchers: set flag to restart without continue args if (conf.restartWithoutContinueArg?.some((rx: RegExp) => line.match(rx))) { logger.debug(`restart-without-continue|${line}`); ctx.shouldRestartWithoutContinue = true; ctx.isFatal = true; // also set fatal to trigger exit await exitAgent(); } // session ID capture for codex if (isCodexFamily(cli)) { const sessionId = extractSessionId(line); if (sessionId) { logger.debug(`session|captured session ID: ${sessionId}`); await storeSessionForCwd(workingDir, sessionId); } } }) ); }) // auto-response // .forkTo(function autoResponse(e) { // return ( // e // .map((e) => removeControlCharacters(e)) // // .map((e) => e.replaceAll("\r", "")) // remove carriage return // .by((s) => { // if (conf.noEOL) return s; // codex use cursor-move csi code insteadof \n to move lines, so the output have no \n at all, this hack prevents stuck on unended line // return s.lines({ EOL: "NONE" }); // other clis use ink, which is rerendering the block based on \n lines // }) // // Generic auto-response handler driven by CLI_CONFIGURES // .forEach(async (line, lineIndex) => // createAutoResponseHandler(line, lineIndex, { ctx, conf, cli, workingDir, exitAgent }), // ) // .run() // ); // }) .by((s) => (removeControlCharactersFromStdout ? s.map((e) => removeControlCharacters(e)) : s)) // terminate whole stream when shell did exited (already crash-handled) .by(createTerminatorStream(pendingExitCode.promise)) .to(fromWritable(process.stdout)); const renderedSaved = await saveLogFile(ctx.logPaths.logPath, xtermProxy.render()); // The raw byte log exists for live tailing during the run. Once a non-empty // rendered log is durably written — and the session stayed on the normal // screen buffer, so the render holds the full scrollback — the raw log is // redundant: drop it and repoint the index at the rendered log. On crash / // empty render / alt-screen we keep the raw log; the startup sweep prunes it. if (renderedSaved && !usedAltScreen && ctx.logPaths.rawLogPath && ctx.logPaths.logPath) { await unlink(ctx.logPaths.rawLogPath).catch(() => null); await pidStore.markRendered(shell.pid, ctx.logPaths.logPath).catch(() => null); } // and then get its exitcode const exitCode = await pendingExitCode.promise; logger.info(`[${cli}-yes] ${cli} exited with code ${exitCode}`); // Final report to the parent — awaited (not fire-and-forget) so the message is // actually handed to `ay send` before this process goes away. Bounded by // deliverPing's own timeout and never throws. await parentPingLoop.pingExit(exitCode ?? null); // Final pidStore cleanup await pidStore.close(); // Capture final render before disposing xterm proxy const finalRender = xtermProxy.render(); xtermProxy.dispose(); // deprecated logFile option, we have logPath now, but keep for backward compatibility await saveDeprecatedLogFile(logFile, finalRender, verbose); return { exitCode, logs: finalRender }; async function exitAgent() { ctx.robust = false; // disable robust to avoid auto restart // send exit command to the shell, must sleep a bit to avoid claude treat it as pasted input for (const cmd of cliConf.exitCommands ?? ["/exit"]) await sendMessage(ctx.messageContext, cmd); // wait for shell to exit or kill it with a timeout let exited = false; await Promise.race([ pendingExitCode.promise.then(() => (exited = true)), // resolve when shell exits // if shell doesn't exit in 2 seconds, kill it. Rust's equivalent (rs/src/context.rs) // doesn't wait for the child's own exit at all — it sends the exit command(s) and tears // down immediately. 2s keeps a real grace window for a CLI to actually process `/exit` // (flush session state, close connections) while still bounding the worst case — down // from a previous 5s that mostly just delayed force-killing CLIs that never respond. new Promise((resolve) => setTimeout(() => { if (exited) return; // if shell already exited, do nothing shell.kill(); // kill the shell process if it doesn't exit in time resolve(); }, 2000), ), // 2 seconds timeout ]); } function getTerminalDimensions() { if (!process.stdout.isTTY) return { cols: 80, rows: 24 }; // default size when not tty return { // Enforce minimum 20 columns to avoid layout issues cols: Math.max(20, process.stdout.columns), // Clamp rows too: a tty whose winsize was never set (e.g. a pty created by // `script`/a daemon) reports 0, and bun-pty rejects rows <= 0 outright — // the spawn then fails with a bare "PTY spawn failed" and no clue why. // Mirrors get_terminal_size() in rs/src/pty_spawner.rs. rows: Math.max(4, process.stdout.rows), }; } } function sleep(ms: number) { return new Promise((resolve) => setTimeout(resolve, ms)); }