Files
mempalace-toolkit/extensions/pi/mempalace.ts
T
joakimp 2167a1b033 fix(mailbox): pass order: "desc" on every cursor-less event_list call
Three of the five `mempalace_event_list` calls in the pi extension relied on
the server's default ordering: deriveOwed's `status: "open"` candidate query
(limit 50) and both windows in deriveClosed (from_agent limit 100, to_agent
limit 50). deriveOwed's other two calls already say `order: "desc"` and carry
the comment explaining why: on mempalace <= 3.9.0 the default is `asc`, so a
cursor-less `limit: N` returns the OLDEST N events, and once a device passes N
authored events its newest asks and the replies that close them fall outside
the join window. deriveClosed had exactly that latent truncation.

Why now: mempalace 3.10.0 (2026-09-15) flips the cursor-less default to
newest-first ("Logstream listings default to the newest events when no cursor
is given"). That change is server-side -- the extension talks to the fleet
hub over MEMPALACE_REMOTE_URL, so it lands when the hub upgrades, not when a
client does -- and it would have silently FIXED deriveClosed on 3.10.0 while
leaving it broken against any 3.9.0 hub. A mailbox verdict that depends on
which server version answers is the wrong shape. Saying `order` explicitly
makes both derive* functions read the same window on either version.

The selection in deriveClosed (isStrictlyAfter per correlation) is
order-independent, so this changes WHICH events are in the window, not how
the winner is picked. No behaviour change on a device with < 50 inbound and
< 100 authored events; the fleet hub is at 176 events total today, so the
window was about to matter.

Verified: file compiles under pi's own loader (jiti 2.7.0 from the installed
pi-coding-agent) before and after; `order:"desc"` literal count 3 -> 7, which
is the 3 added arguments plus 1 comment mention. `node --check` and a bare
`tsc` were tried first and both fail identically on HEAD (inline `type`
imports / no @types/node) -- checker faults, not this change.
2026-09-22 14:15:22 +02:00

