Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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 plugin/opencode/engram.v2.real-server.test.mjs

perf-ratchet:
name: Performance Ratchet
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
Expand Down
2 changes: 2 additions & 0 deletions docs/PLUGINS.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,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, 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.

### What the Plugin Does
Expand Down
243 changes: 222 additions & 21 deletions internal/setup/plugins/opencode/engram.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>

export const Engram: Plugin = async (ctx) => {
let project = "unknown"
let projectResolutionError = ""
Expand Down Expand Up @@ -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 <private> block straddling the
// limit would otherwise lose its closing tag and leak.
content: truncate(stripPrivateTags(content), 2000),
project,
...(sourceInboxID ? { source_inbox_id: sourceInboxID } : {}),
Comment thread
Copilot marked this conversation as resolved.
},
})
}

// Try to start engram server if not running
try {
const expectedID = CONFIGURED_ENGRAM_URL ? "" : localInstanceID()
Expand Down Expand Up @@ -579,6 +612,8 @@ export const Engram: Plugin = async (ctx) => {
}

return {
[CAPTURE_PROMPT]: capturePrompt,

dispose: async () => {
disposed = true
if (!localReady) return
Expand Down Expand Up @@ -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")
Expand All @@ -658,22 +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,
content: stripPrivateTags(truncate(finalContent, 2000)),
project,
},
})
}
await capturePrompt(input.sessionID, content || fallback)
},

// ─── Tool Execution Hook ─────────────────────────────────────
Expand Down Expand Up @@ -860,4 +876,189 @@ 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<void> }
type V2Hook = (name: string, callback: (input: any) => Promise<void> | void) => Promise<V2Registration>
type V2Context = {
location: { directory: string; project?: { id?: string } }
event: { subscribe: (options?: { signal?: AbortSignal }) => AsyncIterable<any> }
session: { get: (input: { sessionID: string }) => Promise<any>; 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<void>): Promise<void> {
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 }
})
}

// V2 `Tool.Result.content` is `string | Content[]`; `subagent` returns a string.
function v2ToolResultText(result: any): string {
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")
: ""
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`.
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
}

// 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
const V2_EVENT_RETRY_MAX_MS = 5000

function delayUnlessAborted(ms: number, signal: AbortSignal): Promise<void> {
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<void>> {
const hooks: Record<string | symbol, any> = 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<void> = Promise.resolve()
const cleanup = async () => {
abort.abort()
await Promise.all(registrations.map((registration) => registration.dispose()))
await listening
await hooks.dispose?.()
}

try {
// 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 }))
}))

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)
}))

listening = (async () => {
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
// 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)
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.
}
if (abort.signal.aborted) break
await delayUnlessAborted(retryMs, abort.signal)
retryMs = Math.min(retryMs * 2, V2_EVENT_RETRY_MAX_MS)
}
})()
} 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 }
11 changes: 11 additions & 0 deletions plugin/opencode/engram.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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)}<private>PIN=42</private> 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(
Expand Down
Loading
Loading