/** * `ay ls / read / cat / tail / head / send` subcommand implementations. * * Mirrors the principles of koho's `terminal-ws-lib.ts` (session list, render * via @xterm/headless, keyword-keyed input) — but file-based instead of via * a daemon. Reads ~/.agent-yes/pids.jsonl (cross-runtime global index, written * by both the TS PidStore and the Rust pid_store::PidStore) and the per-pid * raw log files. * * Returns null when argv[2] is not a known subcommand so cli.ts falls through * to the normal agent-spawning flow. */ import { randomBytes } from "crypto"; import { closeSync, constants as fsConstants, openSync, realpathSync } from "fs"; import { execFileSync } from "node:child_process"; import { appendFile, mkdir, open, readFile, stat, writeFile } from "fs/promises"; import ms from "ms"; import { homedir } from "os"; import path from "path"; import { type GlobalPidRecord, readGlobalPids, updateGlobalPidStatus } from "./globalPidIndex.ts"; import { formatIdentity, localHost, localUser } from "./identity.ts"; import { buildAgentForest, flattenForest } from "./agentTree.ts"; import { parseTaskCounts, type TaskCounts } from "./todoParse.ts"; import { agentYesHome } from "./agentYesHome.ts"; import { PidStore } from "./pidStore.ts"; import { type MailParty, type MessageRecord, partyMatches, readMailbox, recordMessage, recordOutbox, senderLabel, type ObservedSender, } from "./messageLog.ts"; import { badgeLabel, matchBadges, TYPING_BADGE } from "./badges.ts"; import { classifyNeedsInput, isWorkingScreen, parseMenu, type MenuState, type NeedsInput, } from "./needsInput.ts"; import { diffLsStates, type LiveState, type LsAgentState } from "./lsWatch.ts"; import { filterSinceSeq, filterSinceTs, filterUnread, maxSeq, postmortemStartedAt, type NotifyEvent, } from "./notifyInbox.ts"; import { clearWatcher, getCursor, heartbeatWatcher, hostId, readInbox, setCursor, } from "./notifyStore.ts"; import { buildStoredResult, normalizeEnvelope, resultPath, resultsDir, type StoredResult, } from "./resultEnvelope.ts"; import { loadSharedCliDefaults } from "./configShared.ts"; import { invokedCliName } from "./invokedCli.ts"; import type { AgentCliConfig } from "./index.ts"; import { framePaste as frameAsPaste, shouldFramePaste } from "./bracketedPaste.ts"; import { ageMatchesRegistration, findAgentAncestor, pidOwnershipVerdict, readAncestryTable, type SenderVia, } from "./senderAncestry.ts"; import yargs from "yargs"; import { type ResolvedRemote, readRemotes, resolveRemoteSpec } from "./remotes.ts"; import { noteRemoteResult, pruneRemoteHealth, readRemoteHealth, remoteBackoffMs, shouldSkipRemote, writeRemoteHealth, } from "./remoteHealth.ts"; import { isWebrtcSpec } from "./webrtcLink.ts"; import { withIpcLock } from "./ipcLock.ts"; // --------------------------------------------------------------------------- // notes store (~/.agent-yes/notes.jsonl) // --------------------------------------------------------------------------- function notesPath(): string { const dir = process.env.AGENT_YES_HOME ?? path.join(homedir(), ".agent-yes"); return path.join(dir, "notes.jsonl"); } export async function readNotes(): Promise> { let raw: string; try { raw = await readFile(notesPath(), "utf-8"); } catch { return new Map(); } const map = new Map(); for (const line of raw.split("\n")) { const t = line.trim(); if (!t) continue; try { const { pid, note } = JSON.parse(t); if (typeof pid === "number") { if (note) map.set(pid, note); else map.delete(pid); } } catch { /* skip */ } } return map; } async function writeNote(pid: number, note: string): Promise { const p = notesPath(); await mkdir(path.dirname(p), { recursive: true }); await appendFile(p, JSON.stringify({ pid, note, updated_at: Date.now() }) + "\n"); } async function compactNotes(): Promise { const map = await readNotes(); const lines = Array.from(map.entries()) .map(([pid, note]) => JSON.stringify({ pid, note, updated_at: Date.now() })) .join("\n"); await writeFile(notesPath(), lines ? lines + "\n" : ""); } // --------------------------------------------------------------------------- // read-recency store (~/.agent-yes/reads.jsonl) // // Records that some sender "read" (tailed/cat'd) a target agent at a time, so // `ay send` can refuse to fire at an agent the sender hasn't actually looked at // recently. This is the guard against a fuzzy keyword silently resolving to the // wrong agent (e.g. you tail "babaiban" but a send resolves to "qq-cli"). // --------------------------------------------------------------------------- export const READ_WINDOW_MS = 60_000; // "read recently" = within the last minute // Max time writeToIpc will keep retrying a backed-up FIFO before erroring. A live // agent drains its stdin in milliseconds; only a wedged reader hits this. const IPC_WRITE_TIMEOUT_MS = 10_000; const READS_KEY_SEP = "\0"; function readsPath(): string { const dir = process.env.AGENT_YES_HOME ?? path.join(homedir(), ".agent-yes"); return path.join(dir, "reads.jsonl"); } // Each line: {"by":"agent:123"|"human","target":456,"at":}. Append-only, // last-per-(by,target) wins; compacted opportunistically so a long `tail -f` // (which refreshes its marker) can't grow the file without bound. async function readReads(): Promise> { let raw: string; try { raw = await readFile(readsPath(), "utf-8"); } catch { return new Map(); } const map = new Map(); for (const line of raw.split("\n")) { const t = line.trim(); if (!t) continue; try { const { by, target, at } = JSON.parse(t); if (typeof by === "string" && typeof target === "number" && typeof at === "number") map.set(`${by}${READS_KEY_SEP}${target}`, at); } catch { /* skip corrupt */ } } return map; } async function recordRead(by: string, target: number): Promise { const p = readsPath(); try { await mkdir(path.dirname(p), { recursive: true }); await appendFile(p, JSON.stringify({ by, target, at: Date.now() }) + "\n"); // Opportunistic compaction once the append-only log grows past a small cap. const raw = await readFile(p, "utf-8").catch(() => ""); if (raw.split("\n").length > 200) { const map = await readReads(); const lines = [...map.entries()] .map(([k, at]) => { const i = k.indexOf(READS_KEY_SEP); return JSON.stringify({ by: k.slice(0, i), target: Number(k.slice(i + 1)), at }); }) .join("\n"); await writeFile(p, lines ? lines + "\n" : ""); } } catch { /* best-effort: the guard degrades to a warning if state can't be written */ } } export async function lastReadAt(by: string, target: number): Promise { const map = await readReads(); return map.get(`${by}${READS_KEY_SEP}${target}`) ?? null; } /** Recent agent→agent read/tail edges (skips "human" readers), for the /rgui * relationship-wire view. `by`/`target` are pids. */ export interface ReadEdge { by: number; target: number; at: number; } export async function recentReadEdges(windowMs = READ_WINDOW_MS): Promise { const now = Date.now(); const map = await readReads(); const out: ReadEdge[] = []; for (const [key, at] of map) { if (now - at > windowMs) continue; const i = key.indexOf(READS_KEY_SEP); const by = key.slice(0, i); if (!by.startsWith("agent:")) continue; // agent→agent only const byPid = Number(by.slice("agent:".length)); const target = Number(key.slice(i + 1)); if (byPid && target && byPid !== target) out.push({ by: byPid, target, at }); } return out; } /** Recent agent→agent MESSAGE edges (a delivered `ay send`/`key`/`select`), for * the /rgui + /w wire view. Like {@link ReadEdge} but sourced from the per-cwd * outbox logs and carrying the send `kind`. `by`/`target` are the sender/recipient * pids AT SEND TIME (a restarted agent keeps a new pid — matched best-effort). */ export interface MessageEdge { by: number; target: number; at: number; /** Derived, not re-listed: a hand-copied union here silently broke when * `MessageRecord["kind"]` grew a member (#453 added "terminal"). */ kind?: MessageRecord["kind"]; } /** * Scan every live agent's outbox for sends within `windowMs` and return them as * directional edges (newest `at` per by→target pair wins). Bounded by the number * of distinct agent cwds — the outbox is per-cwd, so each dir is read once. */ export async function recentMessageEdges(windowMs = READ_WINDOW_MS): Promise { const now = Date.now(); const records = await listRecords(undefined, { all: true, active: false, json: false, latest: false, cwdScope: null, }); const cwds = [...new Set(records.map((r) => r.cwd).filter(Boolean))]; const best = new Map(); await Promise.all( cwds.map(async (cwd) => { for (const rec of await readMailbox(cwd, "outbox")) { if (now - rec.at > windowMs) continue; const by = rec.from?.pid; const target = rec.to?.pid; if (!by || !target || by === target) continue; // agent→agent only const key = `${by}\0${target}`; const prev = best.get(key); if (!prev || rec.at > prev.at) best.set(key, { by, target, at: rec.at, kind: rec.kind }); } }), ); return [...best.values()]; } // Identify the sender. An agent launched by `ay` inherits AGENT_YES_PID=; the registered agent record carries that same wrapper_pid, so we map the // env value back to the agent's own canonical record. Falls back to a direct pid // match (back-compat), then null when there's no agent context (a human shell). /** * What can be measured about THIS process, with no cooperation from anyone. * * Recorded on every send whether or not an agent was identified, so a receiver * facing an unattributed message still has facts to act on instead of a blank. * Deliberately cheap and total — no probing, no failure mode, nothing a caller * can influence. */ function observedSender(): ObservedSender { return { user: localUser(), host: localHost(), cwd: process.cwd(), pid: process.pid, }; } /** * The repository a directory belongs to, as git itself defines it: the common * dir, which every linked worktree of one repository shares and no separate * clone does. Null when the path is not in a repository or git cannot answer. * * Synchronous and cached: it is consulted at most twice per send, and only on * the path where the cheaper directory check already failed. */ const _commonDirCache = new Map(); function gitCommonDir(dir: string): string | null { const hit = _commonDirCache.get(dir); if (hit !== undefined) return hit; let out: string | null = null; try { out = execFileSync("git", ["-C", dir, "rev-parse", "--path-format=absolute", "--git-common-dir"], { encoding: "utf8", timeout: 5000, stdio: ["ignore", "pipe", "ignore"], }) .trim() .replace(/\/+$/, "") || null; } catch { out = null; } _commonDirCache.set(dir, out); return out; } async function resolveSender(): Promise { return (await resolveSenderVia()).agent; } /** * The calling agent AND how that was established. * * `AGENT_YES_PID` is a claim the wrapper passes along; when it is absent — an * SDK / `claude -p` session shelling out, anything re-exec'd through a scrubbed * environment — this used to answer "no agent", which a receiver cannot tell * apart from an anonymous stranger. The process tree is the fallback because it * is a FACT the caller cannot forge: a lane's `ay send` is a descendant of that * lane however many shells deep it runs. See ts/senderAncestry.ts. * * `via` is reported, never inferred by the reader, and "observed" (nothing * identified the caller) stays expressible — no plausible sender is invented to * fill the field. */ let _senderVia: Promise<{ agent: GlobalPidRecord | null; via: SenderVia }> | null = null; /** * Memoized for the process lifetime. Who is calling cannot change inside one * CLI invocation, and a single `ay send` asks twice — once to size the envelope * for the cap check, once to attribute the message. Without this the `ps` and * the registry read are both paid twice for one send. */ function resolveSenderVia(): Promise<{ agent: GlobalPidRecord | null; via: SenderVia }> { return (_senderVia ??= computeSenderVia()); } async function computeSenderVia(): Promise<{ agent: GlobalPidRecord | null; via: SenderVia }> { const recs = await listRecords(undefined, { all: true, active: false, json: false, latest: false, cwdScope: null, }); const byPid = (pid: number) => recs.find((r) => r.wrapper_pid === pid) ?? recs.find((r) => r.pid === pid) ?? null; const envPid = process.env.AGENT_YES_PID ? Number(process.env.AGENT_YES_PID) : null; const declared = envPid && !Number.isNaN(envPid) ? byPid(envPid) : null; // ONE `ps` for both jobs below — the ancestry walk and the pid-reuse check — // rather than a spawn each. ~40ms for 887 processes, and it is skipped // entirely on the fast path when the env already resolved and is corroborated // by the cheapest possible evidence. const table = await readAncestryTable(); // pids are reused: a number matching a registered agent is not proof the // process at that number IS that agent. Require it to be at least as old as // the registration, so a pid handed to something new cannot inherit an // identity. Unknown age ⇒ not attributed. const isLiveAgent = (pid: number) => { const rec = byPid(pid); if (!rec) return null; return ageMatchesRegistration(table?.get(pid)?.ageSecs, rec.started_at) ? rec : null; }; const inherited = table ? await findAgentAncestor(process.pid, isLiveAgent, { readTable: async () => table }) : null; if (declared) { // The honest path is unchanged in OUTCOME: the wrapper that injects // AGENT_YES_PID is an ancestor of the `ay send` it spawns, so a normal lane // corroborates and renders exactly as before. // The claimed lane is our ancestor: claim and kernel agree. if (inherited?.pid === declared.pid) return { agent: declared, via: "env" }; // It is not. That alone does not say which of two very different things is // happening, and the earlier rule — "is some OTHER registered lane my // ancestor" — picked the wrong discriminator. Measured on real traffic it // was inverted: it fired on a lane's own nested processes (same tree, // benign) and stayed silent while a different worktree sent two documents // under this lane's name. // // The working directory separates them. A lane's own detached helper runs // inside that lane's tree; a process in a DIFFERENT tree claiming the lane // is the case the marker exists for. cwd is observed, not asserted — the // body cannot set it. // BOTH sides are realpath'd first. macOS resolves /tmp -> /private/tmp and // /var -> /private/var, so a lane registered under a symlinked path would // otherwise never appear to contain its own processes and every honest send // from it would be accused. Caught by a same-tree test that went loud. // Best-effort: an unreadable path falls back to the raw string rather than // failing the send. const real = (p2: string): string => { try { return realpathSync(p2).replace(/\/+$/, ""); } catch { return p2.replace(/\/+$/, ""); } }; const here = real(process.cwd()); const claimedRoot = declared.cwd ? real(declared.cwd) : ""; const insideClaimedTree = Boolean(claimedRoot) && (here === claimedRoot || here.startsWith(claimedRoot + path.sep)); // "Inside the tree" is not the same as "belongs to that lane". A lane // routinely works in a LINKED WORKTREE that is a sibling directory, not a // child — `~/ws/org/_wt/feature-x` alongside `~/ws/org/repo/tree/dev`. A // bare `ay send` from a shell there inherits AGENT_YES_PID honestly and // would be accused on a path check alone. // // git already answers this exactly: linked worktrees of one repository // share a git common dir, while a separate clone has its own. Measured on // the fleet that reported it — the sibling worktree resolved to the lane's // own `.git`, and the worktree that had been misattributing resolved to a // different one. So the same probe separates the honest case from the case // this marker exists for, which a "same parent directory" heuristic would // not: both live under the same ancestor. // // Only consulted when the path check already failed, so the honest common // path pays nothing. // CONTAINMENT, not equality. A lane's submodule worktrees have a common dir // NESTED inside the lane's own: // // lane …/tree/dev/.git // its submodule wt …/tree/dev/.git/modules/lib/desktop <- descendant // a separate clone …/tree/billings/.git <- neither // // Equality would mark the first case LOUD — the same false positive one // layer down, reported from the fleet that runs that shape. Checked in both // directions because the mirror layout (a lane registered in the submodule // worktree, sending from the parent tree) is equally legitimate and would // otherwise wait to be discovered by firing. A separate clone is under // neither, so detection is unchanged. const sameRepo = () => { if (!claimedRoot) return false; const a = gitCommonDir(here); const b = gitCommonDir(claimedRoot); if (a === null || b === null) return false; return a === b || a.startsWith(b + path.sep) || b.startsWith(a + path.sep); }; if (!insideClaimedTree && !sameRepo()) return { agent: declared, via: "env-uncorroborated" }; // Inside the claimed lane's own tree, or nothing to compare against: // unverifiable, not suspicious, and must not be dressed as it. return { agent: declared, via: "env-unverified" }; } // A set-but-unresolvable AGENT_YES_PID (stale env, aged-out record) is not a // reason to stop: the process tree can still say who this is. if (inherited) return { agent: inherited, via: "ancestry" }; return { agent: null, via: "observed" }; } // The (key, agent) pair used to attribute reads and gate sends. Agents get a // stable per-agent key; a human shell shares the "human" bucket (warn-only). async function senderContext(): Promise<{ key: string; agent: GlobalPidRecord | null; via: SenderVia; }> { const { agent, via } = await resolveSenderVia(); return { key: agent ? `agent:${agent.pid}` : "human", agent, via }; } /** * `ay whoami` — the calling agent's own canonical registration, resolved from * AGENT_YES_PID (see resolveSender). One command answers "which agent am I, * per the registry?": after a fleet restore, several agents can share a cwd * and a resumed conversation can believe it is a different lane than the * process actually registered as — the registry record is the ground truth * every routing surface (console, ay send, heartbeats) actually uses. Also * prints the traceable reply address so an agent can stamp outgoing messages * (` "...">`) without re-deriving its identity. */ async function cmdWhoami(rest: string[]): Promise { const y = yargs(rest) .usage("Usage: ay whoami [--json]") .option("json", { type: "boolean", default: false, description: "Machine-readable output" }) .help(false) .version(false) .exitProcess(false); const argv = await y.parseAsync(); const self = await resolveSender(); if (!self) { // Two distinct failures: no agent context at all (human shell), or a set // AGENT_YES_PID that resolves to nothing (stale env / record aged out). const reason = process.env.AGENT_YES_PID ? "unregistered" : "no-agent-context"; if (argv.json) { process.stdout.write(JSON.stringify({ agent: null, reason }) + "\n"); } else { process.stderr.write( reason === "unregistered" ? `ay whoami: AGENT_YES_PID=${process.env.AGENT_YES_PID} is set but matches no record in the registry (stale env, or the record aged out)\n` : `ay whoami: not inside an agent-yes session — AGENT_YES_PID is unset (human shell)\n`, ); } return 1; } const { state, question } = await deriveLiveState(self); const replyKw = self.agent_id ?? String(self.pid); const reply = `ay send ${replyKw}`; if (argv.json) { process.stdout.write(JSON.stringify({ ...self, state, question, reply }, null, 2) + "\n"); return 0; } const ageMin = Math.max(0, Math.round((Date.now() - self.started_at) / 60_000)); const lines = [ `agent ${self.cli} #${self.pid}${self.agent_id ? ` (agent_id ${self.agent_id})` : ""}`, `identity ${formatIdentity({ cwd: self.cwd, pid: self.pid })}`, `title ${self.title ?? "-"}`, `state ${state}${question ? ` — ${question}` : ""}`, `cwd ${self.cwd}`, `started ${new Date(self.started_at).toISOString()} (${ageMin}m ago)`, `wrapper ${self.wrapper_pid ?? "-"} parent ${self.parent_pid ?? "- (top-level)"}`, `log ${self.log_file ?? "-"}`, `fifo ${self.fifo_file ?? "-"}`, `reply ${reply} "..."`, `envelope …`, ]; process.stdout.write(lines.join("\n") + "\n"); return 0; } /** * Read the per-cwd TS PidStore JSONL and convert to the global record shape, * so pre-existing TS agents that were spawned before the global-index mirror * shipped still show up in `ay ls`. Merging is done in `mergeRecords`. */ async function readLocalTsPids(cwd: string): Promise { const jsonlPath = path.join(cwd, ".agent-yes", "pid-records.jsonl"); let raw: string; try { raw = await readFile(jsonlPath, "utf-8"); } catch { return []; } // Same merge semantics as ts/JsonlStore.ts: last line per _id wins, // tombstones (`$$deleted`) drop the entry. const docs = new Map(); for (const line of raw.split("\n")) { const trimmed = line.trim(); if (!trimmed) continue; try { const doc = JSON.parse(trimmed); if (!doc._id) continue; if (doc.$$deleted) { docs.delete(doc._id); continue; } const prev = docs.get(doc._id); docs.set(doc._id, prev ? { ...prev, ...doc } : doc); } catch { // skip corrupt } } return Array.from(docs.values()).map((d) => ({ pid: d.pid, cli: d.cli, prompt: d.prompt ?? null, cwd: d.cwd, log_file: d.logFile ?? null, fifo_file: d.fifoFile ?? null, status: d.status ?? "active", exit_code: d.exitCode ?? null, exit_reason: d.exitReason ?? null, started_at: d.startedAt ?? 0, title: d.title ?? null, })); } /** Merge by pid; later entries (typically from the global file) win. */ function mergeRecords(...buckets: GlobalPidRecord[][]): GlobalPidRecord[] { const out = new Map(); for (const bucket of buckets) { for (const r of bucket) { const prev = out.get(r.pid); out.set(r.pid, prev ? { ...prev, ...r } : r); } } return Array.from(out.values()); } // Subcommands EVERY *-yes binary accepts — inspection/messaging over the shared // agent registry (`cy ls`, `cy send`, `cy tail`, …). // MIRRORED in rs/src/cli.rs `SUBCOMMANDS` — the Rust runner delegates these to // this JS layer; keep the two lists in sync. const SUBCOMMANDS = new Set([ "ls", "list", "ps", "status", "whoami", "result", "notify", "notifyd", "read", "cat", "tail", "head", "hist", "history", "send", "msgs", "key", "select", "spawn", "attach", "stop", "exit", "restart", "note", "todo", "ask", "answer", "ch", "channels", "term", "widget", "mint", "serve", "tray", "schedule", "remote", "expose", "callback", "reap", "gc", "dsh-legacy", "help", ]); // Subcommands reserved for the GENERIC manager (`ay` / `agent-yes`). A cli-bound // alias like `cy` (= claude-yes = "agent-yes claude") must NOT treat these as // subcommands — `cy setup …` should run claude with that text, not manage the // host. Kept separate from SUBCOMMANDS so a runner alias falls straight through. const MANAGER_SUBCOMMANDS = new Set(["setup", "ws"]); const IDLE_THRESHOLD_MS = 60 * 1000; // `stuck`: alive + the screen still shows a busy marker (config `working`) yet the // log has been silent this long — i.e. wedged mid-stream (a silent API stream // stall), not finished. Deliberately MUCH longer than IDLE_THRESHOLD_MS: a slow // tool call (tests, install) is also "busy + quiet", so only a prolonged silence // is reported as stuck. Detection only — never auto-acts. Override via env. const STUCK_THRESHOLD_MS = (() => { const n = Number(process.env.AGENT_YES_STUCK_MS); return Number.isFinite(n) && n > 0 ? n : 5 * 60 * 1000; })(); // `ay send` submit-confirm tuning. A long/multi-line body pasted via bracketed // paste can take longer than any fixed delay to finish rendering — sending the // trailing Enter before that settles gets swallowed by the CLI's paste handling // (it lands mid-paste instead of submitting). So instead of a blind fixed sleep, // we poll the log for actual quiet, then confirm the Enter landed by watching for // either a `working` busy marker or a meaningful size bump, retrying if not. const SEND_SETTLE_QUIET_MS = 150; // no log growth for this long → paste finished rendering const SEND_SETTLE_MAX_MS = 1500; // cap: don't wait forever on a screen that's busy for other reasons const SEND_CONFIRM_QUIET_MS = 400; // after Enter, no growth for this long → response has settled const SEND_CONFIRM_MAX_MS = 1200; // cap per confirm attempt const SEND_CONFIRM_MIN_GROWTH_BYTES = 8; // filters out pure cursor-blink/frame noise const SEND_SUBMIT_MAX_RETRIES = 2; // total attempts = 1 + this // `ay send` typing-backoff: if the user is actively typing at the target's // terminal, injecting our body mid-line would fuse into their text and submit a // mangled line. Poll until they pause (activity older than TYPING_WINDOW_MS) or // we give up, then send anyway with a warning rather than dropping the message. const SEND_TYPING_POLL_MS = 200; const SEND_TYPING_MAX_WAIT_MS = 10_000; // `ay send` body length cap. Longer text isn't a prompt to type at a live CLI's // stdin — it's a document, and pasting it mid-session fuses with / truncates the // agent's terminal. Reject it outright and point at the stdin/file path so the // full text survives instead of being silently mangled by a bracketed paste. // Lowered 4096 → 1024 (operator 2026-08-20): past ~1KB the bracketed paste keeps // re-rendering long enough that the trailing Enter lands mid-paste and is // swallowed, so the body sits unsent in the target's input line. export const SEND_BODY_MAX_CHARS = 1024; /** * The cap, measured on what is actually WRITTEN to the agent's stdin. * * `ay send` checked the BODY against the cap and then transmitted the body plus * an `` envelope — a header naming the sender's identity and a closing * tag, together ~160 characters on real traffic. So the number the sender was * told was safe was never the number that went down the pipe, and a body * accepted at 1000 arrived as ~1160. The envelope is not a fixed surcharge either: it carries the * sender's cwd, branch and pid, so its length differs per sender and no constant * body budget can be quoted. Hence a check on the sum rather than a smaller * hard-coded limit. * * Returns the error text, or null when the send fits. `envelopeLen` is 0 for a * `--raw` send and for a slash command, both of which go out unwrapped — the * budget is then the whole cap, which is why this is computed per send and not * once. * * Paste framing (ts/bracketedPaste.ts) is deliberately NOT counted: those 12 * bytes are consumed by the receiving terminal as delimiters, never inserted as * text, so charging the sender for them would be charging for something that * does not arrive. */ /** * How many chars the `` envelope will add for THIS sender, or 0 when * none is added. * * Exists so the cap can be enforced before anything about the TARGET is touched. * The envelope is composed from the sender's own identity — cli, cwd, pid, * agent_id — and from the body; it never reads the target's record. So its size * is knowable with no I/O about the recipient, which is what lets the whole * caller-error check run ahead of the reachability probe. * * Deliberately builds the same shape `cmdSend` builds rather than estimating a * constant: the length varies per sender (a deep worktree and a long branch name * measure ~370 chars against ~160 for a short one), so a fixed number would be * wrong for somebody. Kept adjacent to the real construction so the two are * edited together — a divergence would under-count and let an over-cap payload * through the early check, where the authoritative check still catches it. */ /** * How the envelope names the strength of its own attribution. * * The receiving MODEL reads the body, not the mailbox file, so provenance that * exists only in the JSONL cannot inform the decision it is meant to inform. * `senderLabel` says this to a human reading `ay msgs`; this says it to the * agent being asked to act. * * Shared by the real construction and by `envelopeCostFor` so the two cannot * drift — a marker in one and not the other would under-count the cap. */ /** * Build the `` envelope for one send. * * ONE builder, because there are now three callers — `cmdSend`, the remote send * path, and `envelopeCostFor` (which must size exactly what is transmitted). * Two of them were already hand-copied templates whose own comment warned that a * divergence would under-count the cap; the third was missing entirely, which is * the bug this exists to close. * * `remote` names the transport when the message crosses a host boundary. The * reply route has to change with it: the sender's agent id resolves on the * SENDER's host, so a bare `ay send ` on the receiving side would address * nothing, or — worse — a local agent that happens to share the prefix. The * receiver is told what it actually needs: the sender's stable id, the host to * route to, and that its own alias for that host goes in front. No capability * token appears here; the receiver's route back is its own to hold. */ export function buildEnvelope(opts: { nonce: string; cli: string; identity: string; replyTarget: string | number; via: SenderVia; /** Set when the message crossed a host boundary; the sender's own label for * the far side is deliberately NOT used — it is meaningless to the receiver. */ remote?: { senderHost: string } | null; }): { prefix: string; suffix: string } { const attrib = envelopeAttribution(opts.via); const via = opts.remote ? " via remote" : ""; const reply = opts.remote ? `reply: ay send :${opts.replyTarget} "..."` : `reply: ay send ${opts.replyTarget} "..."`; return { prefix: `\n`, suffix: `\n`, }; } export function envelopeAttribution(via: SenderVia): string { switch (via) { case "ancestry": // Derived from the process tree because the wrapper's env was absent. return " via process-tree"; case "env-uncorroborated": // The claim and the process tree name DIFFERENT lanes. Loud, because this // is the case a receiver acting on sender weight must not be handed as a // fact — and loud only here, so it keeps meaning something. return " UNCORROBORATED-SENDER"; case "env-unverified": // Nothing to disagree with. Silent in the envelope: the receiver gains // nothing actionable from "we could not check", and a marker on honest // traffic is how the loud one stops being read. It is still recorded in // from_via for anyone who wants to weigh it. return ""; default: return ""; } } export async function envelopeCostFor(body: string, raw: boolean): Promise { if (raw || isSlashCommand(body)) return 0; const sender = await senderContext(); if (!sender.agent) return 0; const nonce = "00000000"; // 4 random bytes as hex — fixed width, so any value sizes alike const identity = formatIdentity({ cwd: sender.agent.cwd, pid: sender.agent.pid }); const replyTarget = sender.agent.agent_id || sender.agent.pid; const { prefix, suffix } = buildEnvelope({ nonce, cli: sender.agent.cli, identity, replyTarget, via: sender.via, }); return prefix.length + suffix.length; } export function sendPayloadCapError(bodyLen: number, envelopeLen: number): string | null { const transmitted = bodyLen + envelopeLen; if (transmitted <= SEND_BODY_MAX_CHARS) return null; const budget = SEND_BODY_MAX_CHARS - envelopeLen; const remedy = `Write it to a file and send the PATH instead, e.g. ` + `'ay send "details: /path/to/notes.md"'`; // Nothing bounds a cwd's depth or a branch name's length, so an envelope can // in principle reach the cap on its own. Then NO body length works, and // quoting a budget would print a negative number and tell the sender to // shorten to ≤0 — an error naming a remedy that cannot work, which is the // very defect this function exists to remove. Say what is actually true. if (budget <= 0) { return ( `the envelope alone is ${envelopeLen} chars, at or over the ` + `${SEND_BODY_MAX_CHARS}-char limit, so no body length can fit. ${remedy} — ` + `and send it with --raw, which omits the envelope.` ); } return ( `message would transmit ${transmitted} chars, over the ${SEND_BODY_MAX_CHARS}-char limit` + (envelopeLen > 0 ? ` — ${bodyLen} of body plus ${envelopeLen} of envelope, which ` + `ay send adds for you. Your budget for this send is ${budget} chars of body` : "") + `. Longer text isn't a terminal prompt. Piping it in with '-' hits this same ` + `cap — the payload is capped however it arrives. ${remedy}, ` + `or shorten the body to ≤${budget} chars.` ); } /** * `ay send` exit status when the target cannot be written to at all — its stdin * FIFO takes no writer (ENXIO / gone), or it never registered one. Distinct from * 1 (a transport hiccup: a backed-up reader, a lock we could not take, a body * over the cap) because the remedies are different: 1 says try again, this says * the row is dead and needs `ay restart`. `ay ls` shows the same rows as * `unreachable`, so a caller can see it BEFORE it hands out work. * 2 is already "timed out" everywhere else in this CLI, so this is 3. */ export const SEND_EXIT_UNREACHABLE = 3; /** * Whether an errno from a FIFO write means "nobody is on the other end", as * opposed to "the other end is slow" or "the filesystem said no". * * Three, and the third is the one a preflight probe cannot see: * ENXIO — open() found a FIFO with no reader (the probe's case) * ENOENT — the FIFO is gone * EPIPE — the reader was there at open() and vanished DURING the write. Only * this path can produce it, which is why it is not in the probe. * * EAGAIN/EWOULDBLOCK are deliberately absent: a full pipe means a reader exists * and is slow, which is `ay send`'s retry loop, not a dead row. */ export function isUnreachableWriteErrno(code: string | undefined): boolean { return code === "ENXIO" || code === "ENOENT" || code === "EPIPE"; } /** * Whether `name` is a subcommand. `managerCommands` (default true, for the * generic `ay`/`agent-yes` entry) additionally admits manager-only commands * like `setup`; pass false for a cli-bound alias (cy/claude-yes/…) so those * names fall through to running the agent instead. */ export function isSubcommand(name: string | undefined, managerCommands = true): boolean { if (!name) return false; return SUBCOMMANDS.has(name) || (managerCommands && MANAGER_SUBCOMMANDS.has(name)); } /** * Footgun guard for the MANAGER entry (`ay`/`agent-yes`): true when the first arg * is a bare word (not a flag) that is neither a subcommand nor a known CLI — a * typo, or a newer subcommand run on an older build. The caller should error * rather than silently spawn an agent with the word as a prompt (operator 2026-07-26: * bare `ay ` is too dangerous; a spawn must name a CLI — `ay …` or * `--cli`). Never fires for cli-bound aliases (cy/claude-yes/…), where the first * word is legitimately the prompt, nor for flags or an empty invocation. */ export function isUnknownManagerToken( rawArg: string | undefined, managerCommands: boolean, supportedClis: readonly string[], ): boolean { if (!managerCommands || !rawArg || rawArg.startsWith("-")) return false; if (isSubcommand(rawArg, managerCommands)) return false; return !supportedClis.includes(rawArg); } /** * True for a completely bare MANAGER invocation — `ay` / `agent-yes` with no * args at all. `ay` means agent-yes (the fleet manager), so it prints help * instead of silently launching an agent; naming a CLI is what launches one * (`ay claude`), and `cy` stays the zero-argument way to start claude. * * Never fires for a cli-bound alias (managerCommands=false): bare `cy` / * `claude-yes` / `codex-yes` must keep spawning their agent. * * `argv` is process.argv, so length 2 = [runtime, script] with no user args; * a flags-only run like `ay --continue` still launches (it asked for a run). */ export function isBareManagerInvocation(argv: string[], managerCommands: boolean): boolean { return managerCommands && argv.length <= 2; } /** * Write to stdout and wait until it has actually been handed off. * * The CLI ends with `process.exit()`, which DISCARDS bytes still sitting in the * pipe buffer. Writing to a file completes synchronously so this never showed * there, but any consumer that PIPES us silently lost everything past 64KiB. * Measured on `ay ls --json` with a large fleet (symval CTO, 2026-08-05): * * ay ls --json > file 118055 bytes, valid JSON * ay ls --json | consumer 65536 bytes, cut mid-multibyte — will not parse * * The truncation is invisible to the caller: it looks exactly like a small * fleet, and an orchestrator that parses this reported 27 of 71 agents. * * Only capturing the REAL write's completion works. Probing afterwards with an * empty `write("")` — by return value, by drain event, or by callback — reports * ready while the big write is still queued (all three verified failing), so do * not "simplify" this into a flush helper at the exit site. */ async function writeStdoutFlushed(text: string): Promise { // Use the WRITE'S OWN RETURN VALUE, not a probe afterwards. `false` means the // pipe buffer is full and bytes are still queued; only then do we wait. // // Deliberately `=== false`: a stubbed/mocked stdout (tests, embedders) returns // undefined, and treating that as backpressure made this await forever — it // broke 5 specs before the strict compare went in. The timeout is a second // guarantee that a stuck consumer can never hang the CLI. const flushed = process.stdout.write(text); if (flushed === false) { await Promise.race([ new Promise((resolve) => process.stdout.once("drain", () => resolve())), new Promise((resolve) => setTimeout(resolve, 2000)), ]); } } /** * Top-level entry. Returns the desired process exit code, or null if argv * is not a subcommand invocation. */ export async function runSubcommand(argv: string[]): Promise { const sub = argv[2]; // Manager-only subcommands (setup) aren't subcommands for a cli-bound alias // like `cy` — they fall through to running the agent. Computed once from argv // so it holds regardless of caller, and reused to hide manager-only help. const managerCommands = !invokedCliName(argv); if (!isSubcommand(sub, managerCommands)) return null; const rest = argv.slice(3); try { switch (sub) { case "ls": case "list": return await cmdLs(rest); // `ps` was an alias for `ls`; it is now the RESOURCE view — same agents, // rolled up per process tree with box vitals. See ts/cmdPs.ts. case "ps": return await (await import("./cmdPs.ts")).cmdPs(rest); case "status": return await cmdStatus(rest); case "whoami": return await cmdWhoami(rest); case "result": return await cmdResult(rest); case "notify": return await cmdNotify(rest); case "notifyd": return await cmdNotifyd(rest); case "read": case "cat": return await cmdRead(rest, { mode: "cat" }); case "tail": return await cmdRead(rest, { mode: "tail" }); case "head": return await cmdRead(rest, { mode: "head" }); case "hist": case "history": return await (await import("./hist.ts")).cmdHist(rest); case "send": return await cmdSend(rest); case "msgs": return await cmdMsgs(rest); case "key": return await cmdKey(rest); case "select": return await cmdSelect(rest); case "spawn": return await cmdSpawn(rest); case "attach": return await cmdAttach(rest); case "stop": return await cmdStop(rest); case "exit": return await cmdExit(rest); case "restart": return await cmdRestart(rest); case "note": return await cmdNote(rest); case "ask": case "answer": { // `ay ask` needs `ay send`'s delivery path and this file's agent // resolver, both of which live here — so they are handed over rather // than imported, which would make askCli.ts and this module circular. const { runAskSubcommand, runAnswerSubcommand } = await import("./askCli.ts"); const deps = { // `all: true` — a question may legitimately be addressed to an agent // that has since gone idle or exited (that is precisely the case // worth recording), so the answerer's liveness is REPORTED rather // than made a precondition for asking. resolveAgent: (keyword: string) => resolveOne(keyword, { all: true, active: false, json: false, latest: false, cwdScope: null, }), send: cmdSend, }; return sub === "ask" ? await runAskSubcommand(rest, deps) : await runAnswerSubcommand(rest, deps); } case "todo": { const { runTodoSubcommand } = await import("./todoCli.ts"); return runTodoSubcommand(rest); } case "ch": case "channels": { const { cmdCh } = await import("./channels.ts"); return cmdCh(rest); } case "term": { const { cmdTerm } = await import("./terminal.ts"); return cmdTerm(rest); } case "widget": { const { cmdWidget } = await import("./widget.ts"); return cmdWidget(rest); } case "mint": { const { cmdMint } = await import("./widget.ts"); return cmdMint(rest); } case "serve": { const { cmdServe } = await import("./serve.ts"); return cmdServe(rest); } case "tray": { const { cmdTray } = await import("./trayApp.ts"); return cmdTray(rest); } case "setup": { const { cmdSetup } = await import("./setup.ts"); return cmdSetup(rest); } case "ws": { const { cmdWs } = await import("./ws.ts"); return cmdWs(rest); } case "schedule": { const { cmdSchedule } = await import("./schedule.ts"); return cmdSchedule(rest); } case "remote": { const { cmdRemote } = await import("./remotes.ts"); return cmdRemote(rest); } case "expose": { const { cmdExpose } = await import("./expose.ts"); return cmdExpose(rest); } case "callback": { const { cmdCallback } = await import("./callback.ts"); return cmdCallback(rest); } case "reap": { const reaper = await import("./reaper.ts"); await reaper.sweep(); return 0; } case "dsh-legacy": { const { cmdDsh } = await import("./cmdDsh.ts"); return await cmdDsh(rest); } case "gc": { const { gcOldBinaryDirs } = await import("./rustBinary.ts"); const { gcLogs } = await import("./globalPidIndex.ts"); const bins = gcOldBinaryDirs(); // Logs are the bigger leak of the two: binary dirs are ~20-30 MiB per // release, while a single long-lived session's raw log can pass 500. const logs = await gcLogs(); if (bins.removed.length === 0) { process.stdout.write("no old agent-yes binary cache dirs to remove\n"); } else { for (const v of bins.removed) process.stdout.write(`removed ${v}\n`); const mib = (bins.freedBytes / 1024 / 1024).toFixed(1); process.stdout.write( `freed ${mib} MiB (${bins.freedBytes} bytes) across ${bins.removed.length} version dir(s)\n`, ); } if (logs.removed.length === 0) { process.stdout.write("no stale session logs to remove\n"); } else { const mib = (logs.freedBytes / 1024 / 1024).toFixed(1); process.stdout.write( `freed ${mib} MiB (${logs.freedBytes} bytes) across ${logs.removed.length} session log file(s)\n`, ); } return 0; } case "help": return cmdHelp(managerCommands); default: return null; } } catch (err) { const msg = err instanceof Error ? err.message : String(err); process.stderr.write(`ay ${sub}: ${msg}\n`); return 1; } } // --------------------------------------------------------------------------- // ay help // --------------------------------------------------------------------------- /** * The banner shown by `ay help` / `ay -h` when this process is itself running * inside an agent (`AGENT_YES_PID` set — see resolveSender). Answers the three * things a nested agent actually needs: who am I, who spawned me, and how do I * drive sub-agents of my own — so it doesn't have to rediscover the fan-out * primitives (spawn / ay ls forest / ay ls --watch) from scratch every session. */ async function buildAgentContextSection(self: GlobalPidRecord): Promise { const hasParentPid = typeof self.parent_pid === "number" && self.parent_pid > 0; const parent = hasParentPid ? ( await listRecords(undefined, { all: true, active: false, json: false, latest: false, cwdScope: null, }) ).find((r) => r.wrapper_pid === self.parent_pid) : undefined; const whoAmI = `You are agent pid ${self.pid} (${self.cli}) in ${shortenPath(self.cwd)}.`; // Three distinct states: no parent at all (top-level); a parent_pid whose // record we can resolve; or a parent_pid we can't resolve (its record aged // out / lives on a remote) — that last case is still nested, just unknown, // so it must not collapse into the "top-level" line. const parentLine = !hasParentPid ? `Top-level agent — no parent (started from a human shell or scheduler).` : parent ? `Spawned by agent pid ${parent.pid} (${parent.cli}) in ${shortenPath(parent.cwd)}.` : `Nested under a parent (wrapper pid ${self.parent_pid}) whose record isn't in the local registry.`; // The reporting duty is stated in the `` block wrapping this // agent's initial prompt, but that block scrolls out of a long session (or is // compacted away) long before the agent finishes. `ay help` is where an agent // goes when it has lost the thread, so restate the obligation here. const dutyLine = parent ? ` You owe it a report: \`ay send ${parent.agent_id || parent.pid} "..."\` when you finish, and\n` + ` when you are blocked. It is not watching your terminal.\n` : ``; return ( `You are running inside an agent:\n` + ` ${whoAmI}\n` + ` ${parentLine}\n` + dutyLine + `\n` + `As an agent, you can:\n` + ` Spawn a sub-agent:\n` + ` ay -- "" auto-links as your child\n` + ` ay claude --model sonnet --advisor opus -- "" routine task\n` + ` ay claude --model opus --advisor fable -- "" complex task\n` + ` (pick --model by task complexity so easy tasks don't cost like hard ones;\n` + ` --advisor is a claude-cli flag — only takes effect for claude/cy)\n` + ` List agents (your children nest under your own pid in the tree):\n` + ` ay ls --cwd ${shortenPath(self.cwd)}\n` + ` Get notified when a sub-agent finishes / goes idle / crashes (preferred):\n` + ` ay notify watch --unread\n` + ` (one watch loop for your whole fan-out: needs_input / idle / exited edges land\n` + ` in your inbox; a hard child crash is caught by the 2s liveness poll, which\n` + ` nothing push-based can see)\n` + ` Watch agent state changes, scoped to your workspace:\n` + ` ay ls --watch --cwd ${shortenPath(self.cwd)}\n` + ` (NDJSON stream of state changes across every matched agent — one watcher\n` + ` for the whole fan-out instead of N \`ay status --watch\`es)\n` + ` Read one sub-agent's output:\n` + ` ay tail -f follow live output (no single command tails\n` + ` many agents' content at once yet — loop\n` + ` \`ay ls --json\` pids into per-pid \`ay tail\`)\n` + `\n` ); } export async function cmdHelp(managerCommands = true): Promise { // `setup` is manager-only — hide it when invoked through a cli-bound alias // (cy/claude-yes/…), where `cy setup` runs the agent instead of managing the host. const setupLine = managerCommands ? ` ay setup guided setup: pick a workspace, share to agent-yes.com\n` : ``; // `ws` is manager-only for the same reason as `setup`. const wsLines = managerCommands ? ` ay ws ls [--status] list //tree/ workspaces\n` + ` ay ws new /[@branch] clone/refresh a workspace (ay ws help for more)\n` : ``; // Only agents carry AGENT_YES_PID — a human shell never sets it — so this // section is skipped entirely (no async work at all) for interactive use. const self = process.env.AGENT_YES_PID ? await resolveSender() : null; const agentSection = self ? await buildAgentContextSection(self) : ""; process.stdout.write( agentSection + `ay - agent-yes CLI\n` + `\n` + `Management:\n` + ` ay ls [keyword] list running agents\n` + ` ay ps [keyword] per-agent CPU/RSS, rolled up over each\n` + ` agent's whole process tree, + box vitals\n` + ` ay tail [-f] [-n N] last N lines (96), -f to follow\n` + ` ay read [page opts] paginate: --last/--head N, --range A:B,\n` + ` --before-line L [--limit N]\n` + ` ay cat full log\n` + ` ay head first N lines\n` + ` ay hist [-n 6] [--all] [--json] past agent conversations (claude/codex\n` + ` transcripts, incl. exited sessions);\n` + ` this cwd unless --all\n` + ` ay send send a message (keyword '.' = agent in this cwd)\n` + ` ay msgs [keyword] [--in|--out] inter-agent message log (sent + received)\n` + ` ay ch mk|join|send|read|tail local-first E2E channels: AI ↔ humans on a topic (ay ch help)\n` + ` ay term embed