/** * MemPalace ↔ pi bridge. * * Registers every MemPalace MCP tool as a pi tool that proxies to `tools/call`. * Two interchangeable transports, selected at load time: * * - LOCAL (default): spawn the `mempalace-mcp` stdio server as a subprocess * (StdioMcpClient) — hardened with per-request timeouts + generation-tracked * respawn/self-heal for slow cold-opens (see below). * - EXTERNAL: connect to a shared MemPalace over HTTP (RemoteMcpClient) when * $MEMPALACE_REMOTE_URL is set — e.g. one palace serving pi + opencode + * native. Optional bearer auth via $MEMPALACE_REMOTE_TOKEN. No local * `mempalace-mcp` process is spawned in this mode. * * Either way the client performs the MCP `initialize` handshake, lists tools, * and the extension body below is transport-agnostic (depends only on the * IMcpClient interface). * * Lifecycle automation (per ~/.agents/skills/mempalace/SKILL.md): * - Wake-up (auto): on first user prompt of a fresh session, inject * `mempalace_status` + `mempalace_diary_read` output as context so the * agent orients itself the way the mempalace skill describes. Skipped * on resume/fork (palace context is already in the thread). * - Feeding (auto): stage + mine this container's pi transcripts into the * palace on `session_shutdown` and on a debounced `agent_settled`. Needs * no LLM turn (pi transcripts are JSONL on disk), which is why it CAN be * automatic where the diary cannot. The file-side work is delegated to * `mempalace-pi-session --prepare` (export + threshold + staging, plus the * rsync to the palace host when the palace is remote); the mine itself * must run through THIS client, because the palace is single-writer and * this process is the holder — a CLI `mempalace mine` during a live * session dies with "palace ... is held by PID ". Going through the * client also means it automatically targets whichever palace this bridge * is pointed at (local stdio or a shared remote one). * - MEMPALACE_FEED=0 disable feeding entirely * - MEMPALACE_FEED_BIN helper to run (default mempalace-pi-session) * - MEMPALACE_FEED_WING target wing (default wing_conversations) * - MEMPALACE_FEED_DEBOUNCE_MS min gap between mid-session feeds (default 600000) * - Mailbox (auto): read the RFC 003 coordination log and surface the asks * this device actually OWES a reply to — at wake-up, and again on a * debounced `agent_settled` so events arriving mid-session are seen this * session rather than the next one. Owed-ness is DERIVED (see deriveOwed); * it is not the `status="open"` filter, because `status` is written once * into an append-only log and an answered ask keeps matching forever. * Inert unless BOTH $MEMPALACE_PI_DEVICE and $MEMPALACE_REMOTE_URL are set: * on one shared palace `to_agent` carries the entire distinction between * machines, so an unstamped client has no address to be reached at. * - MEMPALACE_MAILBOX=0 disable mailbox reads entirely * - MEMPALACE_MAILBOX_POLL_MS min gap between mid-session polls (default 300000) * - MEMPALACE_MAILBOX_RESURFACE_MS re-announce a still-owed ask after (default 3600000) * - Wind-down (manual): `/mempalace-diary` command prompts the LLM to * write an AAAK-formatted diary entry. Not fully auto because pi * sessions are typically short/tactical and session_shutdown is too * late to drive an LLM turn. * * Identity: `agent_name` for diary calls comes from $MEMPALACE_AGENT_NAME, * defaulting to "pi". First diary write creates `wing_pi`. * * Fail-soft: if the MCP subprocess can't start, pi keeps working without * palace tools (warning on stderr only). * * Stall protection: every JSON-RPC request carries a timeout. If * `mempalace-mcp` wedges (e.g. OrbStack virtiofs cold-open of a large * chroma.sqlite3 / HNSW load), the awaiting promise would otherwise hang * forever and freeze the pi TUI (ESC cancels the LLM stream, not a pending * tool execution). On timeout we reject the request AND kill the wedged * child, so pi gets an error instead of hanging and later calls fail fast. * This is a per-REQUEST timeout, not a process-lifetime one — the * long-lived server is only killed when a request genuinely stalls. * - MEMPALACE_MCP_TIMEOUT_MS tool-call/request timeout (default 60000) * - MEMPALACE_MCP_INIT_TIMEOUT_MS initialize+tools/list timeout (default 300000) * Set either to 0 to disable (legacy unbounded behavior). * * Self-heal (respawn): a stall-kill (or any crash) is no longer a permanent * latch. The next tool call transparently respawns `mempalace-mcp` and * retries, with capped exponential backoff so a persistently-broken server * can't hot-loop. The respawn budget resets on ANY successful JSON-RPC * response (proof the server is actually live), so a recovered server * regains full patience. The two timeouts above are deliberately split: * the long INIT timeout lets a genuine first cold-open finish without being * killed (the original incident), while the short per-call timeout still * aggressively kills a stuck query. After a server has opened the palace * once, the OS page cache is warm, so respawn cold-opens are fast — which * is exactly why a generous INIT timeout + bounded respawn compose well * instead of overlapping. * - MEMPALACE_MCP_MAX_RESPAWNS respawn attempts before giving up (default 2; 0 disables self-heal) * - MEMPALACE_MCP_RESPAWN_BACKOFF_MS base backoff, doubled per attempt (default 1000) * * Debug: set MEMPALACE_EXT_DEBUG=1 to surface mempalace-mcp stderr. */ import { type ChildProcessWithoutNullStreams, spawn } from "node:child_process"; import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { Type } from "typebox"; // Minimal MCP stdio JSON-RPC client. MCP uses newline-delimited JSON. interface McpTool { name: string; description?: string; inputSchema?: unknown; } interface Pending { resolve: (v: any) => void; reject: (e: Error) => void; timer: ReturnType | null; } // Transport-agnostic client contract. The extension body depends only on this, // so LOCAL (StdioMcpClient) and EXTERNAL (RemoteMcpClient) are drop-in swaps. // Each implementation owns its own reliability model — stdio uses timeout + // kill + respawn; http uses fetch + AbortController timeout (+ 404 re-init if // a future streamable-HTTP server issues sessions). interface IMcpClient { tools: McpTool[]; readonly alive: boolean; onExit: (() => void) | null; start(): Promise; callTool(name: string, args: Record): Promise; ensureAlive(): Promise; stop(): void | Promise; } const num = (envVal: string | undefined, fallback: number): number => { const n = envVal !== undefined ? Number(envVal) : Number.NaN; return Number.isFinite(n) && n >= 0 ? n : fallback; }; // One event as returned by `mempalace_event_list`. Every field is optional on // purpose: this is parsed from a tool's JSON text, so it is untrusted input and // a missing/renamed field must degrade to "not owed" rather than throw. type LogEvent = { id?: string; seq?: number; hlc?: string; type?: string; status?: string; from_agent?: string; to_agent?: string; correlation_id?: string | null; created_at?: string; body?: string; metadata?: Record | null; }; const sleep = (ms: number): Promise => new Promise((res) => { const t = setTimeout(res, ms); if (typeof t.unref === "function") t.unref(); }); class StdioMcpClient implements IMcpClient { private proc: ChildProcessWithoutNullStreams | null = null; private nextId = 1; private pending = new Map(); private stdoutBuf = ""; private ready: Promise | null = null; public tools: McpTool[] = []; // Per-request timeouts (ms). 0 = disabled (unbounded, legacy behavior). private requestTimeoutMs = num(process.env.MEMPALACE_MCP_TIMEOUT_MS, 60_000); private initTimeoutMs = num(process.env.MEMPALACE_MCP_INIT_TIMEOUT_MS, 300_000); // Self-heal (respawn) controls. private maxRespawns = num(process.env.MEMPALACE_MCP_MAX_RESPAWNS, 2); private respawnBackoffMs = num(process.env.MEMPALACE_MCP_RESPAWN_BACKOFF_MS, 1_000); private respawns = 0; // consecutive respawn attempts; reset on any success private reviving: Promise | null = null; // Liveness: true only between a completed init and the matching death. // Inferring from `proc` is racy (non-null during the SIGTERM→exit window). private healthy = false; // Spawn generation. Each (re)spawn bumps it; death handlers carry the gen // they were attached for and no-op if a newer server has since taken over // (prevents a stale OLD-proc 'exit' from clobbering a freshly respawned one). private gen = 0; // Spawn args remembered so a respawn can reuse them (set via constructor). private command: string; private args: string[]; // Fired when the child process dies (exit or stall-kill). Lets the // extension flip `available` so later tool calls fail fast. public onExit: (() => void) | null = null; constructor(command = "mempalace-mcp", args: string[] = []) { this.command = command; this.args = args; } /** True only when a server has completed init and not since died. */ get alive(): boolean { return this.healthy; } async start(): Promise { if (this.ready) return this.ready; this.ready = (async () => { const myGen = ++this.gen; const child = spawn(this.command, this.args, { stdio: ["pipe", "pipe", "pipe"] }); this.proc = child; child.on("error", (err) => this.handleDeath(myGen, err)); child.on("exit", (code) => this.handleDeath(myGen, new Error(`mempalace-mcp exited (code=${code})`)), ); this.proc.stdout.setEncoding("utf8"); this.proc.stdout.on("data", (chunk: string) => this.onStdout(chunk)); // Drain stderr silently. Re-enable by setting MEMPALACE_EXT_DEBUG=1. this.proc.stderr.setEncoding("utf8"); if (process.env.MEMPALACE_EXT_DEBUG) { this.proc.stderr.on("data", (chunk: string) => { process.stderr.write(`[mempalace-mcp stderr] ${chunk}`); }); } else { this.proc.stderr.resume(); // drain without logging } // MCP initialize handshake. Cold-open over virtiofs can be slow, so // use the (longer) init timeout here rather than the per-call one. await this.request( "initialize", { protocolVersion: "2024-11-05", capabilities: {}, clientInfo: { name: "pi-mempalace-ext", version: "0.1.0" }, }, this.initTimeoutMs, ); this.notify("notifications/initialized", {}); const listed = await this.request("tools/list", {}, this.initTimeoutMs); this.tools = (listed?.tools as McpTool[]) ?? []; // Init complete and the process is still ours — mark live so ensureAlive() // reports true. Guard against a death that raced in during init. if (myGen === this.gen) this.healthy = true; })(); return this.ready; } private onStdout(chunk: string) { this.stdoutBuf += chunk; let nl: number; while ((nl = this.stdoutBuf.indexOf("\n")) !== -1) { const line = this.stdoutBuf.slice(0, nl).trim(); this.stdoutBuf = this.stdoutBuf.slice(nl + 1); if (!line) continue; let msg: any; try { msg = JSON.parse(line); } catch { continue; } if (typeof msg.id === "number" && this.pending.has(msg.id)) { const p = this.pending.get(msg.id)!; this.settle(msg.id); if (msg.error) p.reject(new Error(msg.error.message ?? "MCP error")); else { // A successful response proves the server is live — restore the // full respawn budget for any future (unrelated) stall. this.respawns = 0; p.resolve(msg.result); } } // notifications (no id) ignored for now. } } /** Remove a pending request and clear its timeout timer. */ private settle(id: number) { const p = this.pending.get(id); if (p?.timer) clearTimeout(p.timer); this.pending.delete(id); } private failAll(err: Error) { for (const id of [...this.pending.keys()]) { const p = this.pending.get(id)!; this.settle(id); p.reject(err); } } /** * Child died (exit, spawn error, or stall-kill): reject everything. * Ignored if a newer generation has already taken over — a late 'exit' * from a killed old process must not tear down a fresh respawn. */ private handleDeath(gen: number, err: Error) { if (gen !== this.gen) return; // stale handler from a superseded process this.proc = null; this.healthy = false; // Clear the memoized start() promise so a subsequent start()/ensureAlive() // spawns a fresh server instead of returning the dead one's promise. this.ready = null; this.failAll(err); try { this.onExit?.(); } catch {} } /** * Ensure a live, initialized server — respawning a dead one with capped * exponential backoff. Returns whether the server is alive afterwards. * Concurrent callers share a single in-flight revive. The attempt counter * is reset by `onStdout` on any successful response, so a server that comes * back and works regains its full budget; a server that keeps dying hits * `maxRespawns` and stays down (restart pi) rather than hot-looping. */ async ensureAlive(): Promise { if (this.alive) return true; if (this.reviving) return this.reviving; this.reviving = (async () => { while (!this.alive && this.respawns < this.maxRespawns) { this.respawns++; await sleep(this.respawnBackoffMs * 2 ** (this.respawns - 1)); // Force a fresh spawn: drop any settled (rejected/dead) start() // promise so start() doesn't short-circuit on a stale memo. Safe // because ensureAlive is serialized (single `reviving`) and only // runs while not healthy — no useful in-flight start() can exist. this.ready = null; try { await this.start(); } catch { // start() rejected (e.g. respawn cold-open also stalled and was // killed). Loop until the budget is exhausted. } } return this.alive; })(); try { return await this.reviving; } finally { this.reviving = null; } } private write(obj: unknown) { if (!this.proc) throw new Error("MCP process not started"); this.proc.stdin.write(`${JSON.stringify(obj)}\n`); } private notify(method: string, params: unknown) { this.write({ jsonrpc: "2.0", method, params }); } request(method: string, params: unknown, timeoutMs = this.requestTimeoutMs): Promise { const id = this.nextId++; return new Promise((resolve, reject) => { let timer: ReturnType | null = null; if (timeoutMs > 0) { timer = setTimeout(() => { if (!this.pending.has(id)) return; this.settle(id); reject( new Error( `mempalace-mcp request '${method}' timed out after ${timeoutMs}ms ` + `(server wedged — likely cold storage open). Terminating the ` + `stalled server; it will be respawned on the next call.`, ), ); // Kill the wedged child so subsequent calls fail fast instead of // stacking up behind a dead server. The 'exit' handler // (handleDeath) rejects any other pending requests. this.kill(); }, timeoutMs); // Don't let a pending MCP timer keep the event loop alive. if (typeof timer.unref === "function") timer.unref(); } this.pending.set(id, { resolve, reject, timer }); try { this.write({ jsonrpc: "2.0", id, method, params }); } catch (err) { this.settle(id); reject(err as Error); } }); } async callTool(name: string, args: Record): Promise { return this.request("tools/call", { name, arguments: args }); } /** SIGTERM then SIGKILL grace, for stall recovery. */ private kill() { const proc = this.proc; if (!proc) return; try { proc.kill("SIGTERM"); } catch {} const t = setTimeout(() => { try { proc.kill("SIGKILL"); } catch {} }, 2_000); if (typeof t.unref === "function") t.unref(); } stop() { if (this.proc) { try { this.proc.kill("SIGTERM"); } catch {} this.proc = null; } } } // ─────────────────────────────────────────────────────────────────────── // RemoteMcpClient — EXTERNAL transport (streamable-HTTP / sessionless JSON-RPC). // // VENDORED from pi-extensions/extensions/mcp-loader.ts (class RemoteMcpClient). // Kept as a copy rather than a shared import because the two repos ship // independently; the canonical implementation lives in mcp-loader.ts. Port // protocol fixes in both directions. // // MCP-STREAMABLE-HTTP-CLIENT-SYNC: v1 // This token must match the canonical block's. Bump it in BOTH files whenever // the streamable-HTTP protocol handling changes; scripts/check-mcp-client-sync.sh // fails when they diverge (skips gracefully if pi-extensions isn't checked out). // // Deltas from the canonical copy: // • protocolVersion pinned to "2024-11-05" to match mempalace-mcp's stdio // handshake (the server echoes whatever it is sent, so cosmetic today, // but keeps both transports consistent). // • per-request AbortController timeout honouring MEMPALACE_MCP_TIMEOUT_MS / // MEMPALACE_MCP_INIT_TIMEOUT_MS, mirroring StdioMcpClient's timeout ethos. // • alive / ensureAlive / onExit to satisfy IMcpClient. // // NOTE: mempalace-mcp --transport http is a SESSIONLESS, stateless JSON-RPC // server (no Mcp-Session-Id, always application/json, Connection: close), so // the session-id / SSE / 404-reinit branches below are never exercised against // it today. They are retained so this client also works unchanged if mempalace // later moves to a full streamable-HTTP (FastMCP) transport. // ─────────────────────────────────────────────────────────────────────── const REMOTE_PROTOCOL_VERSION = "2024-11-05"; const REMOTE_CLIENT_INFO = { name: "pi-mempalace-ext", version: "0.1.0" }; class RemoteMcpClient implements IMcpClient { private url: string; private extraHeaders: Record; private sessionId: string | null = null; private nextId = 1; public tools: McpTool[] = []; // Unused for HTTP (no child process to die) — present to satisfy IMcpClient. public onExit: (() => void) | null = null; private requestTimeoutMs = num(process.env.MEMPALACE_MCP_TIMEOUT_MS, 60_000); private initTimeoutMs = num(process.env.MEMPALACE_MCP_INIT_TIMEOUT_MS, 300_000); private healthy = false; private reviving: Promise | null = null; constructor(url: string, headers?: Record) { this.url = url; this.extraHeaders = headers ?? {}; } get alive(): boolean { return this.healthy; } async start(): Promise { await this.request( "initialize", { protocolVersion: REMOTE_PROTOCOL_VERSION, capabilities: {}, clientInfo: REMOTE_CLIENT_INFO }, { timeoutMs: this.initTimeoutMs }, ); await this.notify("notifications/initialized", {}); const listed = await this.request("tools/list", {}, { timeoutMs: this.initTimeoutMs }); this.tools = (listed?.tools as McpTool[]) ?? []; this.healthy = true; } async callTool(name: string, args: Record): Promise { return this.request("tools/call", { name, arguments: args }); } /** * Re-establish connectivity (and re-list tools) after a fetch failure. * For the stateless server this is just a fresh initialize+tools/list; * concurrent callers share one in-flight attempt. */ async ensureAlive(): Promise { if (this.healthy) return true; if (this.reviving) return this.reviving; this.reviving = (async () => { try { this.sessionId = null; await this.start(); } catch { // stays unhealthy } return this.healthy; })(); try { return await this.reviving; } finally { this.reviving = null; } } async stop(): Promise { this.healthy = false; if (!this.sessionId) return; try { await fetch(this.url, { method: "DELETE", headers: this.buildHeaders() }); } catch {} } private buildHeaders(): Record { const h: Record = { "Content-Type": "application/json", Accept: "application/json, text/event-stream", "MCP-Protocol-Version": REMOTE_PROTOCOL_VERSION, ...this.extraHeaders, }; if (this.sessionId) h["Mcp-Session-Id"] = this.sessionId; return h; } private signalFor(timeoutMs: number): AbortSignal | undefined { return timeoutMs > 0 ? AbortSignal.timeout(timeoutMs) : undefined; } private async reinitialize(): Promise { await this.request( "initialize", { protocolVersion: REMOTE_PROTOCOL_VERSION, capabilities: {}, clientInfo: REMOTE_CLIENT_INFO }, { timeoutMs: this.initTimeoutMs, allowReinitOn404: false }, ); await this.notify("notifications/initialized", {}); } private async request( method: string, params: unknown, opts: { timeoutMs?: number; allowReinitOn404?: boolean } = {}, ): Promise { const timeoutMs = opts.timeoutMs ?? this.requestTimeoutMs; const allowReinitOn404 = opts.allowReinitOn404 ?? true; const id = this.nextId++; const body = JSON.stringify({ jsonrpc: "2.0", id, method, params }); const sessionAtStart = this.sessionId; let res: Response; try { res = await fetch(this.url, { method: "POST", headers: this.buildHeaders(), body, signal: this.signalFor(timeoutMs), }); } catch (err) { // Network failure / timeout (AbortError) — mark unreachable so the next // tool call triggers ensureAlive(). this.healthy = false; const e = err as Error; const reason = e.name === "AbortError" || e.name === "TimeoutError" ? `timed out after ${timeoutMs}ms` : e.message; throw new Error(`mempalace remote request '${method}' failed: ${reason}`); } const sid = res.headers.get("Mcp-Session-Id") ?? res.headers.get("mcp-session-id"); if (sid) this.sessionId = sid; // 404 + we sent a session id → server forgot us. Drop id, re-init, retry once. if (res.status === 404 && sessionAtStart && allowReinitOn404) { try { await res.arrayBuffer(); } catch {} this.sessionId = null; await this.reinitialize(); return this.request(method, params, { timeoutMs, allowReinitOn404: false }); } if (!res.ok) { let detail = ""; try { detail = (await res.text()).slice(0, 200); } catch {} throw new Error(`mempalace remote: HTTP ${res.status} on ${method}${detail ? `: ${detail}` : ""}`); } const ct = (res.headers.get("content-type") ?? "").toLowerCase(); if (ct.includes("text/event-stream")) { return await this.readSseForResponse(res, id); } if (ct.includes("application/json")) { const json: any = await res.json(); if (json.error) throw new Error(json.error.message ?? "MCP error"); return json.result; } if (res.status === 202) return undefined; throw new Error(`mempalace remote: unexpected content-type "${ct}" on ${method}`); } private async notify(method: string, params: unknown): Promise { const body = JSON.stringify({ jsonrpc: "2.0", method, params }); let res: Response; try { res = await fetch(this.url, { method: "POST", headers: this.buildHeaders(), body, signal: this.signalFor(this.requestTimeoutMs), }); } catch (err) { throw new Error(`mempalace remote: fetch failed on notify ${method}: ${(err as Error).message}`); } try { await res.arrayBuffer(); } catch {} if (!res.ok && res.status !== 202) { throw new Error(`mempalace remote: HTTP ${res.status} on notify ${method}`); } } private async readSseForResponse(res: Response, expectedId: number): Promise { if (!res.body) throw new Error("mempalace remote: SSE response had no body"); const reader = res.body.getReader(); const decoder = new TextDecoder(); let buf = ""; let dataLines: string[] = []; const tryDispatch = (): { matched: boolean; result?: any } => { if (dataLines.length === 0) return { matched: false }; const data = dataLines.join("\n"); dataLines = []; let msg: any; try { msg = JSON.parse(data); } catch { return { matched: false }; } if (typeof msg.id === "number" && msg.id === expectedId) { if (msg.error) throw new Error(msg.error.message ?? "MCP error"); return { matched: true, result: msg.result }; } return { matched: false }; }; try { while (true) { const { done, value } = await reader.read(); if (done) break; buf += decoder.decode(value, { stream: true }); let nl: number; while ((nl = buf.indexOf("\n")) !== -1) { const rawLine = buf.slice(0, nl); buf = buf.slice(nl + 1); const line = rawLine.endsWith("\r") ? rawLine.slice(0, -1) : rawLine; if (line === "") { const out = tryDispatch(); if (out.matched) return out.result; } else if (line.startsWith(":")) { // SSE comment / keepalive — ignore } else if (line.startsWith("data:")) { dataLines.push(line.slice(5).replace(/^ /, "")); } } } const out = tryDispatch(); if (out.matched) return out.result; } finally { try { reader.cancel(); } catch {} } throw new Error(`mempalace remote: SSE stream closed without response for id=${expectedId}`); } } /** * Pick the transport from the environment: MEMPALACE_REMOTE_URL selects the * EXTERNAL (HTTP) client; unset falls back to the LOCAL stdio server. */ function createClient(): IMcpClient { const remoteUrl = process.env.MEMPALACE_REMOTE_URL?.trim(); if (remoteUrl) { return new RemoteMcpClient(remoteUrl, authHeaders()); } return new StdioMcpClient("mempalace-mcp"); } /** Bearer auth header for the remote transport, if MEMPALACE_REMOTE_TOKEN is set. */ function authHeaders(): Record | undefined { const token = process.env.MEMPALACE_REMOTE_TOKEN?.trim(); return token ? { Authorization: `Bearer ${token}` } : undefined; } export default async function mempalaceExtension(pi: ExtensionAPI) { const client = createClient(); let available = false; const agentName = process.env.MEMPALACE_AGENT_NAME ?? "pi"; // --- Edge provenance (RFC 001 §7.3.2, Phase 2) ------------------------- // Provenance belongs to the sync boundary, not to the agent. §7.3.2 ranks // "agent via skill" as the ❌ worst possible stamper — per-call boilerplate, // forgettable, improvisable — and the client/edge as the ⚠️ acceptable // interim until per-device tokens let the primary stamp authoritatively // (Phase 4). So the bridge stamps, uniformly, from host-supplied env, and no // skill instruction or LLM discipline is involved. Being self-asserted it is // ADVISORY: a hint, never load-bearing for authz or destructive scoping. // // R1 (solitary operation stays unchanged) is why this is doubly gated: it is // inert unless the host both labels the device AND points this client at a // shared palace. A solitary devbox is single-origin by definition (§7.3.3), // so it stamps nothing and loses nothing. const device = process.env.MEMPALACE_PI_DEVICE?.trim(); const shared = Boolean(process.env.MEMPALACE_REMOTE_URL?.trim()); const stampProvenance = device && shared; // Which parameter carries "who wrote this", per tool. An ALLOWLIST, never a // blanket default: mempalace 3.8.0's dispatcher rejects undeclared arguments // with -32602 (mcp_server.py:6440 — it used to drop them silently), so // injecting into a tool that has no such property would break every call. // diary_write and kg_add have no writer slot at all and are handled below. const WRITER_PARAM: Record = { mempalace_add_drawer: "added_by", mempalace_checkpoint: "added_by", mempalace_mine: "agent", mempalace_event_append: "from_agent", mempalace_artifact_put: "created_by", }; // `mine` produces bulk machine-extracted content, not agent-authored memory. // Labelling its harness segment `miner` keeps the existing pi/opencode/miner // taxonomy honest while still recording which box did the mining. const MINER_TOOLS = new Set(["mempalace_mine"]); const identity = (toolName: string) => `${MINER_TOOLS.has(toolName) ? "miner" : agentName}@${device}`; // diary_write has NO usable writer parameter (§7.3.1: "none usable") and the // device must never go in agent_name — `wing = f"wing_{agent_name}"`, so that // splinters the diary into one wing per host and hides entries from // diary_read. The entry TEXT is the only channel left, and it is also the // only one a *reader* ever sees: search results project a fixed key set and // diary_read returns content, so neither ever shows metadata. A metadata-only // fix would not have prevented the 2026-08-25 cross-host misattribution. // Prefixing (not appending) keeps the marker in chunk 0 and in the reader's // first line. AAAK is pipe-separated, so `HOST:x|SESSION:…` stays in-dialect. const HOST_MARKER = /(^|\|)\s*HOST:/; const markEntry = (v: unknown): unknown => typeof v === "string" && v.trim() !== "" && !HOST_MARKER.test(v) ? `HOST:${device}|${v}` : v; /** Stamp origin into a tool's arguments in place. Never overwrites a value * the caller set explicitly — an explicit argument wins, so a deliberate * re-file on behalf of another device stays possible. */ const applyProvenance = (toolName: string, args: Record) => { const param = WRITER_PARAM[toolName]; if (param) { const current = args[param]; if (typeof current !== "string" || current.trim() === "") { args[param] = identity(toolName); } } if (toolName === "mempalace_diary_write") { if ("entry" in args) args.entry = markEntry(args.entry); if ("content" in args) args.content = markEntry(args.content); } // checkpoint files drawers AND writes one diary entry through the same // server path, so its nested diary entry needs the marker too. if (toolName === "mempalace_checkpoint" && args.diary && typeof args.diary === "object") { const diary = args.diary as Record; if ("entry" in diary) diary.entry = markEntry(diary.entry); } }; // Gate: inject wake-up context only on the first before_agent_start of a // fresh session. Set true on resume/fork (context already in thread). let wokeUp = false; try { client.onExit = () => { // Child died (stall-kill or crash): mark unavailable. The next tool // call will attempt a bounded respawn via client.ensureAlive(). available = false; }; await client.start(); available = true; } catch (err) { // First cold-open stalled/crashed. Give the bounded self-heal a chance // before giving up entirely — a respawn often succeeds against a now-warm // page cache. Only fail-soft if even that is exhausted. process.stderr.write( `[mempalace ext] mempalace-mcp start failed: ${(err as Error).message} — attempting respawn\n`, ); available = await client.ensureAlive(); if (!available) { process.stderr.write( "[mempalace ext] mempalace-mcp unavailable after retries; continuing without palace tools\n", ); return; // fail-soft: pi keeps working without palace tools } } // Register MCP tools as pi tools. Pass the MCP `inputSchema` through as // the pi `parameters` schema so the LLM sees the real parameter names // (e.g. `agent_name`, not guessed `agent`). TypeBox schemas are plain // JSON Schema at runtime, so `Type.Unsafe` is sufficient to wrap an // externally-sourced JSON Schema — no conversion needed. for (const tool of client.tools) { const schema = tool.inputSchema && typeof tool.inputSchema === "object" ? (Type.Unsafe>(tool.inputSchema as object) as unknown as ReturnType) : Type.Object({}, { additionalProperties: true }); pi.registerTool({ name: tool.name, label: tool.name, description: tool.description ?? `MemPalace tool: ${tool.name}`, parameters: schema, async execute(_toolCallId, params) { if (!available) { // Stall-kill/crash is not a permanent latch: try a bounded respawn. available = await client.ensureAlive(); } if (!available) { return { content: [ { type: "text", text: "mempalace-mcp not available (respawn budget exhausted) — restart pi to retry palace tools", }, ], details: {}, isError: true, }; } try { const args = { ...((params ?? {}) as Record) }; if (stampProvenance) applyProvenance(tool.name, args); const result = await client.callTool(tool.name, args); // MCP tool results use { content: [...], isError?: boolean } return { content: result?.content ?? [{ type: "text", text: JSON.stringify(result) }], details: { raw: result }, isError: result?.isError === true, }; } catch (err) { return { content: [{ type: "text", text: `MCP call failed: ${(err as Error).message}` }], details: {}, isError: true, }; } }, }); } // --- Automatic transcript feeding --- // // Split deliberately: `mempalace-pi-session --prepare` does the palace-free // file work (export + quality threshold + staging, plus the rsync to the // palace host in remote mode) and prints the path to mine; we then mine it // through this client. See the header note on single-writer contention. const feedEnabled = (process.env.MEMPALACE_FEED ?? "1") !== "0"; const feedBin = process.env.MEMPALACE_FEED_BIN || "mempalace-pi-session"; const feedWing = process.env.MEMPALACE_FEED_WING ?? "wing_conversations"; const feedDebounceMs = num(process.env.MEMPALACE_FEED_DEBOUNCE_MS, 600_000); const feedPrepareTimeoutMs = num(process.env.MEMPALACE_FEED_PREPARE_TIMEOUT_MS, 120_000); const feedMineTimeoutMs = num(process.env.MEMPALACE_FEED_MINE_TIMEOUT_MS, 30_000); let lastFeedAt = 0; // 0 => the first settled turn also acts as a catch-up let feedInFlight: Promise | null = null; /** Run `mempalace-feed --prepare`; resolve the path to mine, or null. */ function prepareFeed(reason: string): Promise { return new Promise((resolve) => { // A missing helper surfaces as an async 'error' event (ENOENT), not a // throw, so the handler below is the fail-soft path. const child = spawn(feedBin, ["--prepare", "--reason", reason, "--wing", feedWing], { stdio: ["ignore", "pipe", "pipe"], }); let out = ""; let settled = false; const finish = (value: string | null) => { if (settled) return; settled = true; clearTimeout(timer); resolve(value); }; const timer = setTimeout(() => { try { child.kill("SIGKILL"); } catch { /* already gone */ } finish(null); }, feedPrepareTimeoutMs); child.stdout.on("data", (chunk) => { out += String(chunk); }); child.stderr.on("data", () => { /* the script keeps its own log */ }); child.on("error", () => finish(null)); child.on("exit", (code) => { if (code !== 0) return finish(null); const match = out.match(/^MINE_SOURCE=(.+)$/m); finish(match ? match[1].trim() : null); }); }); } /** * Stage + mine this container's transcripts. Never throws, and coalesces: * an overlapping trigger joins the in-flight run instead of racing it. */ function feedPalace(reason: string): Promise { if (!feedEnabled || !available) return Promise.resolve(); if (feedInFlight) return feedInFlight; const run = (async () => { try { const source = await prepareFeed(reason); if (!source) return; await Promise.race([ client.callTool("mempalace_mine", { source, mode: "convos", wing: feedWing, // Internal call: it does not pass through the registered tool's // execute(), so it stamps itself. These ARE this harness's own // transcripts from this device, so the harness segment is the // agent (not `miner`) even though the tool is `mine`. agent: stampProvenance ? `${agentName}@${device}` : agentName, }), new Promise((_resolve, reject) => setTimeout( () => reject(new Error(`mine timed out after ${feedMineTimeoutMs}ms`)), feedMineTimeoutMs, ), ), ]); lastFeedAt = Date.now(); } catch (err) { process.stderr.write( `[mempalace ext] feed (${reason}) failed: ${(err as Error).message}\n`, ); } })(); feedInFlight = run.finally(() => { feedInFlight = null; }); return feedInFlight; } // Mid-session feed. A hard container kill runs no handler at all, so this is // what bounds crash loss to one debounce window instead of a whole session. // Re-mining a grown transcript purges and refiles that source_file, so // repeated ticks refresh a session's drawers rather than duplicating them. pi.on("agent_settled", async () => { if (Date.now() - lastFeedAt < feedDebounceMs) return; void feedPalace("tick"); // deliberately not awaited: never stall a turn }); pi.on("session_shutdown", async () => { // Feed before stopping the client: we are the palace holder, so nothing // else can mine while we live. pi awaits this handler, so the mine really // does complete; feedMineTimeoutMs keeps a wedged palace from hanging exit. await feedPalace("shutdown"); client.stop(); }); pi.on("session_start", async (event, ctx) => { // On resume/fork, the previous session's palace context is already in // the thread — skip the wake-up injection. if (event.reason === "resume" || event.reason === "fork") { wokeUp = true; } ctx.ui.notify( `mempalace bridge: ${client.tools.length} tools registered (agent=${agentName})`, "info", ); }); // --- Automatic mailbox delivery (RFC 003 coordination log) ------------- // // Doubly gated, exactly like the stamper above: a device label AND a shared // palace. On one shared palace every client reports the same origin_replica, // so `to_agent` carries the whole distinction between machines and an // unstamped client has no address to be reached at — there is nothing to read. const mailboxEnabled = Boolean(stampProvenance) && process.env.MEMPALACE_MAILBOX !== "0"; const mailboxPollMs = num(process.env.MEMPALACE_MAILBOX_POLL_MS, 300_000); const mailboxResurfaceMs = num(process.env.MEMPALACE_MAILBOX_RESURFACE_MS, 3_600_000); // The address is the NON-miner identity: `mine` is the only tool that stamps // itself `miner@device`, and nobody addresses a mailbox to the miner. const mailboxAddress = `${agentName}@${device}`; // Only these close an obligation. `claimed` and `ready` are deliberately NOT // terminal — that is what expresses "taken, but not finished", so picked-up // work keeps resurfacing until it is actually closed out. const TERMINAL_STATUS = new Set(["applied", "superseded", "failed", "blocked"]); let lastMailboxPollAt = 0; // event id -> when we last surfaced it. In-memory ON PURPOSE. After a restart // this forgets, so an ask already seen can be shown again; that is the SAFE // failure direction. A spuriously resurfacing item is visible noise a human // corrects in one turn, whereas a suppressed unanswered ask is silent and // permanent. Never "fix" the noise by persisting this into a decision that // can hide an obligation. const surfaced = new Map(); /** Parse an `event_list` result into a plain array. Never throws. */ const parseEvents = (raw: unknown): LogEvent[] => { try { const text = extractText(raw); if (!text) return []; const parsed = JSON.parse(text); const list = Array.isArray(parsed) ? parsed : parsed?.events; return Array.isArray(list) ? (list as LogEvent[]) : []; } catch { // Malformed or non-JSON payload: report nothing owed rather than guessing. return []; } }; /** * Is `later` strictly after `earlier` in the fleet's ordering? * * Prefer `hlc` — a hybrid logical clock rendered fixed-width * (`--`), so a plain string * comparison IS the causal comparison, and it stays correct once a second * replica exists. `seq` is the arrival rowid of THIS database, so on a mesh * the same event carries a different `seq` per replica and a reply can land * before the ask it answers. * * Fall back to `seq` only when either side lacks an `hlc` (a server predating * the field, or a row whose backfill did not run). On a single replica the two * agree, so the fallback is not a downgrade today — it is the single-replica * case still being handled after the mesh case became primary. * * NEVER compare `created_at`: it is server-generated at second precision, so * ties are routine, and a tie or a skew there can suppress an UNANSWERED ask * outright. Every failure mode here is deliberately kept on the * noisy-but-visible side — an answered item resurfacing is annoying, an * unanswered ask going silent defeats the mailbox. */ const isStrictlyAfter = (later: LogEvent, earlier: LogEvent): boolean => { if ( typeof later.hlc === "string" && typeof earlier.hlc === "string" && later.hlc !== "" && earlier.hlc !== "" ) { return later.hlc > earlier.hlc; } if (typeof later.seq !== "number" || typeof earlier.seq !== "number") return false; return later.seq > earlier.seq; }; /** * Has one of MY events closed this candidate? * * Owed-ness cannot be read off `status`: event_ack APPENDS and never mutates, * and `status` is written once, so a directed "open" keeps matching the * mailbox query forever — answered or not. It has to be derived by joining * the candidate to my own replies. */ const isAnswered = (candidate: LogEvent, mine: LogEvent[]): boolean => mine.some((m) => { // (c) terminal status only. if (!TERMINAL_STATUS.has(String(m.status))) return false; // (a) LOAD-BEARING ordering test, and the reason it must exist: without // it one terminal reply suppresses every LATER ask on the same // correlation_id forever — silently, permanently, and worst on exactly // the long-running threads the correlation join is for. // See isStrictlyAfter for why the key is `hlc` and not `seq`/`created_at`. if (!isStrictlyAfter(m, candidate)) return false; // (b) join on the exact ack (written for us by event_ack), else on a // shared correlation_id — which is why correlation_id is mandatory on a // directed open: without it there is no key to join a reply back to. const ackOf = m.metadata?.ack_of; if (typeof ackOf === "string" && candidate.id && ackOf === candidate.id) return true; return Boolean(m.correlation_id && m.correlation_id === candidate.correlation_id); }); /** * Has the ORIGINAL REQUESTER explicitly withdrawn this candidate? * * RFC 003 §3.3 clause 3 clears an ask only on "no event **of yours**", and the * skill states the consequence outright: "there is nothing anyone can do about * it from the other end". That asymmetry is deliberate and mostly right — owed * ness is a statement about the RECIPIENT's accountability, and a requester * must not be able to delete an obligation the recipient genuinely has. * * It is wrong in exactly one case: the requester withdrawing its OWN ask. * * MEASURED COST, 2026-09-08/09. pi@mbp-m1-2020 withdrew a v1.8.13 rollout ask * to pi@tor-ms22 (evt seq 119: terminal `superseded`, same correlation_id, * `metadata.closes` naming it, body "DO NOT SPEND A MINUTE ON v1.8.13") and * recorded the withdrawal as done. It had no effect: seq 119's `from_agent` is * mbp, so it could never satisfy a join that only looks at tor-ms22's own * events. tor-ms22's next wake-up — 8h later, on the first boot of the image * shipping this very derivation — still listed the ask as owed, 41h old, for a * release that had been superseded and never installed there. It had to spend a * write closing an ask nobody wanted answered. The asymmetry is also INVISIBLE * from the sender's side: mbp did everything a sender is told to do and got a * result indistinguishable from success. * * WHY AN EXPLICIT MARKER AND NOT "ANY TERMINAL EVENT FROM THE REQUESTER". * This file's standing rule is that every failure mode stays on the * noisy-but-visible side, because a spuriously resurfacing item is one turn of * human correction whereas a suppressed unanswered ask is silent and permanent. * Inferring release from any terminal event would breach that: a requester * appending `applied` for its own bookkeeping — on the correlation, addressed * to me, before I ever replied — would silently delete a real obligation. * So release must be STATED, not inferred. `withdraws` is the canonical key; * `closes` is honoured because it is already this fleet's de-facto marker (mbp * seq 119, emb seq 117/118 all use it), and either must name this exact ask — * its `correlation_id` or its event `id`. Prose does not count. * * WHY THIS DOES NOT BREAK THE FIXED POINT IN RFC 003 §9.2. §9.2 rejects letting * terminal directed events into the owed set, because then every closure mints a * fresh obligation and the loop never terminates. That argument is about * CANDIDATES. This adds a CLEARER. A withdrawal carries a terminal status, so it * can never be a candidate, and nothing new becomes owed. The asserting shape * (`open`) and the clearing shape (terminal) stay disjoint. A reader who reaches * for §9.2 to object is answering a different question. */ const isWithdrawn = (candidate: LogEvent, inbound: LogEvent[]): boolean => { const requester = candidate.from_agent; if (!requester) return false; // Names THIS ask specifically: its correlation or its id. A marker naming // something else, or carrying prose, is not a withdrawal of this ask. const names = (v: unknown): boolean => typeof v === "string" && v.length > 0 && ((Boolean(candidate.correlation_id) && v === candidate.correlation_id) || (Boolean(candidate.id) && v === candidate.id)); return inbound.some((e) => { // Only the party that ASKED may retract. A third party writing a terminal // event on a shared correlation must never clear my obligation. if (e.from_agent !== requester) return false; // Addressed to the party being released. `to_agent=` also matches '*' // at the SQL level (RFC 003 §7.4), and a broadcast must not be able to // quietly empty every machine's mailbox at once. if (e.to_agent !== mailboxAddress) return false; if (!TERMINAL_STATUS.has((e.status ?? "").toLowerCase())) return false; // Same ordering guard, same reason as isAnswered: a withdrawal must not // retire an ask the requester sent LATER on the same correlation. if (!isStrictlyAfter(e, candidate)) return false; const ackOf = e.metadata?.ack_of; const joins = (typeof ackOf === "string" && Boolean(candidate.id) && ackOf === candidate.id) || Boolean(e.correlation_id && e.correlation_id === candidate.correlation_id); if (!joins) return false; return names(e.metadata?.withdraws) || names(e.metadata?.closes); }); }; /** * Directed asks addressed to this device with no terminal reply from it, and * not explicitly withdrawn by whoever sent them (see isWithdrawn). */ async function deriveOwed(): Promise { if (!mailboxEnabled || !available) return []; try { const [candidatesRaw, mineRaw, inboundRaw] = await Promise.all([ client.callTool("mempalace_event_list", { to_agent: mailboxAddress, status: "open", limit: 50, }), // `order: "desc"` is a fix, not a flourish. event_list defaults to `asc` // (append order), so this asked for the OLDEST 100 events this device // ever wrote — meaning that once a device passes 100 authored events its // most RECENT replies fall out of the join window and every ask it just // answered resurfaces as owed. Latent rather than theoretical: tor-ms22 // was at ~20 authored events when this was written. The window has to be // anchored at the newest end for the same reason the join uses `hlc`. client.callTool("mempalace_event_list", { from_agent: mailboxAddress, limit: 100, order: "desc", }), // Inbound with NO status filter, deliberately: a withdrawal carries a // TERMINAL status and is therefore structurally invisible to the // `status: "open"` candidate query above — the same omission deriveClosed // is built on. No amount of polling the open set could ever see one. client.callTool("mempalace_event_list", { to_agent: mailboxAddress, limit: 100, order: "desc", }), ]); const candidates = parseEvents(candidatesRaw).filter((c) => { // `to_agent: ` ALSO matches `*` broadcasts, per the tool contract. // The protocol says a broadcast owes nobody a reply, so a broadcast // written with status="open" must not enter anyone's owed set -- // otherwise this code contradicts the skill that documents it, and // every machine would think it personally owed the same answer. // This also gives the "don't broadcast an ask" anti-pattern teeth: // broadcasting one now demonstrably reaches no owed set at all. if (c.to_agent === "*") return false; // You cannot owe yourself. return c.from_agent !== mailboxAddress; }); if (candidates.length === 0) return []; const mine = parseEvents(mineRaw); const inbound = parseEvents(inboundRaw); return candidates.filter((c) => !isAnswered(c, mine) && !isWithdrawn(c, inbound)); } catch { // Fail silent and open: a mailbox read must never break a session. return []; } } /** * Terminal replies to asks THIS device sent. NOT work — news. * * Why this is a second query rather than a widened deriveOwed(): deriveOwed() * asks the log for `status: "open"`, and a reply that CLOSES an ask is by * definition not open, so it is structurally invisible to that filter — no * amount of polling or waiting could ever have surfaced it. * * Measured 2026-09-07: emb-7kj4vr4g closed a v1.8.13 rollout ask with a * task.reply at status=applied, and the operator reasonably expected to be * told. The mailbox stayed silent and was CORRECT to — "owed" means "you must * reply", and nothing was owed. But the single most useful thing a fleet can * tell a human is "the thing you asked for is done" (here: done, AND your * premise was wrong), and that was the one category the feature could not * report. Silence was right by its own definition and wrong by the user's. */ async function deriveClosed(): Promise { if (!mailboxEnabled || !available) return []; try { const [mineRaw, inboundRaw] = await Promise.all([ client.callTool("mempalace_event_list", { from_agent: mailboxAddress, limit: 100 }), // NO status filter, deliberately: that omission is the entire fix. client.callTool("mempalace_event_list", { to_agent: mailboxAddress, limit: 50 }), ]); // Correlations this device opened as a DIRECTED ask. A broadcast owes // nobody a reply, so it cannot be closed by one either. const myAsks = new Map(); for (const e of parseEvents(mineRaw)) { if (!e.correlation_id || !e.to_agent) continue; if (e.to_agent === "*" || e.to_agent === mailboxAddress) continue; if (e.type !== "task.request") continue; myAsks.set(e.correlation_id, e); } if (myAsks.size === 0) return []; // Newest terminal reply per correlation only. A peer that appends // applied-then-superseded should cost one line, not a wall of them. const best = new Map(); for (const e of parseEvents(inboundRaw)) { if (e.from_agent === mailboxAddress) continue; // cannot inform myself const cid = e.correlation_id; if (!cid || !myAsks.has(cid)) continue; // NOTE the deliberate asymmetry with deriveOwed: a broadcast is excluded // there (it owes nobody) but allowed here, because a peer answering on MY // correlation_id is news to me regardless of how widely it was addressed. if (!TERMINAL_STATUS.has((e.status ?? "").toLowerCase())) continue; const prev = best.get(cid); if (!prev || isStrictlyAfter(e, prev)) best.set(cid, e); } return [...best.values()]; } catch { // Same contract as deriveOwed: a mailbox read never breaks a session. return []; } } const relAge = (iso: string | undefined): string => { if (!iso) return "age unknown"; const then = Date.parse(iso); if (!Number.isFinite(then)) return "age unknown"; const mins = Math.max(0, Math.round((Date.now() - then) / 60_000)); if (mins < 60) return `${mins}m ago`; const hours = Math.round(mins / 60); return hours < 48 ? `${hours}h ago` : `${Math.round(hours / 24)}d ago`; }; const excerpt = (body: string | undefined, max = 200): string => { if (typeof body !== "string") return ""; const flat = body.replace(/\s+/g, " ").trim(); return flat.length > max ? `${flat.slice(0, max)}\u2026` : flat; }; /** One terse line per ask — this goes into a live context window. */ const formatOwed = (items: LogEvent[]): string => items .map((e) => { const head = [ e.id ?? "(no id)", `from ${e.from_agent ?? "?"}`, e.type ?? "?", e.correlation_id ? `correlation_id=${e.correlation_id}` : "NO correlation_id", relAge(e.created_at), ].join(" \u00b7 "); const ex = excerpt(e.body); return ex ? `- ${head}\n ${ex}` : `- ${head}`; }) .join("\n"); const OWED_HOWTO = "To close one, append an event on the SAME correlation_id with a terminal status " + "(applied/superseded/failed/blocked). An ack alone does NOT clear it, and neither " + "does claimed/ready — those deliberately keep it visible as taken-but-unfinished."; // Written for the ONE reader the rest of this text ignores: the human watching // the window. Measured 2026-08-26 on two devices — a delivery lands, the agent // is idle, nothing happens, and the operator asks "do I have to nudge you for // you to read this?". Yes, and it is by design (see the deliverAs comment // below): at agent_settled no inference is running, so the text is queued for // the next turn. The message explained how to CLOSE an ask but never when it // would be SEEN, so the person who needed that fact was the only one not told. // Cheapest possible fix, and it changes no behaviour: say so in the text. const QUEUED_NOTE = "DELIVERY NOTE — this is a QUEUED message, not an action: it arrived while this " + "agent was idle, so no model call was made and nothing woke it. It is read at the " + "start of the next turn. If you are the human watching and want it handled now, " + "send any message to start a turn; the agent is not ignoring the ask, it is not running."; const CLOSED_NOTE = "NO ACTION IS OWED on these — they are replies to asks THIS device sent, shown " + "once each because 'the thing you asked for is done' is news you wanted and the " + "owed-set could never carry it. Read the reply before assuming your original ask " + "was right: a peer that did the work is the most likely party to have found your " + "premise wrong."; // Mid-session mailbox. Tier 1 (the wake-up injection below) owns the first // look; this exists because arrivals are bursty and correlate with our own // activity — measured 2026-08-26, 11 of 22 events on this log landed inside a // single 5h38m working window, so session-start-only delivery sleeps through // most of a burst. `agent_settled` is already activity-coupled, which is why // it is the trigger rather than a wall-clock timer. // // POLLING AND DELIVERY ARE DECOUPLED. We poll on a floor, but speak only when // the owed set holds something not already surfaced. Re-announcing the same // ask every poll is precisely how a channel teaches its reader to ignore it, // which is the failure this whole feature exists to reverse. // --- A: tell the human at poll time (RFC 003 §7.11) ----------------------- // // B (the QUEUED_NOTE above) fixes the confusion of someone who IS looking at // the window. It does nothing for the case that actually loses an ask: nobody // is looking. The delivered text is queued for a turn that only a human starts, // so an ask can wait indefinitely on an idle session while its recipient makes // coffee. A notification at poll time is the only part of this feature that // reaches a person who is not watching. // // Two surfaces, deliberately split by risk: // default in-TUI ctx.ui.notify, exactly as session_start already does. // Zero new output channels; visible only if you are looking. // =desktop additionally emit a terminal-native notification, reusing // the detection the fleet's own notify.ts already proves in // this harness (Kitty OSC 99 / Windows toast / OSC 777). This // escapes the container without notify-send, DBus or any host // access: the escape sequence is interpreted by the terminal // emulator on the human's machine. Opt-in because writing raw // escapes to stdout is a behaviour change on shared machines, // not because it is unreliable. // =kitty =osc777 force one protocol, skipping detection entirely. // =0 silent, for anyone who wants the mailbox without pings. // // WHY FORCING EXISTS — measured on tor-ms22 2026-08-26, and it invalidates // detection in exactly the deployment this ships in. Terminal identity lives in // env vars set by the emulator (KITTY_WINDOW_ID, TERM_PROGRAM) and `docker exec` // does NOT forward them: inside the container pi sees only TERM=xterm-256color // no matter what is rendering it. So detection can never see Kitty from in here // and always falls through to OSC 777, which Kitty does not implement — on a // containerised client the notification would silently do nothing, the worst // possible failure for a feature whose whole job is to break a silence. Naming // the protocol ends the guessing. // // Fired only when something is actually due, i.e. the same condition as the // delivery itself: a notification that fires on an empty poll would teach its // reader to ignore it, which is the failure this whole feature exists to undo. const mailboxNotifyMode = (process.env.MEMPALACE_MAILBOX_NOTIFY ?? "").trim().toLowerCase(); const mailboxNotifyEnabled = mailboxNotifyMode !== "0" && mailboxNotifyMode !== "off"; const mailboxNotifyTerminal = mailboxNotifyMode === "desktop" || mailboxNotifyMode === "kitty" || mailboxNotifyMode === "osc777"; /** Terminal-native notification. Best-effort; never throws into a handler. */ const notifyTerminal = (title: string, body: string): void => { try { const esc = (s: string): string => s.replace(/[\x00-\x1f\x7f;]/g, " "); if ( mailboxNotifyMode === "kitty" || (mailboxNotifyMode === "desktop" && !!process.env.KITTY_WINDOW_ID) ) { process.stdout.write(`\x1b]99;i=1:d=0;${esc(title)}\x1b\\`); process.stdout.write(`\x1b]99;i=1:p=body;${esc(body)}\x1b\\`); return; } process.stdout.write(`\x1b]777;notify;${esc(title)};${esc(body)}\x07`); } catch { /* a notification is never worth an exception */ } }; pi.on("agent_settled", async (_event, ctx) => { if (!mailboxEnabled || !available) return; if (!wokeUp) return; // the wake-up injection has not run yet; it does look #1 if (Date.now() - lastMailboxPollAt < mailboxPollMs) return; // Claim the slot BEFORE awaiting anything, so overlapping settles cannot // stampede past the floor into concurrent polls. lastMailboxPollAt = Date.now(); void (async () => { try { const [owed, closed] = await Promise.all([deriveOwed(), deriveClosed()]); const now = Date.now(); const unseen = (list: LogEvent[]): LogEvent[] => list.filter((e) => { if (!e.id) return false; const last = surfaced.get(e.id); return last === undefined || now - last >= mailboxResurfaceMs; }); const due = unseen(owed); // Closing replies are announced ONCE and never resurface: an ask you // already know is finished is not a nag, and re-announcing it hourly is // the train-the-reader-to-ignore-it failure this window exists to stop. const newsRaw = closed.filter((e) => e.id && !surfaced.has(e.id)); if (due.length === 0 && newsRaw.length === 0) return; // silence is correct for (const e of [...due, ...newsRaw]) if (e.id) surfaced.set(e.id, now); const blocks: string[] = []; if (due.length > 0) { blocks.push( `${due.length} directed ask(s) addressed to "${mailboxAddress}" with no ` + `terminal reply from this device. ${OWED_HOWTO}\n\n${formatOwed(due)}`, ); } if (newsRaw.length > 0) { blocks.push( `${newsRaw.length} repl(y|ies) CLOSING an ask this device sent. ` + `${CLOSED_NOTE}\n\n${formatOwed(newsRaw)}`, ); } pi.sendMessage( { customType: "mempalace-mailbox", content: `MemPalace logstream mailbox.\n\n${QUEUED_NOTE}\n\n` + blocks.join("\n\n---\n\n"), display: true, }, // "steer" is queue-on-arrival, NOT an interrupt: at agent_settled the // agent is idle, so this lands at the start of the next turn without // spending an LLM call. Chosen over "nextTurn" only because if a settle // ever races a freshly-started turn, "steer" delivers within it instead // of waiting for the next prompt. Deliberately NO triggerTurn: waking // the model on inbound fleet traffic is a far bigger behavioural change // than auto-delivery and is not what was approved. { deliverAs: "steer" }, ); // AFTER the queueing call, so a notification can never be the only // thing that happened: the message is in the thread first, then we // point at it. Wording names the nudge explicitly, because "you have // mail" without "press a key" reproduces the exact confusion B fixes. if (mailboxNotifyEnabled) { const parts = [ due.length > 0 ? `${due.length} ask(s) owed` : "", newsRaw.length > 0 ? `${newsRaw.length} closed` : "", ].filter(Boolean); const summary = `${parts.join(" + ")} for ${mailboxAddress} — send any message to handle`; try { if (ctx?.hasUI) ctx.ui.notify(`MemPalace mailbox: ${summary}`, "info"); } catch { /* ignore: UI may be gone by the time the poll resolves */ } if (mailboxNotifyTerminal) notifyTerminal("MemPalace mailbox", summary); } } catch { /* best-effort delivery: never break a turn over a mailbox read */ } })(); // deliberately not awaited: never stall a turn }); // --- Auto wake-up (mempalace skill Phase 1) --- pi.on("before_agent_start", async (_event, _ctx) => { if (wokeUp || !available) return; wokeUp = true; // one-shot, even if the calls below fail const sections: string[] = []; try { const status = await client.callTool("mempalace_status", {}); const text = extractText(status); if (text) sections.push(`## mempalace_status\n\n${text}`); } catch (err) { sections.push(`## mempalace_status\n\n(error: ${(err as Error).message})`); } try { const diary = await client.callTool("mempalace_diary_read", { agent_name: agentName, last_n: 5, }); const text = extractText(diary); if (text) sections.push(`## mempalace_diary_read (agent=${agentName}, last_n=5)\n\n${text}`); } catch (err) { sections.push(`## mempalace_diary_read\n\n(error: ${(err as Error).message})`); } // Tier 1 mailbox: what this device owes a reply to, derived rather than read // off a filter. deriveOwed() swallows its own failures and returns [], so a // broken palace costs a missing section, never a broken wake-up. if (mailboxEnabled) { const [owed, closed] = await Promise.all([deriveOwed(), deriveClosed()]); if (owed.length > 0) { // Record what this injection showed, so the first mid-session poll does // not re-announce the identical list minutes later. Without this the // two tiers double up at every session start, which is the same // train-the-reader-to-ignore-it failure the resurface window prevents. const now = Date.now(); for (const e of owed) if (e.id) surfaced.set(e.id, now); sections.push( `## logstream mailbox (${owed.length} owed)\n\n` + `Directed asks addressed to "${mailboxAddress}" with no terminal reply from ` + `this device. ${OWED_HOWTO}\n\n${formatOwed(owed)}`, ); } const news = closed.filter((e) => e.id && !surfaced.has(e.id)); if (news.length > 0) { const now = Date.now(); for (const e of news) if (e.id) surfaced.set(e.id, now); sections.push( `## logstream mailbox (${news.length} closed — no action owed)\n\n` + `${CLOSED_NOTE}\n\n${formatOwed(news)}`, ); } } if (sections.length === 0) return; const body = `MemPalace wake-up context (auto-injected by the mempalace extension). ` + `This is your palace orientation for this session — do not announce it to the user, ` + `just use it to inform your answers. Agent identity for diary tools: "${agentName}".\n\n` + (device ? `You are running on device "${device}"${shared ? "" : " (solitary palace)"}. ` + `Drawers and diary entries below may have been written by a DIFFERENT machine — ` + `diary_read returns every device's "${agentName}" diary interleaved. Check the ` + `HOST: marker or the drawer's device metadata before treating a past session as ` + `this machine's history.\n\n` : "") + sections.join("\n\n---\n\n"); return { message: { customType: "mempalace-wakeup", content: body, display: true, }, }; }); // --- Manual wind-down (mempalace skill Phase 3) --- pi.registerCommand("mempalace-diary", { description: "Ask the LLM to write an AAAK diary entry summarizing this session", handler: async (args, ctx) => { if (!available) { ctx.ui.notify("mempalace bridge not available", "warning"); return; } const topic = args.trim() || "session-summary"; const prompt = `Write a MemPalace diary entry for this session using the AAAK format ` + `described in the mempalace skill. Call mempalace_diary_write with ` + `agent_name="${agentName}", topic="${topic}", and a compressed AAAK entry ` + `that summarizes what we worked on, what was discovered, and any open ` + `threads. Then confirm the write succeeded. Do not ask me for ` + `clarification — draft from the session so far.`; pi.sendUserMessage(prompt); }, }); } /** Flatten MCP tool result content into plain text for context injection. */ function extractText(mcpResult: any): string { const parts = mcpResult?.content; if (!Array.isArray(parts)) return typeof mcpResult === "string" ? mcpResult : ""; return parts .map((p: any) => (typeof p?.text === "string" ? p.text : "")) .filter(Boolean) .join("\n"); }