From 707059f72fa91fcb5cc44c6026e8714ab6dac535 Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Mon, 28 Sep 2026 23:57:30 +0200 Subject: [PATCH 1/9] fix(opencode): add V2 setup entrypoint to the plugin OpenCode V2 rejects plugins without a default definition exposing setup or effect, so the Engram adapter never loaded there. Keep the V1 server entry unchanged and add a thin setup that maps V2 session, tool and event hooks onto the existing handlers. --- internal/setup/plugins/opencode/engram.ts | 140 ++++++++++- plugin/opencode/engram.ts | 140 ++++++++++- plugin/opencode/engram.v2.test.mjs | 293 ++++++++++++++++++++++ 3 files changed, 571 insertions(+), 2 deletions(-) create mode 100644 plugin/opencode/engram.v2.test.mjs diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index 2f8dcd30f..cdb02edb6 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -860,4 +860,142 @@ export const Engram: Plugin = async (ctx) => { } } -export default { id: "engram", server: Engram } +// ─── OpenCode V2 Adapter ───────────────────────────────────────────────────── +// OpenCode V2 calls `setup(ctx)` instead of `server`. It exposes hooks through +// per-domain registrations and session lifecycle through an event stream, so +// this adapter translates them onto the V1 hooks above and adds no behavior. +// Types are declared structurally: V1 hosts may not ship `@opencode/plugin`. +// +// Not exported by name on purpose: older V1 loaders call every exported +// function as a plugin factory. + +type V2SystemPart = { type: "text"; text: string } +type V2Registration = { dispose: () => Promise } +type V2Hook = (name: string, callback: (input: any) => Promise | void) => Promise +type V2Context = { + location: { directory: string; project?: { id?: string } } + event: { subscribe: (options?: { signal?: AbortSignal }) => AsyncIterable } + session: { get: (input: { sessionID: string }) => Promise; hook: V2Hook } + tool: { hook: V2Hook } +} + +// V1 hooks append to the last system string; V2 carries system text parts. +function appendSystemText(system: V2SystemPart[], text: string): void { + const last = system[system.length - 1] + if (last) system[system.length - 1] = { ...last, text: `${last.text}\n\n${text}` } + else system.push({ type: "text", text }) +} + +async function withSystemStrings(system: V2SystemPart[], run: (texts: string[]) => Promise): Promise { + const texts = system.map((part) => part.text) + await run(texts) + texts.forEach((text, index) => { + if (index >= system.length) system.push({ type: "text", text }) + else if (text !== system[index].text) system[index] = { ...system[index], text } + }) +} + +function v2ToolResultText(result: any): string { + const text = Array.isArray(result?.content) + ? result.content.filter((part: any) => part?.type === "text").map((part: any) => part.text ?? "").join("\n") + : "" + return text || (result?.output === undefined ? "" : JSON.stringify(result.output)) +} + +// V2 session events carry `data.sessionID`; V1 hooks expect `properties.info.id`. +function v1SessionEvent(event: any, directory: string): any { + const data = event?.data + if (typeof data?.sessionID !== "string") return undefined + if (event.type === "session.created") { + // The V2 server is shared across locations; V1 only saw its own instance. + if (data.location?.directory && data.location.directory !== directory) return undefined + return { type: event.type, properties: { info: { id: data.sessionID, parentID: data.parentID, projectID: data.projectID } } } + } + if (event.type === "session.deleted") { + return { type: event.type, properties: { info: { id: data.sessionID } } } + } + return undefined +} + +async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { + const hooks: Record = await Engram({ + directory: ctx.location.directory, + project: { id: ctx.location.project?.id }, + client: { + session: { + // V1 SDK results carry `{ data, error }`; the V2 client throws instead. + async get({ path }: { path: { id: string } }) { + try { + return { data: await ctx.session.get({ sessionID: path.id }) } + } catch (error) { + return { error } + } + }, + }, + }, + } as any) + + const abort = new AbortController() + const registrations: V2Registration[] = [] + let listening: Promise = Promise.resolve() + const cleanup = async () => { + abort.abort() + await Promise.all(registrations.map((registration) => registration.dispose())) + await listening + await hooks.dispose?.() + } + + try { + registrations.push(await ctx.session.hook("prompt", async (prompt) => { + await hooks["chat.message"]({ sessionID: prompt.sessionID }, { + message: {}, + parts: [{ type: "text", text: prompt.prompt?.text ?? "" }], + }) + })) + + registrations.push(await ctx.session.hook("context", async (request) => { + await withSystemStrings(request.system, (system) => + hooks["experimental.chat.system.transform"]({ sessionID: request.sessionID, model: request.model }, { system })) + })) + + registrations.push(await ctx.session.hook("compaction", async (request) => { + const context: string[] = [] + await hooks["experimental.session.compacting"]({ sessionID: request.sessionID }, { context }) + if (context.length > 0) appendSystemText(request.system, context.join("\n\n")) + })) + + registrations.push(await ctx.tool.hook("execute.before", async (call) => { + const args = call.input && typeof call.input === "object" ? call.input : {} + await hooks["tool.execute.before"]({ tool: call.tool, sessionID: call.sessionID, callID: call.id }, { args }) + if (args !== call.input && Object.keys(args).length > 0) call.input = args + })) + + registrations.push(await ctx.tool.hook("execute.after", async (call) => { + // V2 renamed the V1 `Task` delegation tool to `subagent`. + const tool = call.tool === "subagent" ? "Task" : call.tool + const output = call.status === "completed" ? v2ToolResultText(call.result) : "" + await hooks["tool.execute.after"]({ tool, sessionID: call.sessionID, callID: call.id }, output) + })) + + const events = ctx.event.subscribe({ signal: abort.signal }) + listening = (async () => { + try { + for await (const event of events) { + if (abort.signal.aborted) break + const translated = v1SessionEvent(event, ctx.location.directory) + if (translated) await hooks.event({ event: translated }).catch(() => {}) + } + } catch { + // Stream failure loses lifecycle events; hooks still bind sessions lazily. + } + })() + } catch (cause) { + await cleanup() + throw cause + } + + return cleanup +} + +// V1 (1.18.29+) calls server(); V2 calls setup(). +export default { id: "engram", server: Engram, setup: setupEngramV2 } diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index 2f8dcd30f..cdb02edb6 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -860,4 +860,142 @@ export const Engram: Plugin = async (ctx) => { } } -export default { id: "engram", server: Engram } +// ─── OpenCode V2 Adapter ───────────────────────────────────────────────────── +// OpenCode V2 calls `setup(ctx)` instead of `server`. It exposes hooks through +// per-domain registrations and session lifecycle through an event stream, so +// this adapter translates them onto the V1 hooks above and adds no behavior. +// Types are declared structurally: V1 hosts may not ship `@opencode/plugin`. +// +// Not exported by name on purpose: older V1 loaders call every exported +// function as a plugin factory. + +type V2SystemPart = { type: "text"; text: string } +type V2Registration = { dispose: () => Promise } +type V2Hook = (name: string, callback: (input: any) => Promise | void) => Promise +type V2Context = { + location: { directory: string; project?: { id?: string } } + event: { subscribe: (options?: { signal?: AbortSignal }) => AsyncIterable } + session: { get: (input: { sessionID: string }) => Promise; hook: V2Hook } + tool: { hook: V2Hook } +} + +// V1 hooks append to the last system string; V2 carries system text parts. +function appendSystemText(system: V2SystemPart[], text: string): void { + const last = system[system.length - 1] + if (last) system[system.length - 1] = { ...last, text: `${last.text}\n\n${text}` } + else system.push({ type: "text", text }) +} + +async function withSystemStrings(system: V2SystemPart[], run: (texts: string[]) => Promise): Promise { + const texts = system.map((part) => part.text) + await run(texts) + texts.forEach((text, index) => { + if (index >= system.length) system.push({ type: "text", text }) + else if (text !== system[index].text) system[index] = { ...system[index], text } + }) +} + +function v2ToolResultText(result: any): string { + const text = Array.isArray(result?.content) + ? result.content.filter((part: any) => part?.type === "text").map((part: any) => part.text ?? "").join("\n") + : "" + return text || (result?.output === undefined ? "" : JSON.stringify(result.output)) +} + +// V2 session events carry `data.sessionID`; V1 hooks expect `properties.info.id`. +function v1SessionEvent(event: any, directory: string): any { + const data = event?.data + if (typeof data?.sessionID !== "string") return undefined + if (event.type === "session.created") { + // The V2 server is shared across locations; V1 only saw its own instance. + if (data.location?.directory && data.location.directory !== directory) return undefined + return { type: event.type, properties: { info: { id: data.sessionID, parentID: data.parentID, projectID: data.projectID } } } + } + if (event.type === "session.deleted") { + return { type: event.type, properties: { info: { id: data.sessionID } } } + } + return undefined +} + +async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { + const hooks: Record = await Engram({ + directory: ctx.location.directory, + project: { id: ctx.location.project?.id }, + client: { + session: { + // V1 SDK results carry `{ data, error }`; the V2 client throws instead. + async get({ path }: { path: { id: string } }) { + try { + return { data: await ctx.session.get({ sessionID: path.id }) } + } catch (error) { + return { error } + } + }, + }, + }, + } as any) + + const abort = new AbortController() + const registrations: V2Registration[] = [] + let listening: Promise = Promise.resolve() + const cleanup = async () => { + abort.abort() + await Promise.all(registrations.map((registration) => registration.dispose())) + await listening + await hooks.dispose?.() + } + + try { + registrations.push(await ctx.session.hook("prompt", async (prompt) => { + await hooks["chat.message"]({ sessionID: prompt.sessionID }, { + message: {}, + parts: [{ type: "text", text: prompt.prompt?.text ?? "" }], + }) + })) + + registrations.push(await ctx.session.hook("context", async (request) => { + await withSystemStrings(request.system, (system) => + hooks["experimental.chat.system.transform"]({ sessionID: request.sessionID, model: request.model }, { system })) + })) + + registrations.push(await ctx.session.hook("compaction", async (request) => { + const context: string[] = [] + await hooks["experimental.session.compacting"]({ sessionID: request.sessionID }, { context }) + if (context.length > 0) appendSystemText(request.system, context.join("\n\n")) + })) + + registrations.push(await ctx.tool.hook("execute.before", async (call) => { + const args = call.input && typeof call.input === "object" ? call.input : {} + await hooks["tool.execute.before"]({ tool: call.tool, sessionID: call.sessionID, callID: call.id }, { args }) + if (args !== call.input && Object.keys(args).length > 0) call.input = args + })) + + registrations.push(await ctx.tool.hook("execute.after", async (call) => { + // V2 renamed the V1 `Task` delegation tool to `subagent`. + const tool = call.tool === "subagent" ? "Task" : call.tool + const output = call.status === "completed" ? v2ToolResultText(call.result) : "" + await hooks["tool.execute.after"]({ tool, sessionID: call.sessionID, callID: call.id }, output) + })) + + const events = ctx.event.subscribe({ signal: abort.signal }) + listening = (async () => { + try { + for await (const event of events) { + if (abort.signal.aborted) break + const translated = v1SessionEvent(event, ctx.location.directory) + if (translated) await hooks.event({ event: translated }).catch(() => {}) + } + } catch { + // Stream failure loses lifecycle events; hooks still bind sessions lazily. + } + })() + } catch (cause) { + await cleanup() + throw cause + } + + return cleanup +} + +// V1 (1.18.29+) calls server(); V2 calls setup(). +export default { id: "engram", server: Engram, setup: setupEngramV2 } diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs new file mode 100644 index 000000000..e5b545f4d --- /dev/null +++ b/plugin/opencode/engram.v2.test.mjs @@ -0,0 +1,293 @@ +import assert from "node:assert/strict" +import { readFileSync } from "node:fs" +import { createRequire, syncBuiltinESMExports } from "node:module" +import { test } from "node:test" + +const require = createRequire(import.meta.url) +const childProcess = require("node:child_process") +const fs = require("node:fs") + +const source = readFileSync(new URL("./engram.ts", import.meta.url), "utf8") + +const DIRECTORY = "/work/engram" +const PROJECT_ID = "project-1" +const INSTANCE_ID = "00000000000000000000000000000000" +let runtimeImport = 0 + +function httpResponse(data) { + return { ok: true, async json() { return data } } +} + +// Minimal OpenCode V2 event stream. emit() resolves once the plugin asks for +// the next event, which proves the previous one was fully handled. +function eventStream() { + const queued = [] + let waiting + let pulled + let closed = false + const subscriptions = [] + return { + subscriptions, + subscribe(options) { + subscriptions.push(options) + options?.signal?.addEventListener("abort", () => { + closed = true + waiting?.({ value: undefined, done: true }) + }) + return { + [Symbol.asyncIterator]() { + return { + next() { + pulled?.() + pulled = undefined + if (queued.length > 0) return Promise.resolve({ value: queued.shift(), done: false }) + if (closed) return Promise.resolve({ value: undefined, done: true }) + return new Promise((resolve) => { waiting = resolve }) + }, + async return() { + closed = true + return { value: undefined, done: true } + }, + } + }, + } + }, + emit(event) { + const handled = new Promise((resolve) => { pulled = resolve }) + if (waiting) { + const resolve = waiting + waiting = undefined + resolve({ value: event, done: false }) + } else { + queued.push(event) + } + return handled + }, + } +} + +async function setupV2(t, { sessions = new Map() } = {}) { + const originalFetch = globalThis.fetch + const originalBun = globalThis.Bun + const originalEngramURL = process.env.ENGRAM_URL + const originalSpawnSync = childProcess.spawnSync + const originalSpawn = childProcess.spawn + const originalExistsSync = fs.existsSync + delete globalThis.Bun + delete process.env.ENGRAM_URL + childProcess.spawnSync = () => ({ status: 0, stdout: `${INSTANCE_ID}\n` }) + childProcess.spawn = () => ({ on() { return this }, unref() {} }) + fs.existsSync = () => false + syncBuiltinESMExports() + + const requests = [] + globalThis.fetch = async (url, init) => { + const path = new URL(url).pathname + if (path === "/health") return httpResponse({ status: "ok", instance_id: INSTANCE_ID }) + const body = init?.body ? JSON.parse(init.body) : undefined + requests.push({ path, method: init?.method, body }) + if (path === "/project/current") return httpResponse({ project: "engram", project_source: "git_remote" }) + if (path === "/sessions") return httpResponse({ id: body.id, status: "created" }) + if (path === "/context/compaction") return httpResponse({ context: "previous session context" }) + return httpResponse({}) + } + + t.after(() => { + globalThis.fetch = originalFetch + globalThis.Bun = originalBun + if (originalEngramURL === undefined) delete process.env.ENGRAM_URL + else process.env.ENGRAM_URL = originalEngramURL + childProcess.spawnSync = originalSpawnSync + childProcess.spawn = originalSpawn + fs.existsSync = originalExistsSync + syncBuiltinESMExports() + }) + + const hooks = new Map() + const disposedHooks = [] + const sessionGetIDs = [] + const register = (domain) => async (name, callback) => { + hooks.set(`${domain}.${name}`, callback) + return { dispose: async () => { disposedHooks.push(`${domain}.${name}`) } } + } + const events = eventStream() + const ctx = { + location: { directory: DIRECTORY, project: { id: PROJECT_ID, directory: DIRECTORY, canonical: DIRECTORY } }, + event: { subscribe: (options) => events.subscribe(options) }, + session: { + hook: register("session"), + async get({ sessionID }) { + sessionGetIDs.push(sessionID) + const info = sessions.get(sessionID) + if (!info) throw new Error(`session ${sessionID} not found`) + return info + }, + }, + tool: { hook: register("tool") }, + } + + runtimeImport += 1 + const module = await import(new URL(`./engram.ts?v2-runtime=${runtimeImport}`, import.meta.url).href) + const cleanup = await module.default.setup(ctx) + return { + module, + cleanup, + hooks, + disposedHooks, + requests, + sessionGetIDs, + events, + created: (sessionID, parentID) => events.emit({ + type: "session.created", + data: { sessionID, projectID: PROJECT_ID, location: { directory: DIRECTORY }, ...(parentID ? { parentID } : {}) }, + }), + deleted: (sessionID) => events.emit({ type: "session.deleted", data: { sessionID } }), + posts: (path) => requests.filter((request) => request.method === "POST" && request.path === path), + } +} + +function sessionInfo(id, parentID) { + return { id, projectID: PROJECT_ID, ...(parentID ? { parentID } : {}) } +} + +test("default export serves V1 through server and V2 through setup", async () => { + const module = await import(new URL("./engram.ts?v2-shape", import.meta.url).href) + assert.equal(module.default.id, "engram") + assert.strictEqual(module.default.server, module.Engram) + assert.equal(typeof module.default.setup, "function") + assert.doesNotMatch(source, /^import\s+(?!type\b)[^\n]*from\s+"@opencode(-ai)?\/plugin"/m, "V1 hosts may lack the V2 SDK") +}) + +test("V2 setup registers session, tool, and event hooks and cleans them up", async (t) => { + const runtime = await setupV2(t) + assert.deepEqual([...runtime.hooks.keys()].sort(), [ + "session.compaction", + "session.context", + "session.prompt", + "tool.execute.after", + "tool.execute.before", + ]) + assert.equal(runtime.events.subscriptions.length, 1) + assert.equal(typeof runtime.cleanup, "function") + + await runtime.created("ses_root") + await runtime.cleanup() + + assert.equal(runtime.disposedHooks.length, 5) + assert.equal(runtime.events.subscriptions[0].signal.aborted, true) + assert.equal(runtime.posts("/sessions/ses_root/end").length, 1, "cleanup ends registered sessions") +}) + +test("V2 session.created binds root sessions but never child sessions", async (t) => { + const runtime = await setupV2(t) + await runtime.created("ses_root") + await runtime.created("ses_child", "ses_root") + + assert.deepEqual(runtime.posts("/sessions").map(({ body }) => body), [ + { id: "ses_root", project: "engram", directory: DIRECTORY }, + ]) + + await runtime.deleted("ses_root") + assert.equal(runtime.posts("/sessions/ses_root/end").length, 1) +}) + +test("V2 ignores other locations and leaves unrelated tool input untouched", async (t) => { + const runtime = await setupV2(t) + await runtime.events.emit({ + type: "session.created", + data: { sessionID: "ses_elsewhere", projectID: PROJECT_ID, location: { directory: "/work/other" } }, + }) + assert.equal(runtime.posts("/sessions").length, 0) + + const call = { tool: "bash", sessionID: "ses_root", input: undefined } + await runtime.hooks.get("tool.execute.before")(call) + assert.equal(call.input, undefined) +}) + +test("V2 setup releases earlier registrations when a later one fails", async (t) => { + const runtime = await setupV2(t) + await runtime.cleanup() + const ctx = { + location: { directory: DIRECTORY, project: { id: PROJECT_ID } }, + event: { subscribe: () => { throw new Error("unexpected subscribe") } }, + session: { + get: async () => { throw new Error("unexpected get") }, + hook: async (name) => ({ dispose: async () => { disposed.push(name) } }), + }, + tool: { hook: async () => { throw new Error("tool hooks unavailable") } }, + } + const disposed = [] + await assert.rejects(runtime.module.default.setup(ctx), /tool hooks unavailable/) + assert.deepEqual(disposed, ["prompt", "context", "compaction"]) +}) + +test("V2 prompt hook captures user prompts for the authoritative session", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + await runtime.hooks.get("session.prompt")({ + sessionID: "ses_root", + messageID: "msg_1", + prompt: { text: "Please remember the token decision" }, + delivery: "immediate", + }) + + assert.deepEqual(runtime.posts("/prompts").map(({ body }) => body), [ + { session_id: "ses_root", content: "Please remember the [REDACTED] decision", project: "engram" }, + ]) +}) + +test("V2 tool hooks bind Engram writes to the root session and capture subagent output", async (t) => { + const runtime = await setupV2(t, { + sessions: new Map([["ses_root", sessionInfo("ses_root")], ["ses_child", sessionInfo("ses_child", "ses_root")]]), + }) + const call = { tool: "engram_mem_save", sessionID: "ses_child", agent: "build", messageID: "msg_1", id: "call_1", input: { title: "x" } } + await runtime.hooks.get("tool.execute.before")(call) + assert.equal(call.input.session_id, "ses_root") + assert.deepEqual(runtime.sessionGetIDs, ["ses_child", "ses_root"]) + + const output = "Subagent finished: the auth middleware now validates JWT expiry before routing." + await runtime.hooks.get("tool.execute.after")({ + tool: "subagent", + sessionID: "ses_root", + agent: "build", + messageID: "msg_2", + id: "call_2", + input: {}, + status: "completed", + result: { content: [{ type: "text", text: output }] }, + }) + assert.deepEqual(runtime.posts("/observations/passive").map(({ body }) => body), [ + { session_id: "ses_root", content: output, project: "engram", source: "task-complete" }, + ]) +}) + +test("V2 tool hook rejects Engram writes without an authoritative session", async (t) => { + const runtime = await setupV2(t) + const call = { tool: "engram_mem_save", sessionID: "ses_missing", input: {} } + await assert.rejects(runtime.hooks.get("tool.execute.before")(call), /authoritative OpenCode runtime session/) + assert.equal(call.input.session_id, undefined) +}) + +test("V2 context hook appends memory instructions to the last system part", async (t) => { + const runtime = await setupV2(t) + const request = { sessionID: "ses_root", system: [{ type: "text", text: "base" }], messages: [] } + await runtime.hooks.get("session.context")(request) + assert.equal(request.system.length, 1) + assert.match(request.system[0].text, /^base\n\n## Engram Persistent Memory/) + + const empty = { sessionID: "ses_root", system: [], messages: [] } + await runtime.hooks.get("session.context")(empty) + assert.equal(empty.system.length, 1) + assert.equal(empty.system[0].type, "text") + assert.match(empty.system[0].text, /^## Engram Persistent Memory/) +}) + +test("V2 compaction hook injects session context and the summary instruction", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + const compaction = { sessionID: "ses_root", system: [{ type: "text", text: "summarize" }], messages: [] } + await runtime.hooks.get("session.compaction")(compaction) + + assert.equal(compaction.system.length, 1) + assert.match(compaction.system[0].text, /^summarize\n\nprevious session context\n\nCRITICAL INSTRUCTION FOR COMPACTED SUMMARY/) + assert.match(compaction.system[0].text, /Use project: 'engram'/) + assert.equal(runtime.requests.filter(({ path }) => path === "/context/compaction").length, 1) +}) From 2dc243bfa00a1cf5dd8ed5b3daa2d0afe0fc6249 Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 00:33:20 +0200 Subject: [PATCH 2/9] fix(opencode): keep the V2 event stream alive and run plugin tests in CI setupEngramV2 now re-subscribes to the event stream with bounded backoff when it ends or throws, stops on abort, and isolates per-event handler failures so one bad event cannot stop session lifecycle handling. CI runs the OpenCode plugin Node tests, which previously never ran. --- .github/workflows/ci.yml | 4 + internal/setup/plugins/opencode/engram.ts | 44 +++++++-- plugin/opencode/engram.ts | 44 +++++++-- plugin/opencode/engram.v2.test.mjs | 104 +++++++++++++++++++--- 4 files changed, 170 insertions(+), 26 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8ad32fc0d..948eb820d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -194,6 +194,10 @@ jobs: working-directory: plugin/pi run: npm test + # Covers the dual-major (1.x server / 2.x setup) OpenCode adapter. + - name: Run OpenCode plugin tests + run: node --test plugin/opencode/engram.test.mjs plugin/opencode/engram.v2.test.mjs + perf-ratchet: name: Performance Ratchet if: github.event_name == 'push' && github.ref == 'refs/heads/main' diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index cdb02edb6..13ecf3050 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -917,6 +917,24 @@ function v1SessionEvent(event: any, directory: string): any { return undefined } +// The V2 event stream ends or throws when the server restarts; reconnect with +// a bounded doubling delay that resets once events flow again. +const V2_EVENT_RETRY_MIN_MS = 50 +const V2_EVENT_RETRY_MAX_MS = 5000 + +function delayUnlessAborted(ms: number, signal: AbortSignal): Promise { + return new Promise((resolve) => { + if (signal.aborted) return resolve() + const done = () => { + clearTimeout(timer) + signal.removeEventListener("abort", done) + resolve() + } + const timer = setTimeout(done, ms) + signal.addEventListener("abort", done, { once: true }) + }) +} + async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { const hooks: Record = await Engram({ directory: ctx.location.directory, @@ -977,16 +995,26 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { await hooks["tool.execute.after"]({ tool, sessionID: call.sessionID, callID: call.id }, output) })) - const events = ctx.event.subscribe({ signal: abort.signal }) listening = (async () => { - try { - for await (const event of events) { - if (abort.signal.aborted) break - const translated = v1SessionEvent(event, ctx.location.directory) - if (translated) await hooks.event({ event: translated }).catch(() => {}) + let retryMs = V2_EVENT_RETRY_MIN_MS + while (!abort.signal.aborted) { + try { + for await (const event of ctx.event.subscribe({ signal: abort.signal })) { + if (abort.signal.aborted) break + retryMs = V2_EVENT_RETRY_MIN_MS + try { + const translated = v1SessionEvent(event, ctx.location.directory) + if (translated) await hooks.event({ event: translated }) + } catch { + // One failing event must not stop lifecycle tracking. + } + } + } catch { + // Events missed while disconnected are lost; hooks still bind sessions lazily. } - } catch { - // Stream failure loses lifecycle events; hooks still bind sessions lazily. + if (abort.signal.aborted) break + await delayUnlessAborted(retryMs, abort.signal) + retryMs = Math.min(retryMs * 2, V2_EVENT_RETRY_MAX_MS) } })() } catch (cause) { diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index cdb02edb6..13ecf3050 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -917,6 +917,24 @@ function v1SessionEvent(event: any, directory: string): any { return undefined } +// The V2 event stream ends or throws when the server restarts; reconnect with +// a bounded doubling delay that resets once events flow again. +const V2_EVENT_RETRY_MIN_MS = 50 +const V2_EVENT_RETRY_MAX_MS = 5000 + +function delayUnlessAborted(ms: number, signal: AbortSignal): Promise { + return new Promise((resolve) => { + if (signal.aborted) return resolve() + const done = () => { + clearTimeout(timer) + signal.removeEventListener("abort", done) + resolve() + } + const timer = setTimeout(done, ms) + signal.addEventListener("abort", done, { once: true }) + }) +} + async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { const hooks: Record = await Engram({ directory: ctx.location.directory, @@ -977,16 +995,26 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { await hooks["tool.execute.after"]({ tool, sessionID: call.sessionID, callID: call.id }, output) })) - const events = ctx.event.subscribe({ signal: abort.signal }) listening = (async () => { - try { - for await (const event of events) { - if (abort.signal.aborted) break - const translated = v1SessionEvent(event, ctx.location.directory) - if (translated) await hooks.event({ event: translated }).catch(() => {}) + let retryMs = V2_EVENT_RETRY_MIN_MS + while (!abort.signal.aborted) { + try { + for await (const event of ctx.event.subscribe({ signal: abort.signal })) { + if (abort.signal.aborted) break + retryMs = V2_EVENT_RETRY_MIN_MS + try { + const translated = v1SessionEvent(event, ctx.location.directory) + if (translated) await hooks.event({ event: translated }) + } catch { + // One failing event must not stop lifecycle tracking. + } + } + } catch { + // Events missed while disconnected are lost; hooks still bind sessions lazily. } - } catch { - // Stream failure loses lifecycle events; hooks still bind sessions lazily. + if (abort.signal.aborted) break + await delayUnlessAborted(retryMs, abort.signal) + retryMs = Math.min(retryMs * 2, V2_EVENT_RETRY_MAX_MS) } })() } catch (cause) { diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs index e5b545f4d..d3b0e4117 100644 --- a/plugin/opencode/engram.v2.test.mjs +++ b/plugin/opencode/engram.v2.test.mjs @@ -19,20 +19,31 @@ function httpResponse(data) { } // Minimal OpenCode V2 event stream. emit() resolves once the plugin asks for -// the next event, which proves the previous one was fully handled. +// the next event, which proves the previous one was fully handled. end() and +// fail() interrupt only the current subscription, like a server restart. function eventStream() { const queued = [] let waiting + let waitingReject let pulled let closed = false + let interrupt const subscriptions = [] + const settle = (result) => { + const resolve = waiting + const reject = waitingReject + waiting = undefined + waitingReject = undefined + if (result instanceof Error) reject(result) + else resolve(result) + } return { subscriptions, subscribe(options) { subscriptions.push(options) options?.signal?.addEventListener("abort", () => { closed = true - waiting?.({ value: undefined, done: true }) + if (waiting) settle({ value: undefined, done: true }) }) return { [Symbol.asyncIterator]() { @@ -40,9 +51,17 @@ function eventStream() { next() { pulled?.() pulled = undefined + if (interrupt) { + const result = interrupt + interrupt = undefined + return result instanceof Error ? Promise.reject(result) : Promise.resolve(result) + } if (queued.length > 0) return Promise.resolve({ value: queued.shift(), done: false }) if (closed) return Promise.resolve({ value: undefined, done: true }) - return new Promise((resolve) => { waiting = resolve }) + return new Promise((resolve, reject) => { + waiting = resolve + waitingReject = reject + }) }, async return() { closed = true @@ -54,15 +73,37 @@ function eventStream() { }, emit(event) { const handled = new Promise((resolve) => { pulled = resolve }) - if (waiting) { - const resolve = waiting - waiting = undefined - resolve({ value: event, done: false }) - } else { - queued.push(event) - } + if (waiting) settle({ value: event, done: false }) + else queued.push(event) return handled }, + end() { + const result = { value: undefined, done: true } + if (waiting) settle(result) + else interrupt = result + }, + fail(error) { + if (waiting) settle(error) + else interrupt = error + }, + closeForever() { + closed = true + if (waiting) settle({ value: undefined, done: true }) + }, + } +} + +function withTimeout(promise, message, ms = 1000) { + let timer + const timeout = new Promise((_, reject) => { timer = setTimeout(() => reject(new Error(message)), ms) }) + return Promise.race([promise, timeout]).finally(() => clearTimeout(timer)) +} + +async function waitFor(condition, message, ms = 1000) { + const deadline = Date.now() + ms + while (!condition()) { + if (Date.now() > deadline) throw new Error(message) + await new Promise((resolve) => setTimeout(resolve, 5)) } } @@ -291,3 +332,46 @@ test("V2 compaction hook injects session context and the summary instruction", a assert.match(compaction.system[0].text, /Use project: 'engram'/) assert.equal(runtime.requests.filter(({ path }) => path === "/context/compaction").length, 1) }) + +test("V2 re-subscribes when the event stream ends", async (t) => { + const runtime = await setupV2(t) + runtime.events.end() + await waitFor(() => runtime.events.subscriptions.length === 2, "plugin did not re-subscribe after the stream ended") + + await withTimeout(runtime.created("ses_root"), "event after re-subscription was not handled") + assert.equal(runtime.posts("/sessions").length, 1) + await runtime.cleanup() +}) + +test("V2 re-subscribes when the event stream throws", async (t) => { + const runtime = await setupV2(t) + runtime.events.fail(new Error("stream reset")) + await waitFor(() => runtime.events.subscriptions.length === 2, "plugin did not re-subscribe after the stream failed") + + await withTimeout(runtime.created("ses_root"), "event after re-subscription was not handled") + assert.equal(runtime.posts("/sessions").length, 1) + await runtime.cleanup() +}) + +test("V2 keeps listening when handling one event throws", async (t) => { + const runtime = await setupV2(t) + const poisoned = { type: "session.created", get data() { throw new Error("malformed event") } } + await withTimeout(runtime.events.emit(poisoned), "loop stopped pulling after a failing event") + await withTimeout(runtime.created("ses_root"), "event after a failing event was not handled") + + assert.equal(runtime.posts("/sessions").length, 1) + assert.equal(runtime.events.subscriptions.length, 1, "a failing event must not drop the subscription") + await runtime.cleanup() +}) + +test("V2 backs off between re-subscriptions and stops after cleanup", async (t) => { + const runtime = await setupV2(t) + runtime.events.closeForever() + await new Promise((resolve) => setTimeout(resolve, 300)) + const attempts = runtime.events.subscriptions.length + assert.ok(attempts >= 2 && attempts <= 6, `expected bounded re-subscription attempts, got ${attempts}`) + + await withTimeout(runtime.cleanup(), "cleanup waited on the reconnect backoff") + await new Promise((resolve) => setTimeout(resolve, 150)) + assert.equal(runtime.events.subscriptions.length, attempts, "no re-subscription after cleanup") +}) From a508dd5ae953eca2e62bf95d9fe74ecab204f323 Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 00:33:35 +0200 Subject: [PATCH 3/9] docs(opencode): document the dual-major plugin entrypoints --- docs/PLUGINS.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/docs/PLUGINS.md b/docs/PLUGINS.md index f96aa5037..297bd9edb 100644 --- a/docs/PLUGINS.md +++ b/docs/PLUGINS.md @@ -44,6 +44,8 @@ engram setup opencode The plugin auto-starts the HTTP server if it's not already running — no manual `engram serve` needed. +The same `engram.ts` supports OpenCode 1.x (1.18.29+) and 2.x. Its default export provides a V1 `server` entry and a V2 `setup` entry; the V2 entry maps session, prompt, context, compaction, and tool hooks onto the same handlers, so both majors share one behavior. OpenCode 2.x still accepts the MCP entry written by `engram setup opencode`, so the same command covers both majors. V2 has no `session.updated` event, so a late root-to-child reclassification relies on `parentID` at creation and on session lookups in hooks. + > **Local model compatibility:** The plugin works with all models, including local ones served via llama.cpp, Ollama, or similar. The Memory Protocol is concatenated into the existing system prompt (not added as a separate system message), so models with strict Jinja templates (Qwen, Mistral/Ministral) work correctly. ### What the Plugin Does From c4eef17a9551504933bb38c3bd84d19d60c2e21c Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 00:55:55 +0200 Subject: [PATCH 4/9] fix(opencode): capture V2 tool results returned as string content --- internal/setup/plugins/opencode/engram.ts | 5 ++++- plugin/opencode/engram.ts | 5 ++++- plugin/opencode/engram.v2.test.mjs | 18 ++++++++++++++++++ 3 files changed, 26 insertions(+), 2 deletions(-) diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index 13ecf3050..039ac8542 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -895,8 +895,11 @@ async function withSystemStrings(system: V2SystemPart[], run: (texts: string[]) }) } +// V2 `Tool.Result.content` is `string | Content[]`; `subagent` returns a string. function v2ToolResultText(result: any): string { - const text = Array.isArray(result?.content) + const text = typeof result?.content === "string" + ? result.content + : Array.isArray(result?.content) ? result.content.filter((part: any) => part?.type === "text").map((part: any) => part.text ?? "").join("\n") : "" return text || (result?.output === undefined ? "" : JSON.stringify(result.output)) diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index 13ecf3050..039ac8542 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -895,8 +895,11 @@ async function withSystemStrings(system: V2SystemPart[], run: (texts: string[]) }) } +// V2 `Tool.Result.content` is `string | Content[]`; `subagent` returns a string. function v2ToolResultText(result: any): string { - const text = Array.isArray(result?.content) + const text = typeof result?.content === "string" + ? result.content + : Array.isArray(result?.content) ? result.content.filter((part: any) => part?.type === "text").map((part: any) => part.text ?? "").join("\n") : "" return text || (result?.output === undefined ? "" : JSON.stringify(result.output)) diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs index d3b0e4117..b7350b11e 100644 --- a/plugin/opencode/engram.v2.test.mjs +++ b/plugin/opencode/engram.v2.test.mjs @@ -301,6 +301,24 @@ test("V2 tool hooks bind Engram writes to the root session and capture subagent ]) }) +test("V2 tool hook captures subagent output returned as string content", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + const output = "Subagent finished: the rendered transcript arrives as a plain string." + await runtime.hooks.get("tool.execute.after")({ + tool: "subagent", + sessionID: "ses_root", + agent: "build", + messageID: "msg_3", + id: "call_3", + input: {}, + status: "completed", + result: { content: output, output: { ignored: true } }, + }) + assert.deepEqual(runtime.posts("/observations/passive").map(({ body }) => body), [ + { session_id: "ses_root", content: output, project: "engram", source: "task-complete" }, + ]) +}) + test("V2 tool hook rejects Engram writes without an authoritative session", async (t) => { const runtime = await setupV2(t) const call = { tool: "engram_mem_save", sessionID: "ses_missing", input: {} } From a8c57e1eea753e965a1579cee49a531312b30b48 Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 08:04:20 +0200 Subject: [PATCH 5/9] fix(opencode): keep plain string V2 tool output unquoted --- internal/setup/plugins/opencode/engram.ts | 4 +++- plugin/opencode/engram.ts | 4 +++- plugin/opencode/engram.v2.test.mjs | 18 ++++++++++++++++++ 3 files changed, 24 insertions(+), 2 deletions(-) diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index 039ac8542..ccc0af388 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -902,7 +902,9 @@ function v2ToolResultText(result: any): string { : Array.isArray(result?.content) ? result.content.filter((part: any) => part?.type === "text").map((part: any) => part.text ?? "").join("\n") : "" - return text || (result?.output === undefined ? "" : JSON.stringify(result.output)) + if (text) return text + if (typeof result?.output === "string") return result.output + return result?.output === undefined ? "" : JSON.stringify(result.output) } // V2 session events carry `data.sessionID`; V1 hooks expect `properties.info.id`. diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index 039ac8542..ccc0af388 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -902,7 +902,9 @@ function v2ToolResultText(result: any): string { : Array.isArray(result?.content) ? result.content.filter((part: any) => part?.type === "text").map((part: any) => part.text ?? "").join("\n") : "" - return text || (result?.output === undefined ? "" : JSON.stringify(result.output)) + if (text) return text + if (typeof result?.output === "string") return result.output + return result?.output === undefined ? "" : JSON.stringify(result.output) } // V2 session events carry `data.sessionID`; V1 hooks expect `properties.info.id`. diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs index b7350b11e..0ddd4ece8 100644 --- a/plugin/opencode/engram.v2.test.mjs +++ b/plugin/opencode/engram.v2.test.mjs @@ -319,6 +319,24 @@ test("V2 tool hook captures subagent output returned as string content", async ( ]) }) +test("V2 tool hook captures a plain string output without JSON quoting", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + const output = "Subagent finished: only a plain string output was returned by the tool." + await runtime.hooks.get("tool.execute.after")({ + tool: "subagent", + sessionID: "ses_root", + agent: "build", + messageID: "msg_4", + id: "call_4", + input: {}, + status: "completed", + result: { content: "", output }, + }) + assert.deepEqual(runtime.posts("/observations/passive").map(({ body }) => body), [ + { session_id: "ses_root", content: output, project: "engram", source: "task-complete" }, + ]) +}) + test("V2 tool hook rejects Engram writes without an authoritative session", async (t) => { const runtime = await setupV2(t) const call = { tool: "engram_mem_save", sessionID: "ses_missing", input: {} } From 696eaeb5e4872fca58832d6fe5fa8a2ce80122ef Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 08:29:42 +0200 Subject: [PATCH 6/9] fix(opencode): redact private blocks before truncating prompts --- internal/setup/plugins/opencode/engram.ts | 4 +++- plugin/opencode/engram.test.mjs | 11 +++++++++++ plugin/opencode/engram.ts | 4 +++- plugin/opencode/engram.v2.test.mjs | 15 +++++++++++++++ 4 files changed, 32 insertions(+), 2 deletions(-) diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index ccc0af388..f5bffade9 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -669,7 +669,9 @@ export const Engram: Plugin = async (ctx) => { method: "POST", body: { session_id: sessionId, - content: stripPrivateTags(truncate(finalContent, 2000)), + // Redact before truncating: a block straddling the + // limit would otherwise lose its closing tag and leak. + content: truncate(stripPrivateTags(finalContent), 2000), project, }, }) diff --git a/plugin/opencode/engram.test.mjs b/plugin/opencode/engram.test.mjs index b616de4b8..874af2628 100644 --- a/plugin/opencode/engram.test.mjs +++ b/plugin/opencode/engram.test.mjs @@ -937,6 +937,17 @@ test("write tool hook revalidates leaf and ancestor ownership after registration } }) +test("chat.message redacts a private block that straddles the truncation limit", async (t) => { + const runtime = await createRuntime(t) + const text = `${"a".repeat(1980)}PIN=42 trailing` + await runtime.chat({ sessionID: "runtime" }, { message: {}, parts: [{ type: "text", text }] }) + + const prompts = runtime.requests.filter(({ path }) => path === "/prompts") + assert.equal(prompts.length, 1) + assert.equal(JSON.stringify(prompts[0].body).includes("PIN=42"), false) + assert.equal(prompts[0].body.content.includes("[REDACTED]"), true) +}) + test("chat.message resolves an unobserved child and skips its prompt", async (t) => { const runtime = await createRuntime(t, { sessionGet: sdkLookup(CHILD_SESSIONS) }) await runtime.chat( diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index ccc0af388..f5bffade9 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -669,7 +669,9 @@ export const Engram: Plugin = async (ctx) => { method: "POST", body: { session_id: sessionId, - content: stripPrivateTags(truncate(finalContent, 2000)), + // Redact before truncating: a block straddling the + // limit would otherwise lose its closing tag and leak. + content: truncate(stripPrivateTags(finalContent), 2000), project, }, }) diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs index 0ddd4ece8..36fd4c8b2 100644 --- a/plugin/opencode/engram.v2.test.mjs +++ b/plugin/opencode/engram.v2.test.mjs @@ -276,6 +276,21 @@ test("V2 prompt hook captures user prompts for the authoritative session", async ]) }) +test("V2 prompt hook redacts a private block that straddles the truncation limit", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + await runtime.hooks.get("session.prompt")({ + sessionID: "ses_root", + messageID: "msg_1", + prompt: { text: `${"a".repeat(1980)}PIN=42 trailing` }, + delivery: "immediate", + }) + + const prompts = runtime.posts("/prompts").map(({ body }) => body) + assert.equal(prompts.length, 1) + assert.equal(JSON.stringify(prompts[0]).includes("PIN=42"), false) + assert.equal(prompts[0].content.includes("[REDACTED]"), true) +}) + test("V2 tool hooks bind Engram writes to the root session and capture subagent output", async (t) => { const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")], ["ses_child", sessionInfo("ses_child", "ses_root")]]), From 1483c9fe2cc62248b5b050eedf9ce2c4fccfe8e7 Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 08:50:49 +0200 Subject: [PATCH 7/9] feat(opencode): capture V2 prompts with durable inbox identity Capture OpenCode 2.x prompts from user items of session.inbox.enqueued and send the inboxID as source_inbox_id, so a server with prompt inbox identity support (#1464) treats replays as no-ops, keeps equal-text items distinct, and refuses deleted identities. V1 chat.message and V2 share one capture helper; the identity-less V2 prompt hook is no longer used for capture. Adds a real-server regression for replay, restart, and delete. --- docs/PLUGINS.md | 2 +- internal/setup/plugins/opencode/engram.ts | 86 +++++--- plugin/opencode/engram.ts | 86 +++++--- .../opencode/engram.v2.real-server.test.mjs | 200 ++++++++++++++++++ plugin/opencode/engram.v2.test.mjs | 104 +++++++-- 5 files changed, 397 insertions(+), 81 deletions(-) create mode 100644 plugin/opencode/engram.v2.real-server.test.mjs diff --git a/docs/PLUGINS.md b/docs/PLUGINS.md index 5bc9a61af..5e5c111bc 100644 --- a/docs/PLUGINS.md +++ b/docs/PLUGINS.md @@ -45,7 +45,7 @@ engram setup opencode The plugin auto-starts the HTTP server if it's not already running — no manual `engram serve` needed. -The same `engram.ts` supports OpenCode 1.x (1.18.29+) and 2.x. Its default export provides a V1 `server` entry and a V2 `setup` entry; the V2 entry maps session, prompt, context, compaction, and tool hooks onto the same handlers, so both majors share one behavior. OpenCode 2.x still accepts the MCP entry written by `engram setup opencode`, so the same command covers both majors. V2 has no `session.updated` event, so a late root-to-child reclassification relies on `parentID` at creation and on session lookups in hooks. +The same `engram.ts` supports OpenCode 1.x (1.18.29+) and 2.x. Its default export provides a V1 `server` entry and a V2 `setup` entry; the V2 entry maps session, context, compaction, and tool hooks onto the same handlers, so both majors share one behavior. On V2, user prompts are captured from the durable `session.inbox.enqueued` event (user items only) and sent with the item's `inboxID` as `source_inbox_id`, so a replayed admission never creates a second prompt, distinct items with identical text stay distinct, and a deleted prompt is not resurrected (the server answers `409`, which the plugin drops). This idempotency requires a server with prompt inbox identity support (#1464); older servers ignore the field and store every admission. OpenCode 2.x still accepts the MCP entry written by `engram setup opencode`, so the same command covers both majors. V2 has no `session.updated` event, so a late root-to-child reclassification relies on `parentID` at creation and on session lookups in hooks. > **Local model compatibility:** The plugin works with all models, including local ones served via llama.cpp, Ollama, or similar. The Memory Protocol is concatenated into the existing system prompt (not added as a separate system message), so models with strict Jinja templates (Qwen, Mistral/Ministral) work correctly. diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index f5bffade9..4e9395e0e 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -270,6 +270,11 @@ export function shouldNudgeForObservations( // ─── Plugin Export ─────────────────────────────────────────────────────────── +// Hidden hooks-object key through which the V2 adapter reaches the shared prompt +// capture. A symbol keeps it out of the V1 hook names OpenCode enumerates. +const CAPTURE_PROMPT = Symbol("engram.capturePrompt") +type CapturePrompt = (sourceSessionID: string, content: string, sourceInboxID?: string) => Promise + export const Engram: Plugin = async (ctx) => { let project = "unknown" let projectResolutionError = "" @@ -541,6 +546,34 @@ export const Engram: Plugin = async (ctx) => { return await registration && !invalidSessions.has(sessionId) && !closeRequestedSessions.has(sessionId) } + /** + * Capture one user prompt for its authoritative root session. A durable + * `sourceInboxID` lets the server treat replays as no-ops and refuse deleted + * identities (HTTP 409); engramFetch drops that response like any other failure. + */ + const capturePrompt: CapturePrompt = async (sourceSessionID, content, sourceInboxID) => { + const sessionId = await resolveAuthoritativeSessionID(sourceSessionID) + // Skip child prompts even when ownership was discovered through the SDK. + if (!sessionId || subAgentSessions.has(sourceSessionID)) return + + // Only capture non-trivial prompts (>10 chars) + if (content.length <= 10) return + const registered = await ensureSession(sessionId, true) + const confirmedSessionID = await resolveAuthoritativeSessionID(sourceSessionID) + if (!registered || confirmedSessionID !== sessionId) return + await engramFetch("/prompts", { + method: "POST", + body: { + session_id: sessionId, + // Redact before truncating: a block straddling the + // limit would otherwise lose its closing tag and leak. + content: truncate(stripPrivateTags(content), 2000), + project, + ...(sourceInboxID ? { source_inbox_id: sourceInboxID } : {}), + }, + }) + } + // Try to start engram server if not running try { const expectedID = CONFIGURED_ENGRAM_URL ? "" : localInstanceID() @@ -579,6 +612,8 @@ export const Engram: Plugin = async (ctx) => { } return { + [CAPTURE_PROMPT]: capturePrompt, + dispose: async () => { disposed = true if (!localReady) return @@ -642,10 +677,6 @@ export const Engram: Plugin = async (ctx) => { // output.parts contains TextPart[] with the actual message text. "chat.message": async (input, output) => { - const sessionId = await resolveAuthoritativeSessionID(input.sessionID) - // Skip child prompts even when ownership was discovered through the SDK. - if (!sessionId || subAgentSessions.has(input.sessionID)) return - // Extract text from parts (type:"text") const content = output.parts .filter((p) => p.type === "text") @@ -658,24 +689,7 @@ export const Engram: Plugin = async (ctx) => { ? `${output.message.summary.title ?? ""}\n${output.message.summary.body ?? ""}`.trim() : "" - const finalContent = content || fallback - - // Only capture non-trivial prompts (>10 chars) - if (finalContent.length > 10) { - const registered = await ensureSession(sessionId, true) - const confirmedSessionID = await resolveAuthoritativeSessionID(input.sessionID) - if (!registered || confirmedSessionID !== sessionId) return - await engramFetch("/prompts", { - method: "POST", - body: { - session_id: sessionId, - // Redact before truncating: a block straddling the - // limit would otherwise lose its closing tag and leak. - content: truncate(stripPrivateTags(finalContent), 2000), - project, - }, - }) - } + await capturePrompt(input.sessionID, content || fallback) }, // ─── Tool Execution Hook ───────────────────────────────────── @@ -924,6 +938,21 @@ function v1SessionEvent(event: any, directory: string): any { return undefined } +// V2 admits each human prompt as a durable `user` inbox item. Its inboxID is the +// prompt's identity: replays reuse it, distinct items with equal text do not. +function v2InboxPrompt(event: any, directory: string): { sessionID: string; inboxID: string; text: string } | undefined { + if (event?.type !== "session.inbox.enqueued") return undefined + // The envelope location is optional; the authoritative session lookup still + // rejects sessions outside this instance's project. + if (event.location?.directory && event.location.directory !== directory) return undefined + const data = event.data + const item = data?.item + if (typeof data?.sessionID !== "string" || !data.sessionID) return undefined + if (typeof data.inboxID !== "string" || !data.inboxID) return undefined + if (item?.type !== "user" || typeof item.payload?.text !== "string") return undefined + return { sessionID: data.sessionID, inboxID: data.inboxID, text: item.payload.text.trim() } +} + // The V2 event stream ends or throws when the server restarts; reconnect with // a bounded doubling delay that resets once events flow again. const V2_EVENT_RETRY_MIN_MS = 50 @@ -943,7 +972,7 @@ function delayUnlessAborted(ms: number, signal: AbortSignal): Promise { } async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { - const hooks: Record = await Engram({ + const hooks: Record = await Engram({ directory: ctx.location.directory, project: { id: ctx.location.project?.id }, client: { @@ -971,13 +1000,8 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { } try { - registrations.push(await ctx.session.hook("prompt", async (prompt) => { - await hooks["chat.message"]({ sessionID: prompt.sessionID }, { - message: {}, - parts: [{ type: "text", text: prompt.prompt?.text ?? "" }], - }) - })) - + // No `prompt` hook: it lacks a durable identity, so prompts are captured + // from `session.inbox.enqueued` below instead. registrations.push(await ctx.session.hook("context", async (request) => { await withSystemStrings(request.system, (system) => hooks["experimental.chat.system.transform"]({ sessionID: request.sessionID, model: request.model }, { system })) @@ -1010,6 +1034,8 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { if (abort.signal.aborted) break retryMs = V2_EVENT_RETRY_MIN_MS try { + const prompt = v2InboxPrompt(event, ctx.location.directory) + if (prompt) await (hooks[CAPTURE_PROMPT] as CapturePrompt)(prompt.sessionID, prompt.text, prompt.inboxID) const translated = v1SessionEvent(event, ctx.location.directory) if (translated) await hooks.event({ event: translated }) } catch { diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index f5bffade9..4e9395e0e 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -270,6 +270,11 @@ export function shouldNudgeForObservations( // ─── Plugin Export ─────────────────────────────────────────────────────────── +// Hidden hooks-object key through which the V2 adapter reaches the shared prompt +// capture. A symbol keeps it out of the V1 hook names OpenCode enumerates. +const CAPTURE_PROMPT = Symbol("engram.capturePrompt") +type CapturePrompt = (sourceSessionID: string, content: string, sourceInboxID?: string) => Promise + export const Engram: Plugin = async (ctx) => { let project = "unknown" let projectResolutionError = "" @@ -541,6 +546,34 @@ export const Engram: Plugin = async (ctx) => { return await registration && !invalidSessions.has(sessionId) && !closeRequestedSessions.has(sessionId) } + /** + * Capture one user prompt for its authoritative root session. A durable + * `sourceInboxID` lets the server treat replays as no-ops and refuse deleted + * identities (HTTP 409); engramFetch drops that response like any other failure. + */ + const capturePrompt: CapturePrompt = async (sourceSessionID, content, sourceInboxID) => { + const sessionId = await resolveAuthoritativeSessionID(sourceSessionID) + // Skip child prompts even when ownership was discovered through the SDK. + if (!sessionId || subAgentSessions.has(sourceSessionID)) return + + // Only capture non-trivial prompts (>10 chars) + if (content.length <= 10) return + const registered = await ensureSession(sessionId, true) + const confirmedSessionID = await resolveAuthoritativeSessionID(sourceSessionID) + if (!registered || confirmedSessionID !== sessionId) return + await engramFetch("/prompts", { + method: "POST", + body: { + session_id: sessionId, + // Redact before truncating: a block straddling the + // limit would otherwise lose its closing tag and leak. + content: truncate(stripPrivateTags(content), 2000), + project, + ...(sourceInboxID ? { source_inbox_id: sourceInboxID } : {}), + }, + }) + } + // Try to start engram server if not running try { const expectedID = CONFIGURED_ENGRAM_URL ? "" : localInstanceID() @@ -579,6 +612,8 @@ export const Engram: Plugin = async (ctx) => { } return { + [CAPTURE_PROMPT]: capturePrompt, + dispose: async () => { disposed = true if (!localReady) return @@ -642,10 +677,6 @@ export const Engram: Plugin = async (ctx) => { // output.parts contains TextPart[] with the actual message text. "chat.message": async (input, output) => { - const sessionId = await resolveAuthoritativeSessionID(input.sessionID) - // Skip child prompts even when ownership was discovered through the SDK. - if (!sessionId || subAgentSessions.has(input.sessionID)) return - // Extract text from parts (type:"text") const content = output.parts .filter((p) => p.type === "text") @@ -658,24 +689,7 @@ export const Engram: Plugin = async (ctx) => { ? `${output.message.summary.title ?? ""}\n${output.message.summary.body ?? ""}`.trim() : "" - const finalContent = content || fallback - - // Only capture non-trivial prompts (>10 chars) - if (finalContent.length > 10) { - const registered = await ensureSession(sessionId, true) - const confirmedSessionID = await resolveAuthoritativeSessionID(input.sessionID) - if (!registered || confirmedSessionID !== sessionId) return - await engramFetch("/prompts", { - method: "POST", - body: { - session_id: sessionId, - // Redact before truncating: a block straddling the - // limit would otherwise lose its closing tag and leak. - content: truncate(stripPrivateTags(finalContent), 2000), - project, - }, - }) - } + await capturePrompt(input.sessionID, content || fallback) }, // ─── Tool Execution Hook ───────────────────────────────────── @@ -924,6 +938,21 @@ function v1SessionEvent(event: any, directory: string): any { return undefined } +// V2 admits each human prompt as a durable `user` inbox item. Its inboxID is the +// prompt's identity: replays reuse it, distinct items with equal text do not. +function v2InboxPrompt(event: any, directory: string): { sessionID: string; inboxID: string; text: string } | undefined { + if (event?.type !== "session.inbox.enqueued") return undefined + // The envelope location is optional; the authoritative session lookup still + // rejects sessions outside this instance's project. + if (event.location?.directory && event.location.directory !== directory) return undefined + const data = event.data + const item = data?.item + if (typeof data?.sessionID !== "string" || !data.sessionID) return undefined + if (typeof data.inboxID !== "string" || !data.inboxID) return undefined + if (item?.type !== "user" || typeof item.payload?.text !== "string") return undefined + return { sessionID: data.sessionID, inboxID: data.inboxID, text: item.payload.text.trim() } +} + // The V2 event stream ends or throws when the server restarts; reconnect with // a bounded doubling delay that resets once events flow again. const V2_EVENT_RETRY_MIN_MS = 50 @@ -943,7 +972,7 @@ function delayUnlessAborted(ms: number, signal: AbortSignal): Promise { } async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { - const hooks: Record = await Engram({ + const hooks: Record = await Engram({ directory: ctx.location.directory, project: { id: ctx.location.project?.id }, client: { @@ -971,13 +1000,8 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { } try { - registrations.push(await ctx.session.hook("prompt", async (prompt) => { - await hooks["chat.message"]({ sessionID: prompt.sessionID }, { - message: {}, - parts: [{ type: "text", text: prompt.prompt?.text ?? "" }], - }) - })) - + // No `prompt` hook: it lacks a durable identity, so prompts are captured + // from `session.inbox.enqueued` below instead. registrations.push(await ctx.session.hook("context", async (request) => { await withSystemStrings(request.system, (system) => hooks["experimental.chat.system.transform"]({ sessionID: request.sessionID, model: request.model }, { system })) @@ -1010,6 +1034,8 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { if (abort.signal.aborted) break retryMs = V2_EVENT_RETRY_MIN_MS try { + const prompt = v2InboxPrompt(event, ctx.location.directory) + if (prompt) await (hooks[CAPTURE_PROMPT] as CapturePrompt)(prompt.sessionID, prompt.text, prompt.inboxID) const translated = v1SessionEvent(event, ctx.location.directory) if (translated) await hooks.event({ event: translated }) } catch { diff --git a/plugin/opencode/engram.v2.real-server.test.mjs b/plugin/opencode/engram.v2.real-server.test.mjs new file mode 100644 index 000000000..9d7c84000 --- /dev/null +++ b/plugin/opencode/engram.v2.real-server.test.mjs @@ -0,0 +1,200 @@ +// OpenCode V2 durable prompt identity against a real Engram HTTP server. +// Requires server support for `source_inbox_id` on POST /prompts (#1464). +import assert from "node:assert/strict" +import { mkdtempSync, rmSync } from "node:fs" +import { tmpdir } from "node:os" +import { join } from "node:path" +import { createInterface } from "node:readline" +import { spawn, spawnSync } from "node:child_process" +import { test } from "node:test" + +const PROJECT_ID = "project-1" +const PLUGIN_URL = "http://engram.invalid" + +async function bounded(promise, message, ms = 5000) { + let timer + try { + return await Promise.race([promise, new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(message)), ms) + })]) + } finally { clearTimeout(timer) } +} + +// Isolated real server on a persistent store directory; stop() closes stdin, +// which is the bridge's shutdown signal. +async function startServer(executable, dir) { + const child = spawn(executable, [join(dir, "store")], { + env: { + ...process.env, + HOME: dir, + ENGRAM_DATA_DIR: join(dir, "data"), + ENGRAM_CLOUD_AUTOSYNC: "0", + ENGRAM_HTTP_TOKEN: "", + }, + stdio: ["pipe", "pipe", "pipe"], + }) + let stderr = "" + child.stderr.setEncoding("utf8").on("data", (chunk) => { stderr += chunk }) + const lines = createInterface({ input: child.stdout }) + const exited = new Promise((resolve) => child.once("exit", resolve)) + const stop = async () => { + lines.close() + if (child.exitCode !== null || child.signalCode !== null) return + child.stdin.end() + try { await bounded(exited, `server shutdown timed out: ${stderr}`) } + catch (error) { + child.kill() + await bounded(exited, `server kill timed out: ${stderr}`) + throw error + } + } + try { + const url = await bounded(new Promise((resolve, reject) => { + lines.once("line", resolve) + child.once("error", reject) + child.once("exit", (code) => reject(new Error(`server exited ${code}: ${stderr}`))) + }), `server startup timed out: ${stderr}`, 60000) + return { url, stop } + } catch (error) { + await stop() + throw error + } +} + +// Each emit resolves once the plugin pulls the next event, which proves the +// previous one was fully handled. +function eventStream() { + const queued = [] + let waiting + let pulled + return { + subscribe({ signal } = {}) { + signal?.addEventListener("abort", () => waiting?.({ value: undefined, done: true })) + return { + [Symbol.asyncIterator]() { + return { + next() { + pulled?.() + pulled = undefined + if (signal?.aborted) return Promise.resolve({ value: undefined, done: true }) + if (queued.length > 0) return Promise.resolve({ value: queued.shift(), done: false }) + return new Promise((resolve) => { waiting = resolve }) + }, + async return() { return { value: undefined, done: true } }, + } + }, + } + }, + emit(event) { + const handled = new Promise((resolve) => { pulled = resolve }) + const resolve = waiting + waiting = undefined + if (resolve) resolve({ value: event, done: false }) + else queued.push(event) + return bounded(handled, `event ${event.type} was not handled`) + }, + } +} + +test("V2 inbox prompts stay idempotent and deleted across real server restarts", async (t) => { + const dir = mkdtempSync(join(tmpdir(), "engram-opencode-v2-real-")) + const originalFetch = globalThis.fetch + const originalEngramURL = process.env.ENGRAM_URL + let server + let cleanup + // One ordered teardown: plugin, fetch bridge, server, then its directory. + t.after(async () => { + try { + await cleanup?.() + } finally { + globalThis.fetch = originalFetch + if (originalEngramURL === undefined) delete process.env.ENGRAM_URL + else process.env.ENGRAM_URL = originalEngramURL + try { await server?.stop() } + finally { rmSync(dir, { recursive: true, force: true }) } + } + }) + const executable = join(dir, process.platform === "win32" ? "real-server.exe" : "real-server") + const build = spawnSync("go", ["build", "-o", executable, "./plugin/opencode/test/support/real-server"], { + cwd: new URL("../..", import.meta.url), timeout: 120000, encoding: "utf8", + }) + assert.ifError(build.error) + assert.equal(build.status, 0, build.stderr) + + server = await startServer(executable, dir) + + // The plugin keeps one stable URL; the fetch bridge follows server restarts. + // Project detection is stubbed because the temp directory has no git remote. + const promptResponses = [] + globalThis.fetch = async (url, init) => { + const target = new URL(url) + if (target.origin !== PLUGIN_URL) return originalFetch(url, init) + if (target.pathname === "/project/current") { + return new Response(JSON.stringify({ project: "engram", project_source: "git_remote" })) + } + const response = await originalFetch(`${server.url}${target.pathname}${target.search}`, init) + if (target.pathname === "/prompts" && init?.method === "POST") { + promptResponses.push({ ...JSON.parse(init.body), status: response.status, reply: await response.clone().json() }) + } + return response + } + process.env.ENGRAM_URL = PLUGIN_URL + + const events = eventStream() + const module = await import(new URL("./engram.ts?v2-real-server", import.meta.url).href) + cleanup = await module.default.setup({ + location: { directory: dir, project: { id: PROJECT_ID } }, + event: { subscribe: (options) => events.subscribe(options) }, + session: { + async get({ sessionID }) { return { id: sessionID, projectID: PROJECT_ID } }, + hook: async () => ({ dispose: async () => {} }), + }, + tool: { hook: async () => ({ dispose: async () => {} }) }, + }) + + const text = "Rotate the staging credentials before the release" + const enqueue = (inboxID) => events.emit({ + type: "session.inbox.enqueued", + location: { directory: dir }, + data: { sessionID: "ses_root", inboxID, item: { type: "user", payload: { text }, delivery: "queue" } }, + }) + const prompts = async () => { + const response = await originalFetch(`${server.url}/prompts/recent?project=engram&limit=20`) + assert.equal(response.status, 200) + return response.json() + } + const restart = async () => { + await server.stop() + server = await startServer(executable, dir) + } + + await enqueue("msg_inbox_x") + assert.equal((await prompts()).length, 1, "first admission creates one prompt") + const promptX = promptResponses[0].reply.id + assert.equal(promptResponses[0].source_inbox_id, "msg_inbox_x") + + await enqueue("msg_inbox_x") + assert.equal((await prompts()).length, 1, "replayed admission must not create another prompt") + assert.equal(promptResponses[1].reply.id, promptX, "replay resolves to the existing prompt") + + await enqueue("msg_inbox_y") + assert.equal((await prompts()).length, 2, "distinct inbox items with identical text stay distinct") + + await restart() + await enqueue("msg_inbox_x") + await enqueue("msg_inbox_y") + assert.equal((await prompts()).length, 2, "replays after a restart must not create prompts") + + const deleted = await originalFetch(`${server.url}/prompts/${promptX}`, { method: "DELETE" }) + assert.equal(deleted.status, 200, await deleted.text()) + assert.deepEqual((await prompts()).map(({ content }) => content), [text]) + + await enqueue("msg_inbox_x") + assert.equal(promptResponses.at(-1).status, 409, "a deleted inbox identity is refused") + assert.equal((await prompts()).length, 1, "replay after deletion must not resurrect the prompt") + + await restart() + await enqueue("msg_inbox_x") + assert.equal(promptResponses.at(-1).status, 409) + assert.equal((await prompts()).length, 1, "deletion survives a restart") +}) diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs index 36fd4c8b2..2dd744ae3 100644 --- a/plugin/opencode/engram.v2.test.mjs +++ b/plugin/opencode/engram.v2.test.mjs @@ -14,8 +14,8 @@ const PROJECT_ID = "project-1" const INSTANCE_ID = "00000000000000000000000000000000" let runtimeImport = 0 -function httpResponse(data) { - return { ok: true, async json() { return data } } +function httpResponse(data, status = 200) { + return { ok: status >= 200 && status < 300, status, async json() { return data } } } // Minimal OpenCode V2 event stream. emit() resolves once the plugin asks for @@ -107,7 +107,7 @@ async function waitFor(condition, message, ms = 1000) { } } -async function setupV2(t, { sessions = new Map() } = {}) { +async function setupV2(t, { sessions = new Map(), promptResponse } = {}) { const originalFetch = globalThis.fetch const originalBun = globalThis.Bun const originalEngramURL = process.env.ENGRAM_URL @@ -130,6 +130,7 @@ async function setupV2(t, { sessions = new Map() } = {}) { if (path === "/project/current") return httpResponse({ project: "engram", project_source: "git_remote" }) if (path === "/sessions") return httpResponse({ id: body.id, status: "created" }) if (path === "/context/compaction") return httpResponse({ context: "previous session context" }) + if (path === "/prompts" && promptResponse) return promptResponse(body) return httpResponse({}) } @@ -183,6 +184,11 @@ async function setupV2(t, { sessions = new Map() } = {}) { data: { sessionID, projectID: PROJECT_ID, location: { directory: DIRECTORY }, ...(parentID ? { parentID } : {}) }, }), deleted: (sessionID) => events.emit({ type: "session.deleted", data: { sessionID } }), + enqueued: (sessionID, inboxID, item, location) => events.emit({ + type: "session.inbox.enqueued", + ...(location ? { location } : {}), + data: { sessionID, inboxID, item }, + }), posts: (path) => requests.filter((request) => request.method === "POST" && request.path === path), } } @@ -204,7 +210,6 @@ test("V2 setup registers session, tool, and event hooks and cleans them up", asy assert.deepEqual([...runtime.hooks.keys()].sort(), [ "session.compaction", "session.context", - "session.prompt", "tool.execute.after", "tool.execute.before", ]) @@ -214,7 +219,7 @@ test("V2 setup registers session, tool, and event hooks and cleans them up", asy await runtime.created("ses_root") await runtime.cleanup() - assert.equal(runtime.disposedHooks.length, 5) + assert.equal(runtime.disposedHooks.length, 4) assert.equal(runtime.events.subscriptions[0].signal.aborted, true) assert.equal(runtime.posts("/sessions/ses_root/end").length, 1, "cleanup ends registered sessions") }) @@ -259,31 +264,90 @@ test("V2 setup releases earlier registrations when a later one fails", async (t) } const disposed = [] await assert.rejects(runtime.module.default.setup(ctx), /tool hooks unavailable/) - assert.deepEqual(disposed, ["prompt", "context", "compaction"]) + assert.deepEqual(disposed, ["context", "compaction"]) }) -test("V2 prompt hook captures user prompts for the authoritative session", async (t) => { +function userItem(text, delivery = "queue") { + return { type: "user", payload: { text }, delivery } +} + +test("V2 captures user inbox items with their durable inbox identity", async (t) => { const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) - await runtime.hooks.get("session.prompt")({ - sessionID: "ses_root", - messageID: "msg_1", - prompt: { text: "Please remember the token decision" }, - delivery: "immediate", - }) + await runtime.enqueued("ses_root", "msg_inbox_1", userItem("Please remember the token decision", "steer")) assert.deepEqual(runtime.posts("/prompts").map(({ body }) => body), [ - { session_id: "ses_root", content: "Please remember the [REDACTED] decision", project: "engram" }, + { session_id: "ses_root", content: "Please remember the [REDACTED] decision", project: "engram", source_inbox_id: "msg_inbox_1" }, ]) }) -test("V2 prompt hook redacts a private block that straddles the truncation limit", async (t) => { +test("V2 ignores non-user inbox items, trivial text, and malformed events", async (t) => { const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) - await runtime.hooks.get("session.prompt")({ - sessionID: "ses_root", - messageID: "msg_1", - prompt: { text: `${"a".repeat(1980)}PIN=42 trailing` }, - delivery: "immediate", + await runtime.enqueued("ses_root", "msg_synthetic", { type: "synthetic", payload: { text: "Synthetic reminder text for the agent" }, delivery: "queue" }) + await runtime.enqueued("ses_root", "msg_compaction", { type: "compaction", payload: {}, delivery: "queue" }) + await runtime.enqueued("ses_root", "msg_move", { + type: "move", + payload: { location: { directory: "/work/other" }, projectID: PROJECT_ID }, + delivery: "queue", + }) + await runtime.enqueued("ses_root", "msg_short", userItem("too short")) + await runtime.enqueued("ses_root", "", userItem("A prompt without a durable inbox identity")) + await runtime.events.emit({ type: "session.inbox.enqueued", data: { sessionID: "ses_root", inboxID: "msg_x" } }) + + assert.equal(runtime.posts("/prompts").length, 0) +}) + +test("V2 ignores inbox items from subagent sessions and other locations", async (t) => { + const runtime = await setupV2(t, { + sessions: new Map([["ses_root", sessionInfo("ses_root")], ["ses_child", sessionInfo("ses_child", "ses_root")]]), + }) + await runtime.enqueued("ses_child", "msg_child", userItem("Delegated prompt text written by the parent agent")) + await runtime.enqueued("ses_root", "msg_elsewhere", userItem("Prompt admitted in another location"), { directory: "/work/other" }) + await runtime.enqueued("ses_unknown", "msg_unknown", userItem("Prompt for a session this instance cannot resolve")) + + assert.equal(runtime.posts("/prompts").length, 0) +}) + +test("V2 keeps distinct inbox items with identical text distinct", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + await runtime.enqueued("ses_root", "msg_inbox_1", userItem("Run the full test suite again please")) + await runtime.enqueued("ses_root", "msg_inbox_2", userItem("Run the full test suite again please")) + + assert.deepEqual(runtime.posts("/prompts").map(({ body }) => body.source_inbox_id), ["msg_inbox_1", "msg_inbox_2"]) +}) + +test("V2 forwards replayed inbox items with the same identity for the server to deduplicate", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + await runtime.enqueued("ses_root", "msg_inbox_1", userItem("Run the full test suite again please")) + await runtime.enqueued("ses_root", "msg_inbox_1", userItem("Run the full test suite again please")) + + const prompts = runtime.posts("/prompts").map(({ body }) => body) + assert.equal(prompts.length, 2, "no in-memory dedup: the server owns replay identity") + assert.deepEqual(prompts[0], prompts[1]) +}) + +test("V2 treats a deleted inbox identity (409) as a silent no-op without retry", async (t) => { + const runtime = await setupV2(t, { + sessions: new Map([["ses_root", sessionInfo("ses_root")]]), + promptResponse: () => httpResponse({ error: "prompt inbox identity was deleted" }, 409), }) + const errors = [] + const originalError = console.error + const originalWarn = console.warn + console.error = (...args) => errors.push(args) + console.warn = (...args) => errors.push(args) + t.after(() => { console.error = originalError; console.warn = originalWarn }) + + await withTimeout(runtime.enqueued("ses_root", "msg_deleted", userItem("A prompt the user already deleted")), "409 stalled the event loop") + await withTimeout(runtime.created("ses_other"), "event after a 409 was not handled") + + assert.equal(runtime.posts("/prompts").length, 1) + assert.deepEqual(errors, []) + assert.equal(runtime.events.subscriptions.length, 1) +}) + +test("V2 inbox capture redacts a private block that straddles the truncation limit", async (t) => { + const runtime = await setupV2(t, { sessions: new Map([["ses_root", sessionInfo("ses_root")]]) }) + await runtime.enqueued("ses_root", "msg_inbox_1", userItem(`${"a".repeat(1980)}PIN=42 trailing`)) const prompts = runtime.posts("/prompts").map(({ body }) => body) assert.equal(prompts.length, 1) From 4b89b896d45200a6ea9d0885e1af05208fbc9c60 Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 10:46:00 +0200 Subject: [PATCH 8/9] fix(opencode): keep V2 reconnect backoff across server handshakes --- internal/setup/plugins/opencode/engram.ts | 4 +- plugin/opencode/engram.ts | 4 +- plugin/opencode/engram.v2.test.mjs | 50 ++++++++++++++++++++--- 3 files changed, 51 insertions(+), 7 deletions(-) diff --git a/internal/setup/plugins/opencode/engram.ts b/internal/setup/plugins/opencode/engram.ts index 4e9395e0e..93a0dfceb 100644 --- a/internal/setup/plugins/opencode/engram.ts +++ b/internal/setup/plugins/opencode/engram.ts @@ -1032,7 +1032,9 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { try { for await (const event of ctx.event.subscribe({ signal: abort.signal })) { if (abort.signal.aborted) break - retryMs = V2_EVENT_RETRY_MIN_MS + // Every subscription opens with a server.connected handshake; only a + // real event proves the stream is healthy enough to reset the backoff. + if (event?.type !== "server.connected") retryMs = V2_EVENT_RETRY_MIN_MS try { const prompt = v2InboxPrompt(event, ctx.location.directory) if (prompt) await (hooks[CAPTURE_PROMPT] as CapturePrompt)(prompt.sessionID, prompt.text, prompt.inboxID) diff --git a/plugin/opencode/engram.ts b/plugin/opencode/engram.ts index 4e9395e0e..93a0dfceb 100644 --- a/plugin/opencode/engram.ts +++ b/plugin/opencode/engram.ts @@ -1032,7 +1032,9 @@ async function setupEngramV2(ctx: V2Context): Promise<() => Promise> { try { for await (const event of ctx.event.subscribe({ signal: abort.signal })) { if (abort.signal.aborted) break - retryMs = V2_EVENT_RETRY_MIN_MS + // Every subscription opens with a server.connected handshake; only a + // real event proves the stream is healthy enough to reset the backoff. + if (event?.type !== "server.connected") retryMs = V2_EVENT_RETRY_MIN_MS try { const prompt = v2InboxPrompt(event, ctx.location.directory) if (prompt) await (hooks[CAPTURE_PROMPT] as CapturePrompt)(prompt.sessionID, prompt.text, prompt.inboxID) diff --git a/plugin/opencode/engram.v2.test.mjs b/plugin/opencode/engram.v2.test.mjs index 2dd744ae3..d1d853423 100644 --- a/plugin/opencode/engram.v2.test.mjs +++ b/plugin/opencode/engram.v2.test.mjs @@ -28,6 +28,7 @@ function eventStream() { let pulled let closed = false let interrupt + let handshake = false const subscriptions = [] const settle = (result) => { const resolve = waiting @@ -45,10 +46,16 @@ function eventStream() { closed = true if (waiting) settle({ value: undefined, done: true }) }) + // OpenCode 2.x opens every subscription with a server.connected handshake. + let greet = handshake return { [Symbol.asyncIterator]() { return { next() { + if (greet) { + greet = false + return Promise.resolve({ value: { type: "server.connected", data: {} }, done: false }) + } pulled?.() pulled = undefined if (interrupt) { @@ -86,6 +93,9 @@ function eventStream() { if (waiting) settle(error) else interrupt = error }, + withHandshake() { + handshake = true + }, closeForever() { closed = true if (waiting) settle({ value: undefined, done: true }) @@ -93,6 +103,20 @@ function eventStream() { } } +// Records the reconnect delays the plugin schedules, so backoff is asserted as +// a schedule rather than by counting subscriptions within a real-time window. +const RECONNECT_DELAYS = new Set([50, 100, 200, 400, 800, 1600, 3200, 5000]) +function recordReconnectDelays(t) { + const original = globalThis.setTimeout + const delays = [] + globalThis.setTimeout = (callback, ms, ...args) => { + if (RECONNECT_DELAYS.has(ms)) delays.push(ms) + return original(callback, ms, ...args) + } + t.after(() => { globalThis.setTimeout = original }) + return delays +} + function withTimeout(promise, message, ms = 1000) { let timer const timeout = new Promise((_, reject) => { timer = setTimeout(() => reject(new Error(message)), ms) }) @@ -480,13 +504,29 @@ test("V2 keeps listening when handling one event throws", async (t) => { }) test("V2 backs off between re-subscriptions and stops after cleanup", async (t) => { + const delays = recordReconnectDelays(t) const runtime = await setupV2(t) runtime.events.closeForever() - await new Promise((resolve) => setTimeout(resolve, 300)) + try { + await waitFor(() => delays.length >= 3, "plugin did not keep re-subscribing", 5000) + assert.deepEqual(delays.slice(0, 3), [50, 100, 200]) + } finally { + await withTimeout(runtime.cleanup(), "cleanup waited on the reconnect backoff") + } const attempts = runtime.events.subscriptions.length - assert.ok(attempts >= 2 && attempts <= 6, `expected bounded re-subscription attempts, got ${attempts}`) - - await withTimeout(runtime.cleanup(), "cleanup waited on the reconnect backoff") - await new Promise((resolve) => setTimeout(resolve, 150)) + await new Promise((resolve) => setTimeout(resolve, 20)) assert.equal(runtime.events.subscriptions.length, attempts, "no re-subscription after cleanup") }) + +test("V2 keeps backing off when each subscription only delivers the server handshake", async (t) => { + const delays = recordReconnectDelays(t) + const runtime = await setupV2(t) + runtime.events.withHandshake() + runtime.events.closeForever() + try { + await waitFor(() => delays.length >= 3, "plugin did not keep re-subscribing", 5000) + assert.deepEqual(delays.slice(0, 3), [50, 100, 200], "server.connected must not reset the backoff") + } finally { + await withTimeout(runtime.cleanup(), "cleanup waited on the reconnect backoff") + } +}) From 8b5fb3f2b38e727d3e56a7e603bc4dd5d9b520ac Mon Sep 17 00:00:00 2001 From: Alan Buscaglia Date: Tue, 29 Sep 2026 19:44:49 +0200 Subject: [PATCH 9/9] ci(opencode): run the V2 real-server prompt identity regression --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 948eb820d..d197b7999 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -196,7 +196,7 @@ jobs: # Covers the dual-major (1.x server / 2.x setup) OpenCode adapter. - name: Run OpenCode plugin tests - run: node --test plugin/opencode/engram.test.mjs plugin/opencode/engram.v2.test.mjs + run: node --test plugin/opencode/engram.test.mjs plugin/opencode/engram.v2.test.mjs plugin/opencode/engram.v2.real-server.test.mjs perf-ratchet: name: Performance Ratchet