a92c75d070
Two changes, neither touching the no-triggerTurn decision, which stands.
B — a delivery note in the message itself. The mid-session text explained how to
CLOSE an ask and never said when it would be SEEN, so the only reader who needed
that fact — the human watching an idle session — was the one not told. It now
says: this is a queued message, nothing woke the agent, any message starts the
turn that handles it, and the agent is not ignoring the ask, it is not running.
Costs nothing and changes no behaviour; it converts "why is it ignoring me?" into
"right, I nudge it". Deliberately NOT added to the wake-up injection, where a
turn is already starting and the note would be false.
A — a notification at poll time, because B only helps someone already looking and
the case that loses an ask is nobody looking. MEMPALACE_MAILBOX_NOTIFY: unset →
in-TUI ctx.ui.notify, the surface session_start already uses; =desktop →
additionally a terminal-native notification (Kitty OSC 99, else OSC 777), reusing
the detection the fleet's own notify.ts already proves in this harness; =0/off →
silent. The desktop path is how a ping escapes a container with no notify-send,
no DBus and no host access: the escape sequence is written to stdout and
interpreted by the terminal emulator on the human's own machine. Opt-in because
writing raw escapes is a behaviour change on a shared machine, not because it is
unreliable — say the word and the default flips.
Placement and firing conditions are deliberate: the notify call sits AFTER
sendMessage so a ping can never be the only thing that happened, and it fires
only when something is due — the same condition as delivery. A notification on
an empty poll would train its reader to ignore it, which is the failure this
whole feature exists to reverse. The wording names the nudge ("send any message
to handle") because "you have mail" without "press a key" reproduces exactly the
confusion B fixes.
Tested: 12 cases. Mode parsing (unset/empty/desktop/DESKTOP-with-space/0/off/1),
escape hygiene, and the OSC invariant that matters — a hostile title or body
containing ";" or a BEL cannot forge an OSC field or terminate the sequence
early, which is worth asserting because event bodies arrive from other machines.
ctx.hasUI is checked and the notify call is wrapped, since the UI can be gone by
the time an unawaited poll resolves. Syntax checked with node --strip-types.
1365 lines
55 KiB
TypeScript
1365 lines
55 KiB
TypeScript
/**
|
|
* 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)
|
|
* - 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>;
|
|
callTool(name: string, args: Record<string, unknown>): 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>): Promise<any> {
|
|
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<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>): Promise<any> {
|
|
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<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);
|
|
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<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;
|
|
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<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);
|
|
});
|
|
|
|
/** Directed asks addressed to this device with no terminal reply from it. */
|
|
async function deriveOwed(): Promise<LogEvent[]> {
|
|
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: <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);
|
|
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.";
|
|
|
|
// 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.";
|
|
|
|
// 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.
|
|
// =0 silent, for anyone who wants the mailbox without pings.
|
|
//
|
|
// 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 mailboxNotifyDesktop = mailboxNotifyMode === "desktop";
|
|
|
|
/** 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 (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 = 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` +
|
|
`${QUEUED_NOTE}\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" },
|
|
);
|
|
// 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 summary =
|
|
`${due.length} directed ask(s) queued 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 (mailboxNotifyDesktop) 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 = 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");
|
|
}
|