diff --git a/plugin/pi/README.md b/plugin/pi/README.md index 5de41059e..5eb8d4787 100644 --- a/plugin/pi/README.md +++ b/plugin/pi/README.md @@ -94,6 +94,8 @@ Restart Pi after installation, then ask Pi what it remembers about the current p | Pi extension | Captures prompts/session events, injects the Memory Protocol, and exposes compact Pi-native `mem_*` tools over the Engram HTTP server. | | MCP tools | Keeps Engram's MCP surface available through `pi-mcp-adapter` for clients and flows that use MCP directly. | +Pi-native `mem_session_end` accepts only the current Pi host session ID; a different, missing, or empty ID is refused before an end request. If a resumed conversation has a distinct persisted session ID, the matching host request ends that effective ID through the same coordination as session shutdown. A raw host-ID end requires locally confirmed registration; an uncertain registration cannot authorize it. A pending resumed ID whose end acknowledgement was lost can be reconciled without another end request only when Engram confirms that exact ID is already ended under the resolved local project. To end an independent/manual session, use a separate direct client rather than supplying its ID to the Pi-native tool. + ```text Pi events/tools -> gentle-engram extension -> ENGRAM_URL / engram serve -> SQLite Pi MCP tools -> pi-mcp-adapter -> ENGRAM_BIN / engram mcp -> SQLite diff --git a/plugin/pi/index.ts b/plugin/pi/index.ts index ad9800e1e..0450ecf33 100644 --- a/plugin/pi/index.ts +++ b/plugin/pi/index.ts @@ -1236,20 +1236,23 @@ function hasSessionRegistrationInFlight(sessionId: string): boolean { return [...sessionRegistrationsInFlight.keys()].some((key) => key.endsWith(`:${sessionId}`)); } -async function waitForSessionRegistration(sessionId: string): Promise { +async function waitForSessionRegistration(sessionId: string, propagateFailure = false): Promise { const registrations = [...sessionRegistrationsInFlight.entries()] .filter(([key]) => key.endsWith(`:${sessionId}`)) - .map(([, registration]) => registration.catch(() => undefined)); + .map(([, registration]) => propagateFailure ? registration : registration.catch(() => undefined)); await Promise.all(registrations); } -async function endRegisteredSessionOnce(sessionId: string, end: () => Promise, persistedPending = false): Promise { +async function endRegisteredSessionOnce(sessionId: string, end: () => Promise, persistedPending = false, requireConfirmedRegistration = false): Promise { const existing = sessionEndingsInFlight.get(sessionId); if (existing) return existing; const registrationWasInFlight = hasSessionRegistrationInFlight(sessionId); const ending = (async () => { - await waitForSessionRegistration(sessionId); + await waitForSessionRegistration(sessionId, requireConfirmedRegistration); + if (requireConfirmedRegistration && !registeredSessionProjects.has(sessionId)) { + throw new Error(`Cannot end Pi session ${sessionId} without confirmed local project ownership`); + } if (!registrationWasInFlight && !hasKnownSession(sessionId) && !submittedEffectiveSessions.has(sessionId) && !persistedPending) return null; try { return await end(); @@ -1535,7 +1538,7 @@ function slugifyTopicKey(params: Record): string { return slug || "memory"; } -async function callMemoryTool(toolName: string, params: Record, ctx: SessionContext, fetch: EngramFetcher = engramFetch, appendEntry?: ExtensionAPI["appendEntry"]): Promise { +async function callMemoryTool(toolName: string, params: Record, ctx: SessionContext, fetch: EngramFetcher = engramFetch, appendEntry?: ExtensionAPI["appendEntry"], transportFailure?: () => EngramTransportFailure | undefined): Promise { const sessionId = getSessionId(ctx); const requestedProject = typeof params.project === "string" && params.project ? params.project : undefined; const activeProject = requestedProject || project; @@ -1639,15 +1642,60 @@ async function callMemoryTool(toolName: string, params: Record, body: { id: params.id, project, directory: params.directory || directory || ctx.cwd }, }); case "mem_session_end": { - const endedSessionID = String(params.id); + if (typeof params.id !== "string" || !params.id || params.id !== sessionId) { + throw new Error("Pi-native session end requires the current host session ID; end independent sessions directly outside Pi-native tools"); + } + const endedSessionID = effectiveSessionID(ctx, sessionId); + const persistedPending = pendingEffectiveSession(ctx, sessionId, endedSessionID); + const owner = persistedPending ? pendingEffectiveSessionProject(ctx, sessionId, endedSessionID) : undefined; + if (persistedPending) { + const localOwner = registeredSessionProjects.get(endedSessionID) || sessionRegistrationProjects.get(endedSessionID); + if (!appendEntry || !owner || owner !== (localOwner || (project !== "unknown" ? project : undefined))) { + throw new Error("Cannot end a pending Pi session without confirmed local project ownership and entry persistence"); + } + } const pendingEnd = sessionEndingsInFlight.get(endedSessionID); if (pendingEnd) return pendingEnd; - const end = () => fetch(`/sessions/${encodeURIComponent(endedSessionID)}/end`, { - method: "POST", - body: { summary: params.summary || "" }, - }); - return endedSessionID === sessionId && (hasKnownSession(endedSessionID) || hasSessionRegistrationInFlight(endedSessionID)) - ? endRegisteredSessionOnce(endedSessionID, end) + // A persisted reservation from an earlier module graph is not local proof of + // ownership. Confirm it with the server before allowing this graph to end it. + if (persistedPending && !registeredSessionProjects.has(endedSessionID)) { + try { + await ensureSession(endedSessionID, owner, fetch, true); + } catch (error) { + const data = error instanceof EngramHttpError ? error.data as { code?: string; session_id?: string } | null : null; + if (!(error instanceof EngramHttpError && error.status === 409 + && data?.code === "session_already_ended" && data.session_id === endedSessionID + && owner === project && pendingEffectiveSession(ctx, sessionId, endedSessionID))) throw error; + appendEntry!(EFFECTIVE_SESSION_ENTRY, { + runtimeID: sessionId, effectiveID: endedSessionID, pending: false, project: owner, + }); + return { status: "already_ended", session_id: endedSessionID }; + } + } + if (endedSessionID === sessionId && !registeredSessionProjects.has(endedSessionID)) { + await waitForSessionRegistration(endedSessionID, true); + if (!registeredSessionProjects.has(endedSessionID)) { + throw new Error(`Cannot end Pi session ${endedSessionID} without confirmed local registration`); + } + } + const end = async () => { + const result = await fetch(`/sessions/${encodeURIComponent(endedSessionID)}/end`, { + method: "POST", + body: { summary: params.summary || "" }, + }); + if (persistedPending && !transportFailure?.()) appendEntry!(EFFECTIVE_SESSION_ENTRY, { + runtimeID: sessionId, effectiveID: endedSessionID, pending: false, project: owner, + }); + return result; + }; + return endedSessionID !== sessionId || hasKnownSession(endedSessionID) || hasSessionRegistrationInFlight(endedSessionID) + ? endRegisteredSessionOnce(endedSessionID, async () => { + if (endedSessionID !== sessionId && (!registeredSessionProjects.has(endedSessionID) + || (persistedPending && registeredSessionProjects.get(endedSessionID) !== owner))) { + throw new Error(`Cannot end Pi session ${endedSessionID} without confirmed local project ownership`); + } + return end(); + }, persistedPending, endedSessionID !== sessionId) : end(); } case "mem_current_project": { @@ -1747,7 +1795,7 @@ async function executeMemoryTool(toolName: string, params: Record { diff --git a/plugin/pi/test/native-tool-contract.test.mjs b/plugin/pi/test/native-tool-contract.test.mjs index 3d0decaf8..de7818b75 100644 --- a/plugin/pi/test/native-tool-contract.test.mjs +++ b/plugin/pi/test/native-tool-contract.test.mjs @@ -75,6 +75,32 @@ function recordingFetch(routes) { return { calls, fetchStub }; } +test("Pi-native mem_session_end refuses an existing foreign session without an end request", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const { calls, fetchStub } = recordingFetch([ + { method: "GET", path: "/health", body: { status: "ok" } }, + { method: "GET", path: "/project/current", body: { project: "paidosdep" } }, + { method: "POST", path: "/sessions/foreign-existing-session/end", body: { status: "ended" } }, + ]); + globalThis.fetch = fetchStub; + try { + await withPluginSandbox("engram-pi-contract-", async ({ sandbox }) => { + const { registeredTools } = await loadPluginHarness(sandbox); + const result = await registeredTools.get("mem_session_end").execute( + "foreign-end", { id: "foreign-existing-session" }, undefined, undefined, runtimeContext("host-session"), + ); + assert.equal(result.isError, true); + assert.equal(calls.filter((call) => call.method === "POST" && call.path.endsWith("/end")).length, 0); + }); + } finally { + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + test("registered Pi-native mem_save_prompt persists through the Engram /prompts endpoint", async () => { const originalFetch = globalThis.fetch; const originalUrl = process.env.ENGRAM_URL; @@ -1407,6 +1433,170 @@ test("resumed ended conversation registers a distinct persistent identity before } }); +test("Pi-native explicit end targets the registered resumed identity, not the host ID", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const calls = []; + const entries = []; + const ctx = runtimeContext("explicit-resume"); + ctx.sessionManager.getBranch = () => entries; + globalThis.fetch = async (url, init = {}) => { + const path = new URL(url).pathname; + const body = init.body ? JSON.parse(init.body) : undefined; + calls.push({ method: init.method ?? "GET", path, body }); + if (path === "/project/current") return new Response(JSON.stringify({ project: "resume-project" })); + if (path === "/sessions" && body.id === "explicit-resume") { + return new Response(JSON.stringify({ code: "session_already_ended" }), { status: 409 }); + } + return new Response(JSON.stringify({ status: "ok" })); + }; + try { + await withPluginSandbox("engram-pi-explicit-resume-end-", async ({ sandbox }) => { + const { registeredTools, eventHandlers } = await loadPluginHarness(sandbox, (customType, data) => entries.push({ type: "custom", customType, data })); + const save = await registeredTools.get("mem_save").execute("save", { title: "resume", content: "resume" }, undefined, undefined, ctx); + assert.equal(save.isError, undefined, JSON.stringify(save)); + const effectiveID = calls.filter(({ path }) => path === "/sessions").at(-1).body.id; + assert.match(effectiveID, /^explicit-resume:resume:/); + const endTool = registeredTools.get("mem_session_end"); + for (const id of [effectiveID, "foreign", "", undefined]) { + const refused = await endTool.execute("refused", { id }, undefined, undefined, ctx); + assert.equal(refused.isError, true); + } + assert.equal(calls.filter(({ path }) => path.endsWith("/end")).length, 0); + const ended = await endTool.execute("host-end", { id: "explicit-resume" }, undefined, undefined, ctx); + assert.equal(ended.isError, undefined, JSON.stringify(ended)); + assert.deepEqual(calls.filter(({ path }) => path.endsWith("/end")).map(({ path }) => path), + [`/sessions/${encodeURIComponent(effectiveID)}/end`]); + assert.equal(entries.at(-1).data.pending, false, "explicit end must clear its persisted reservation"); + await eventHandlers.get("session_shutdown")({}, ctx); + assert.equal(calls.filter(({ path }) => path.endsWith("/end")).length, 1, "shutdown must not repeat explicit end"); + }); + } finally { + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + +test("ambiguous resumed explicit end leaves its pending reservation intact; JSON null success clears it", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const entries = []; + const ctx = runtimeContext("ambiguous-end"); + const effectiveID = "ambiguous-end:resume:owned"; + ctx.sessionManager.getBranch = () => entries; + entries.push({ type: "custom", customType: "engram-effective-session", data: { + runtimeID: "ambiguous-end", effectiveID, pending: true, project: "resume-project", + } }); + let endAttempts = 0; + globalThis.fetch = async (url, init = {}) => { + const path = new URL(url).pathname; + if (path === "/project/current") return new Response(JSON.stringify({ project: "resume-project" })); + if (path === "/sessions") return new Response(JSON.stringify({ status: "created" })); + if (path === "/observations") return new Response(JSON.stringify({ id: 1 })); + if (path === `/sessions/${encodeURIComponent(effectiveID)}/end`) { + endAttempts++; + if (endAttempts === 1) throw new Error("response lost after dispatch"); + return new Response("null", { status: 200 }); + } + throw new Error(`unexpected request: ${path} ${init.method}`); + }; + try { + await withPluginSandbox("engram-pi-ambiguous-explicit-end-", async ({ sandbox }) => { + const { registeredTools } = await loadPluginHarness(sandbox, (customType, data) => entries.push({ type: "custom", customType, data })); + const save = await registeredTools.get("mem_save").execute("save", { title: "owned", content: "owned" }, undefined, undefined, ctx); + assert.equal(save.isError, undefined, JSON.stringify(save)); + const end = registeredTools.get("mem_session_end"); + const ambiguous = await end.execute("end-unknown", { id: "ambiguous-end" }, undefined, undefined, ctx); + assert.equal(ambiguous.isError, true); + assert.equal(entries.at(-1).data.pending, true, "unknown delivery must not clear the persisted marker"); + const renewed = await registeredTools.get("mem_save").execute("renew", { title: "owned", content: "owned" }, undefined, undefined, ctx); + assert.equal(renewed.isError, undefined, JSON.stringify(renewed)); + const confirmed = await end.execute("end-confirmed", { id: "ambiguous-end" }, undefined, undefined, ctx); + assert.equal(confirmed.isError, undefined, JSON.stringify(confirmed)); + assert.equal(entries.at(-1).data.pending, false, "HTTP JSON null is confirmed delivery"); + assert.equal(endAttempts, 2); + }); + } finally { + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + +test("Pi-native explicit end in a new graph confirms prior pending ownership before closing", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const runtimeID = "prior-graph-end"; + const effectiveID = `${runtimeID}:resume:owned`; + const entries = [{ type: "custom", customType: "engram-effective-session", data: { + runtimeID, effectiveID, pending: true, project: "local-project", + } }]; + const ctx = runtimeContext(runtimeID); + ctx.sessionManager.getBranch = () => entries; + const calls = []; + let conflict = false; + globalThis.fetch = async (url, init = {}) => { + const path = new URL(url).pathname; + calls.push(path); + if (path === "/project/current") return new Response(JSON.stringify({ project: "local-project" })); + if (path === "/sessions") return conflict + ? new Response(JSON.stringify({ code: "session_project_conflict", session_id: effectiveID, + owner_project: "foreign-project", requested_project: "local-project" }), { status: 409 }) + : new Response(JSON.stringify({ status: "created" })); + if (path === `/sessions/${encodeURIComponent(effectiveID)}/end`) return new Response(JSON.stringify({ status: "ended" })); + throw new Error(`unexpected request: ${path} ${init.method}`); + }; + try { + for (const shouldConflict of [true, false]) { + conflict = shouldConflict; + await withPluginSandbox("engram-pi-prior-graph-end-", async ({ sandbox }) => { + const graph = await loadPluginHarness(sandbox, (customType, data) => entries.push({ type: "custom", customType, data })); + const before = calls.length; + const result = await graph.registeredTools.get("mem_session_end").execute("end", { id: runtimeID }, undefined, undefined, ctx); + assert.deepEqual(calls.slice(before).filter((path) => path === "/sessions"), ["/sessions"], "new graph must confirm server ownership"); + assert.equal(result.isError, shouldConflict ? true : undefined, JSON.stringify(result)); + assert.equal(calls.slice(before).filter((path) => path.endsWith("/end")).length, shouldConflict ? 0 : 1); + assert.equal(entries.at(-1).data.pending, shouldConflict); + }); + } + } finally { + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + +test("Pi-native explicit end refuses a foreign-owned pending reservation", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const { calls, fetchStub } = recordingFetch([ + { method: "GET", path: "/project/current", body: { project: "local-project" } }, + ]); + globalThis.fetch = fetchStub; + const ctx = runtimeContext("foreign-reservation"); + ctx.sessionManager.getBranch = () => [{ type: "custom", customType: "engram-effective-session", + data: { runtimeID: "foreign-reservation", effectiveID: "foreign-reservation:resume:other", pending: true, project: "other-project" } }]; + try { + await withPluginSandbox("engram-pi-foreign-explicit-end-", async ({ sandbox }) => { + const { registeredTools, eventHandlers } = await loadPluginHarness(sandbox); + await eventHandlers.get("session_start")({}, ctx); + const result = await registeredTools.get("mem_session_end").execute("foreign-reservation-end", + { id: "foreign-reservation" }, undefined, undefined, ctx); + assert.equal(result.isError, true); + assert.equal(calls.filter(({ path }) => path.endsWith("/end")).length, 0); + }); + } finally { + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + test("ambiguous fresh registration reuses its submitted identity and ends it on shutdown", async () => { const originalFetch = globalThis.fetch; const originalUrl = process.env.ENGRAM_URL; @@ -2197,9 +2387,104 @@ test("Pi session shutdown serializes end delivery and waits for registration", a uncertainRegistration.resolve(); const [writeResult, endResult] = await Promise.all([write, explicitEnd]); assert.equal(writeResult.isError, true, "the registration outcome is uncertain"); - assert.equal(endResult.isError, undefined, "explicit end must still attempt delivery"); - assert.deepEqual(uncertainEndCalls, [{ method: "POST", body: { summary: "" } }]); + assert.equal(endResult.isError, true, "uncertain registration cannot authorize explicit end"); + assert.deepEqual(uncertainEndCalls, []); + }); + } finally { + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + +test("Pi-native explicit end awaits a pending effective registration conflict before sending end", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const gate = deferred(); + const started = deferred(); + const calls = []; + const entries = []; + const runtimeID = "effective-conflict"; + const effectiveID = `${runtimeID}:resume:reserved`; + const ctx = runtimeContext(runtimeID); + ctx.sessionManager.getBranch = () => entries; + entries.push({ type: "custom", customType: "engram-effective-session", data: { + runtimeID, effectiveID, pending: true, project: "local-project", + } }); + globalThis.fetch = async (url, init = {}) => { + const path = new URL(url).pathname; + calls.push(path); + if (path === "/project/current") return new Response(JSON.stringify({ project: "local-project" })); + if (path === "/sessions") { + started.resolve(); + await gate.promise; + return new Response(JSON.stringify({ code: "session_project_conflict", session_id: effectiveID, + owner_project: "foreign-project", requested_project: "local-project" }), { status: 409 }); + } + if (path.endsWith("/end")) return new Response(JSON.stringify({ status: "ended" })); + throw new Error(`unexpected request: ${path} ${init.method}`); + }; + try { + await withPluginSandbox("engram-pi-effective-conflict-end-", async ({ sandbox }) => { + const { registeredTools } = await loadPluginHarness(sandbox, (customType, data) => entries.push({ type: "custom", customType, data })); + const save = registeredTools.get("mem_save").execute("save", { title: "conflict", content: "conflict" }, undefined, undefined, ctx); + await started.promise; + const end = registeredTools.get("mem_session_end").execute("end", { id: runtimeID }, undefined, undefined, ctx); + assert.equal(calls.filter((path) => path.endsWith("/end")).length, 0); + gate.resolve(); + const [saveResult, endResult] = await Promise.all([save, end]); + assert.equal(saveResult.isError, true); + assert.equal(endResult.isError, true, "registration conflict must propagate to explicit end"); + assert.match(endResult.content[0].text, /belongs to Engram project foreign-project/); + assert.equal(calls.filter((path) => path.endsWith("/end")).length, 0); + assert.equal(entries.at(-1).customType, "engram-rejected-effective-session"); }); + } finally { + gate.resolve(); + globalThis.fetch = originalFetch; + if (originalUrl === undefined) delete process.env.ENGRAM_URL; + else process.env.ENGRAM_URL = originalUrl; + } +}); + +test("explicit end reconciles only an owned pending effective ID with matching already-ended evidence", async () => { + const originalFetch = globalThis.fetch; + const originalUrl = process.env.ENGRAM_URL; + process.env.ENGRAM_URL = "http://127.0.0.1:17437"; + const runtimeID = "pending-explicit"; + const effectiveID = `${runtimeID}:resume:reserved`; + const calls = []; + let responseID = effectiveID; + let owner = "local-project"; + globalThis.fetch = async (url) => { + const path = new URL(url).pathname; + calls.push(path); + if (path === "/project/current") return new Response(JSON.stringify({ project: "local-project" })); + if (path === "/sessions") return new Response(JSON.stringify({ code: "session_already_ended", session_id: responseID }), { status: 409 }); + if (path.endsWith("/end")) return new Response(JSON.stringify({ status: "ended" })); + throw new Error(`unexpected request: ${path}`); + }; + try { + for (const [response, markerOwner, success] of [ + [effectiveID, "local-project", true], ["other-id", "local-project", false], + [undefined, "local-project", false], [effectiveID, "foreign-project", false], + ]) { + responseID = response; + owner = markerOwner; + const entries = [{ type: "custom", customType: "engram-effective-session", data: { + runtimeID, effectiveID, pending: true, project: owner, + } }]; + const ctx = runtimeContext(runtimeID); + ctx.sessionManager.getBranch = () => entries; + await withPluginSandbox("engram-pi-pending-explicit-", async ({ sandbox }) => { + const { registeredTools } = await loadPluginHarness(sandbox, (customType, data) => entries.push({ type: "custom", customType, data })); + const result = await registeredTools.get("mem_session_end").execute("end", { id: runtimeID }, undefined, undefined, ctx); + assert.equal(result.isError, success ? undefined : true); + assert.equal(entries.at(-1).data.pending, !success); + assert.equal(calls.filter((path) => path.endsWith("/end") || path.includes(":resume:")).length, 0); + }); + } } finally { globalThis.fetch = originalFetch; if (originalUrl === undefined) delete process.env.ENGRAM_URL;