/** * 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; 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 []; } }; /** * 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. // // Compare `seq`, NEVER `created_at`. `seq` is replica-local (it equals // origin_seq only while a single replica authors for every machine); the // durable key once `mempalace_mesh_peers` reports real peers is `hlc`, // which is total and causally consistent. Local-seq skew can only make an // ANSWERED item resurface (noise, visible), whereas a timestamp // comparison can suppress an UNANSWERED ask outright. if (typeof m.seq !== "number" || typeof candidate.seq !== "number") return false; if (m.seq <= candidate.seq) 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); }); /** Directed asks addressed to this device with no terminal reply from it. */ async function deriveOwed(): Promise { if (!mailboxEnabled || !available) return []; try { const [candidatesRaw, mineRaw] = await Promise.all([ client.callTool("mempalace_event_list", { to_agent: mailboxAddress, status: "open", limit: 50, }), client.callTool("mempalace_event_list", { from_agent: mailboxAddress, limit: 100 }), ]); 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); return candidates.filter((c) => !isAnswered(c, mine)); } catch { // Fail silent and open: a mailbox read must never break 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."; // 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. pi.on("agent_settled", async () => { 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 = await deriveOwed(); const now = Date.now(); const due = owed.filter((e) => { if (!e.id) return false; const last = surfaced.get(e.id); return last === undefined || now - last >= mailboxResurfaceMs; }); if (due.length === 0) return; // silence is the correct output here for (const e of due) if (e.id) surfaced.set(e.id, now); pi.sendMessage( { customType: "mempalace-mailbox", content: `MemPalace logstream mailbox: ${due.length} directed ask(s) addressed to ` + `"${mailboxAddress}" with no terminal reply from this device. ${OWED_HOWTO}\n\n` + formatOwed(due), 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" }, ); } 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 = await deriveOwed(); 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)}`, ); } } 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"); }