diff --git a/extensions/pi/README.md b/extensions/pi/README.md index 3876d55..359bc74 100644 --- a/extensions/pi/README.md +++ b/extensions/pi/README.md @@ -301,13 +301,42 @@ identity — `from_agent` / `to_agent` carry the entire distinction. Measured So the stamping above is what makes `to_agent="pi@tor-ms22"` mean anything, and an *unstamped* client addressed as bare `pi` is unreachable. -**2. The bridge is write-only today.** It stamps events on the way out and never -reads the log: there is no mailbox, no poll, no delivery. An event addressed to -this machine by name reaches the agent only if the agent runs -`mempalace_event_list` itself — which is why that query is a wake-up step in the -*consumer* skill (`~/.agents/skills/mempalace/SKILL.md`), and why the protocol -norms live there rather than here. **This file documents the mechanism; the -skill is normative for behaviour.** +**2. The bridge now READS the log too — auto-delivered mailbox.** It derives what +this device owes and injects it, so an event addressed to this machine no longer +waits for the agent to think of asking. Two delivery points, both fail-silent: + +- **Session start** — one more section in the existing `before_agent_start` + wake-up injection, alongside `mempalace_status` and `diary_read`. +- **Mid-session** — an `agent_settled` poll, floored at + `MEMPALACE_MAILBOX_POLL_MS` (default 300000, i.e. 5 min), delivering via + `pi.sendMessage(..., { deliverAs: "steer" })`. Note this **queues, it does not + interrupt**: at `agent_settled` the agent is idle, so the message lands at the + start of the next turn and spends no LLM call. There is deliberately no + `triggerTurn` — waking the model on inbound fleet traffic is a much larger + behavioural change than auto-delivery. + +Owed-ness is **derived, never read off a field**, because `event_ack` appends and +`status` is written once: a directed `open` event matches the mailbox query +*forever*, answered or not. Two calls (`to_agent= status=open`, and +`from_agent=`), then a candidate counts as answered only when one of this +device's own events has a **strictly higher `seq`**, joins via +`metadata.ack_of` or a shared `correlation_id`, *and* carries a terminal status +(`applied`/`superseded`/`failed`/`blocked`). The `seq` test is load-bearing: +without it one terminal reply suppresses every later ask on that correlation +forever. `seq` is replica-local, so `hlc` is the correct key once `mesh_peers` +reports actual peers. + +`*` broadcasts are excluded even though `to_agent=` matches them, because the +protocol says a broadcast owes nobody a reply — which also means broadcasting an +ask demonstrably reaches no owed set, giving that documented anti-pattern teeth. + +Gated on `MEMPALACE_PI_DEVICE` **and** `MEMPALACE_REMOTE_URL` (the same pair as +the stamper, since an unstamped client has no address to be reached at), and +disabled outright with `MEMPALACE_MAILBOX=0`. Inert on a solitary palace: no +calls, no injection. Delivery is the mechanism; the *norms* — what a reply owes, +and that only a terminal event closes a thread — remain normative in the +*consumer* skill (`~/.agents/skills/mempalace/SKILL.md`). **This file documents +the mechanism; the skill is normative for behaviour.** **3. Live push is a deployment question, not a code one.** The palace implements an SSE endpoint (`GET /logstream/stream`, `text/event-stream` in @@ -315,8 +344,9 @@ an SSE endpoint (`GET /logstream/stream`, `text/event-stream` in through its reverse proxy — verified 2026-08-26 against `https://mempalace.jordbo.se`, where `/logstream/events`, `/logstream/stream` and `/sync/peers` all return 404 while `/mcp` serves normally. Where that is the -case, polling through the existing MCP client is the only available path, and -enabling SSE means a proxy route plus an auth decision — not an extension change. +case, polling through the existing MCP client is the only available path — which +is what the mailbox in §2 does — and enabling SSE means a proxy route plus an +auth decision, not an extension change. As with stamping, all of this is inert unless `MEMPALACE_REMOTE_URL` points at a shared palace. On a solitary palace the event tools work fine and the log diff --git a/extensions/pi/mempalace.ts b/extensions/pi/mempalace.ts index 43f0699..5f91959 100644 --- a/extensions/pi/mempalace.ts +++ b/extensions/pi/mempalace.ts @@ -36,6 +36,18 @@ * - 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 @@ -114,6 +126,22 @@ const num = (envVal: string | undefined, fallback: number): number => { return Number.isFinite(n) && n >= 0 ? n : fallback; }; +// One event as returned by `mempalace_event_list`. Every field is optional on +// purpose: this is parsed from a tool's JSON text, so it is untrusted input and +// a missing/renamed field must degrade to "not owed" rather than throw. +type LogEvent = { + id?: string; + seq?: number; + type?: string; + status?: string; + from_agent?: string; + to_agent?: string; + correlation_id?: string | null; + created_at?: string; + body?: string; + metadata?: Record | null; +}; + const sleep = (ms: number): Promise => new Promise((res) => { const t = setTimeout(res, ms); @@ -940,6 +968,201 @@ export default async function mempalaceExtension(pi: ExtensionAPI) { ); }); + // --- Automatic mailbox delivery (RFC 003 coordination log) ------------- + // + // Doubly gated, exactly like the stamper above: a device label AND a shared + // palace. On one shared palace every client reports the same origin_replica, + // so `to_agent` carries the whole distinction between machines and an + // unstamped client has no address to be reached at — there is nothing to read. + const mailboxEnabled = Boolean(stampProvenance) && process.env.MEMPALACE_MAILBOX !== "0"; + const mailboxPollMs = num(process.env.MEMPALACE_MAILBOX_POLL_MS, 300_000); + const mailboxResurfaceMs = num(process.env.MEMPALACE_MAILBOX_RESURFACE_MS, 3_600_000); + // The address is the NON-miner identity: `mine` is the only tool that stamps + // itself `miner@device`, and nobody addresses a mailbox to the miner. + const mailboxAddress = `${agentName}@${device}`; + // Only these close an obligation. `claimed` and `ready` are deliberately NOT + // terminal — that is what expresses "taken, but not finished", so picked-up + // work keeps resurfacing until it is actually closed out. + const TERMINAL_STATUS = new Set(["applied", "superseded", "failed", "blocked"]); + let lastMailboxPollAt = 0; + // event id -> when we last surfaced it. In-memory ON PURPOSE. After a restart + // this forgets, so an ask already seen can be shown again; that is the SAFE + // failure direction. A spuriously resurfacing item is visible noise a human + // corrects in one turn, whereas a suppressed unanswered ask is silent and + // permanent. Never "fix" the noise by persisting this into a decision that + // can hide an obligation. + const surfaced = new Map(); + + /** Parse an `event_list` result into a plain array. Never throws. */ + const parseEvents = (raw: unknown): LogEvent[] => { + try { + const text = extractText(raw); + if (!text) return []; + const parsed = JSON.parse(text); + const list = Array.isArray(parsed) ? parsed : parsed?.events; + return Array.isArray(list) ? (list as LogEvent[]) : []; + } catch { + // Malformed or non-JSON payload: report nothing owed rather than guessing. + return []; + } + }; + + /** + * Has one of MY events closed this candidate? + * + * Owed-ness cannot be read off `status`: event_ack APPENDS and never mutates, + * and `status` is written once, so a directed "open" keeps matching the + * mailbox query forever — answered or not. It has to be derived by joining + * the candidate to my own replies. + */ + const isAnswered = (candidate: LogEvent, mine: LogEvent[]): boolean => + mine.some((m) => { + // (c) terminal status only. + if (!TERMINAL_STATUS.has(String(m.status))) return false; + // (a) LOAD-BEARING ordering test, and the reason it must exist: without + // it one terminal reply suppresses every LATER ask on the same + // correlation_id forever — silently, permanently, and worst on exactly + // the long-running threads the correlation join is for. + // + // Compare `seq`, NEVER `created_at`. `seq` is replica-local (it equals + // origin_seq only while a single replica authors for every machine); the + // durable key once `mempalace_mesh_peers` reports real peers is `hlc`, + // which is total and causally consistent. Local-seq skew can only make an + // ANSWERED item resurface (noise, visible), whereas a timestamp + // comparison can suppress an UNANSWERED ask outright. + if (typeof m.seq !== "number" || typeof candidate.seq !== "number") return false; + if (m.seq <= candidate.seq) return false; + // (b) join on the exact ack (written for us by event_ack), else on a + // shared correlation_id — which is why correlation_id is mandatory on a + // directed open: without it there is no key to join a reply back to. + const ackOf = m.metadata?.ack_of; + if (typeof ackOf === "string" && candidate.id && ackOf === candidate.id) return true; + return Boolean(m.correlation_id && m.correlation_id === candidate.correlation_id); + }); + + /** Directed asks addressed to this device with no terminal reply from it. */ + async function deriveOwed(): Promise { + if (!mailboxEnabled || !available) return []; + try { + const [candidatesRaw, mineRaw] = await Promise.all([ + client.callTool("mempalace_event_list", { + to_agent: mailboxAddress, + status: "open", + limit: 50, + }), + client.callTool("mempalace_event_list", { from_agent: mailboxAddress, limit: 100 }), + ]); + const candidates = parseEvents(candidatesRaw).filter((c) => { + // `to_agent: ` ALSO matches `*` broadcasts, per the tool contract. + // The protocol says a broadcast owes nobody a reply, so a broadcast + // written with status="open" must not enter anyone's owed set -- + // otherwise this code contradicts the skill that documents it, and + // every machine would think it personally owed the same answer. + // This also gives the "don't broadcast an ask" anti-pattern teeth: + // broadcasting one now demonstrably reaches no owed set at all. + if (c.to_agent === "*") return false; + // You cannot owe yourself. + return c.from_agent !== mailboxAddress; + }); + if (candidates.length === 0) return []; + const mine = parseEvents(mineRaw); + return candidates.filter((c) => !isAnswered(c, mine)); + } catch { + // Fail silent and open: a mailbox read must never break a session. + return []; + } + } + + const relAge = (iso: string | undefined): string => { + if (!iso) return "age unknown"; + const then = Date.parse(iso); + if (!Number.isFinite(then)) return "age unknown"; + const mins = Math.max(0, Math.round((Date.now() - then) / 60_000)); + if (mins < 60) return `${mins}m ago`; + const hours = Math.round(mins / 60); + return hours < 48 ? `${hours}h ago` : `${Math.round(hours / 24)}d ago`; + }; + + const excerpt = (body: string | undefined, max = 200): string => { + if (typeof body !== "string") return ""; + const flat = body.replace(/\s+/g, " ").trim(); + return flat.length > max ? `${flat.slice(0, max)}\u2026` : flat; + }; + + /** One terse line per ask — this goes into a live context window. */ + const formatOwed = (items: LogEvent[]): string => + items + .map((e) => { + const head = [ + e.id ?? "(no id)", + `from ${e.from_agent ?? "?"}`, + e.type ?? "?", + e.correlation_id ? `correlation_id=${e.correlation_id}` : "NO correlation_id", + relAge(e.created_at), + ].join(" \u00b7 "); + const ex = excerpt(e.body); + return ex ? `- ${head}\n ${ex}` : `- ${head}`; + }) + .join("\n"); + + const OWED_HOWTO = + "To close one, append an event on the SAME correlation_id with a terminal status " + + "(applied/superseded/failed/blocked). An ack alone does NOT clear it, and neither " + + "does claimed/ready — those deliberately keep it visible as taken-but-unfinished."; + + // Mid-session mailbox. Tier 1 (the wake-up injection below) owns the first + // look; this exists because arrivals are bursty and correlate with our own + // activity — measured 2026-08-26, 11 of 22 events on this log landed inside a + // single 5h38m working window, so session-start-only delivery sleeps through + // most of a burst. `agent_settled` is already activity-coupled, which is why + // it is the trigger rather than a wall-clock timer. + // + // POLLING AND DELIVERY ARE DECOUPLED. We poll on a floor, but speak only when + // the owed set holds something not already surfaced. Re-announcing the same + // ask every poll is precisely how a channel teaches its reader to ignore it, + // which is the failure this whole feature exists to reverse. + pi.on("agent_settled", async () => { + if (!mailboxEnabled || !available) return; + if (!wokeUp) return; // the wake-up injection has not run yet; it does look #1 + if (Date.now() - lastMailboxPollAt < mailboxPollMs) return; + // Claim the slot BEFORE awaiting anything, so overlapping settles cannot + // stampede past the floor into concurrent polls. + lastMailboxPollAt = Date.now(); + void (async () => { + try { + const owed = await deriveOwed(); + const now = Date.now(); + const due = owed.filter((e) => { + if (!e.id) return false; + const last = surfaced.get(e.id); + return last === undefined || now - last >= mailboxResurfaceMs; + }); + if (due.length === 0) return; // silence is the correct output here + for (const e of due) if (e.id) surfaced.set(e.id, now); + pi.sendMessage( + { + customType: "mempalace-mailbox", + content: + `MemPalace logstream mailbox: ${due.length} directed ask(s) addressed to ` + + `"${mailboxAddress}" with no terminal reply from this device. ${OWED_HOWTO}\n\n` + + formatOwed(due), + display: true, + }, + // "steer" is queue-on-arrival, NOT an interrupt: at agent_settled the + // agent is idle, so this lands at the start of the next turn without + // spending an LLM call. Chosen over "nextTurn" only because if a settle + // ever races a freshly-started turn, "steer" delivers within it instead + // of waiting for the next prompt. Deliberately NO triggerTurn: waking + // the model on inbound fleet traffic is a far bigger behavioural change + // than auto-delivery and is not what was approved. + { deliverAs: "steer" }, + ); + } catch { + /* best-effort delivery: never break a turn over a mailbox read */ + } + })(); // deliberately not awaited: never stall a turn + }); + // --- Auto wake-up (mempalace skill Phase 1) --- pi.on("before_agent_start", async (_event, _ctx) => { if (wokeUp || !available) return; @@ -963,6 +1186,25 @@ export default async function mempalaceExtension(pi: ExtensionAPI) { } 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;