309980b62c
"[mempalace ext] feed (tick) failed: mine timed out after 30000ms" was parked as
cosmetic on 2026-08-27. It is not cosmetic: the tight deadline was the trigger,
but the defect is a positive feedback loop that puts multiple writers on a
single-writer palace.
lastFeedAt was assigned only AFTER a successful await. The Promise.race abandons
our WAIT and cannot cancel the server's work, so on a mine that takes longer
than the deadline -- measured at 30-60s in normal operation, against a 30s
deadline -- the catch ran with lastFeedAt UNCHANGED. That left both guards in
the agent_settled handler open at once: the debounce test
(`Date.now() - lastFeedAt < feedDebounceMs`) passed because lastFeedAt was still
stale, and feedInFlight was already null because `run` had settled. Every
subsequent settled turn therefore launched another mine on top of the one still
running, each making the next slower and the next timeout likelier -- which is
why operators saw the message many times per session instead of at most once per
10-minute debounce window.
Fix, three lines:
- move `lastFeedAt = Date.now()` to before the await, so a timeout still starts
the debounce clock. A timeout is not a "did not happen": the mine is running
server-side and `mine --mode convos` dedups by source_file and is idempotent.
- raise MEMPALACE_FEED_MINE_TIMEOUT_MS from 30_000 to 300_000. The mine is the
slowest thing this extension does yet carried the tightest deadline: 4x
tighter than the prepare step before it (120_000) and 10x tighter than the
init handshake (300_000), a fast call. All three were introduced together in
29e660e and this one was never revisited. 300_000 matches the init timeout
because liveness is the only job left for this deadline -- it cannot cancel
the server's work, so it must sit far above the slowest honest completion.
Simulated both guards over 10 minutes of settled turns at 20s intervals with a
60s mine: BEFORE 16 mines launched, 15 of them overlapping an already-running
mine; AFTER 2 launched, 0 overlapping. With a mine that exceeds even the new
deadline (400s): BEFORE 16/15, AFTER still 2/0 -- the lastFeedAt move is what
actually fixes it, and it holds even when the timeout still fires. The raise
stops the spurious message; the move stops the pile-up.
NOT fixed here: the message still goes to process.stderr, which pi renders into
the TUI input field. That needs a pi-side channel or a log file, and is tracked
separately. After this change the message should be rare, and when it does
appear it means something real: a mine exceeding five minutes.
1601 lines
68 KiB
TypeScript
1601 lines
68 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);
|
||
// 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();
|
||
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,
|
||
),
|
||
),
|
||
]);
|
||
} 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([
|
||
client.callTool("mempalace_event_list", {
|
||
to_agent: mailboxAddress,
|
||
status: "open",
|
||
limit: 50,
|
||
}),
|
||
// `order: "desc"` is a fix, not a flourish. event_list defaults to `asc`
|
||
// (append order), so this asked for the OLDEST 100 events this device
|
||
// ever wrote — meaning that once a device passes 100 authored events its
|
||
// most RECENT replies fall out of the join window and every ask it just
|
||
// answered resurfaces as owed. Latent rather than theoretical: tor-ms22
|
||
// was at ~20 authored events when this was written. The window has to be
|
||
// anchored at the newest end for the same reason the join uses `hlc`.
|
||
client.callTool("mempalace_event_list", {
|
||
from_agent: mailboxAddress,
|
||
limit: 100,
|
||
order: "desc",
|
||
}),
|
||
// Inbound with NO status filter, deliberately: a withdrawal carries a
|
||
// TERMINAL status and is therefore structurally invisible to the
|
||
// `status: "open"` candidate query above — the same omission deriveClosed
|
||
// is built on. No amount of polling the open set could ever see one.
|
||
client.callTool("mempalace_event_list", {
|
||
to_agent: mailboxAddress,
|
||
limit: 100,
|
||
order: "desc",
|
||
}),
|
||
]);
|
||
const candidates = parseEvents(candidatesRaw).filter((c) => {
|
||
// `to_agent: <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 {
|
||
const [mineRaw, inboundRaw] = await Promise.all([
|
||
client.callTool("mempalace_event_list", { from_agent: mailboxAddress, limit: 100 }),
|
||
// NO status filter, deliberately: that omission is the entire fix.
|
||
client.callTool("mempalace_event_list", { to_agent: mailboxAddress, limit: 50 }),
|
||
]);
|
||
// Correlations this device opened as a DIRECTED ask. A broadcast owes
|
||
// nobody a reply, so it cannot be closed by one either.
|
||
const myAsks = new Map<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");
|
||
}
|