1640 lines
70 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
* 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 <ours>". 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);
* the feed's mine carries its own, longer
* deadline (MEMPALACE_FEED_MINE_TIMEOUT_MS)
* - 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<typeof setTimeout> | 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<void>;
/**
* `opts.timeoutMs` overrides the transport's generic per-request deadline
* for THIS call only. Callers that knowingly invoke a long server-side job
* (the feed's `mempalace_mine`) pass their own deadline here; everything
* else keeps the short default, which is the wedged-query guard.
*/
callTool(name: string, args: Record<string, unknown>, opts?: { timeoutMs?: number }): Promise<any>;
ensureAlive(): Promise<boolean>;
stop(): void | Promise<void>;
}
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<string, unknown> | null;
};
const sleep = (ms: number): Promise<void> =>
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<number, Pending>();
private stdoutBuf = "";
private ready: Promise<void> | 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<boolean> | 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<void> {
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<boolean> {
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<any> {
const id = this.nextId++;
return new Promise((resolve, reject) => {
let timer: ReturnType<typeof setTimeout> | 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<string, unknown>, opts?: { timeoutMs?: number }): Promise<any> {
// The per-call override matters MORE here than for the HTTP client: on
// timeout this transport kills the server child, so a generic deadline
// that undercuts a long mine does not merely abandon the wait — it
// aborts the mine.
return this.request("tools/call", { name, arguments: args }, opts?.timeoutMs ?? this.requestTimeoutMs);
}
/** 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.
// • callTool() takes an optional per-call `{ timeoutMs }` (IMcpClient
// contract, see there) so the feed's long-running mine is not cut off by
// the generic per-request deadline. Not a protocol change; sync token
// unchanged.
//
// 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<string, string>;
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<boolean> | null = null;
constructor(url: string, headers?: Record<string, string>) {
this.url = url;
this.extraHeaders = headers ?? {};
}
get alive(): boolean {
return this.healthy;
}
async start(): Promise<void> {
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<string, unknown>, opts?: { timeoutMs?: number }): Promise<any> {
return this.request("tools/call", { name, arguments: args }, { timeoutMs: opts?.timeoutMs });
}
/**
* 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<boolean> {
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<void> {
this.healthy = false;
if (!this.sessionId) return;
try {
await fetch(this.url, { method: "DELETE", headers: this.buildHeaders() });
} catch {}
}
private buildHeaders(): Record<string, string> {
const h: Record<string, string> = {
"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<void> {
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<any> {
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<void> {
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<any> {
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<string, string> | 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<string, string> = {
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<string, unknown>) => {
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<string, unknown>;
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<Record<string, unknown>>(tool.inputSchema as object) as unknown as ReturnType<typeof Type.Object>)
: 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<string, unknown>) };
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);
// 300_000, raised from 30_000 on 2026-09-10. The mine is the SLOWEST thing
// this extension does — it offers every qualifying session transcript to a
// SINGLE-WRITER palace, measured at 30–60s in normal operation and growing
// with the corpus — yet it carried by far the TIGHTEST deadline: 4x tighter
// than the prepare step that precedes it (120_000) and 10x tighter than the
// init handshake (300_000), which is a fast call. All three were introduced
// together in 29e660e (2026-08-12) and this one was never revisited, so the
// deadline fired during entirely normal operation and the resulting message
// read as an error when nothing had gone wrong.
//
// Matched to the init timeout because liveness is the ONLY legitimate job
// left for this deadline: the Promise.race below abandons our WAIT, it cannot
// cancel the server's work, so the deadline buys nothing except an escape
// from a permanently hung call. It must therefore sit far above the slowest
// honest completion, not near it.
const feedMineTimeoutMs = num(process.env.MEMPALACE_FEED_MINE_TIMEOUT_MS, 300_000);
let lastFeedAt = 0; // 0 => the first settled turn also acts as a catch-up
let feedInFlight: Promise<void> | null = null;
/** Run `mempalace-feed --prepare`; resolve the path to mine, or null. */
function prepareFeed(reason: string): Promise<string | null> {
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<void> {
if (!feedEnabled || !available) return Promise.resolve();
if (feedInFlight) return feedInFlight;
const run = (async () => {
try {
const source = await prepareFeed(reason);
if (!source) return;
// Record the attempt HERE, before awaiting — not after a successful
// wait. The race below abandons only our WAIT; the mine keeps running
// server-side, and `mine --mode convos` dedups by source_file and is
// idempotent, so a timeout is emphatically not a "did not happen".
//
// Leaving lastFeedAt stale on the timeout path defeated the debounce
// guard in the agent_settled handler below (`Date.now() - lastFeedAt <
// feedDebounceMs`): with lastFeedAt unchanged that guard passed on EVERY
// settled turn, and because `run` had already settled, feedInFlight was
// null too — so BOTH guards stood open. Each settled turn then launched
// another mine while the previous one was still running: overlapping
// writers queueing on a single-writer palace, each making the next one
// slower and the next timeout likelier. That positive feedback loop,
// not the tight deadline by itself, is why operators saw "mine timed out
// after 30000ms" many times per session rather than at most once per
// debounce window.
lastFeedAt = Date.now();
// The deadline is passed DOWN to the transport as well as raced
// here. Until 2026-09-18 it was only raced: callTool() had no way
// to carry it, so the transport's generic per-request timeout
// (MEMPALACE_MCP_TIMEOUT_MS, 60 000) fired first on every honest
// 60 s+ mine — "remote request 'tools/call' failed: timed out
// after 60000ms" — and the 300 000 below was unreachable. The
// race stays as the liveness guard for a transport whose timeout
// is disabled (0).
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,
},
{ timeoutMs: feedMineTimeoutMs },
),
new Promise((_resolve, reject) =>
setTimeout(
() => reject(new Error(`mine timed out after ${feedMineTimeoutMs}ms`)),
feedMineTimeoutMs,
),
),
]);
} 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<string, number>();
/** 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
* (`<unix_ms:13 digits>-<counter:6 hex>-<replica_id>`), 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=<me>` 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<LogEvent[]> {
if (!mailboxEnabled || !available) return [];
try {
const [candidatesRaw, mineRaw, inboundRaw] = await Promise.all([
// Every cursor-less event_list call in this file passes `order` explicitly.
// The server's default is not ours to lean on: mempalace <= 3.9.0 defaults
// to `asc`, 3.10.0 flips a cursor-less listing to newest-first. Both
// derive* joins want the NEWEST window (see the comment below), so say so
// and the verdict stops depending on which server version answers.
client.callTool("mempalace_event_list", {
to_agent: mailboxAddress,
status: "open",
limit: 50,
order: "desc",
}),
// `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: <me>` 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<LogEvent[]> {
if (!mailboxEnabled || !available) return [];
try {
// `order: "desc"` for the same reason as deriveOwed: without it a 3.9.0
// server hands back the OLDEST window, so past 100 authored events this
// device's newest asks (and past 50 inbound, the replies that close them)
// fall outside the join. The selection below is order-independent
// (isStrictlyAfter), so this only fixes WHICH events are in the window.
const [mineRaw, inboundRaw] = await Promise.all([
client.callTool("mempalace_event_list", { from_agent: mailboxAddress, limit: 100, order: "desc" }),
// NO status filter, deliberately: that omission is the entire fix.
client.callTool("mempalace_event_list", { to_agent: mailboxAddress, limit: 50, order: "desc" }),
]);
// 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<string, LogEvent>();
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<string, LogEvent>();
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");
}