From 4fcca7c8c8d1434baef9f23c3009d9dfe743dae6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Carlos=20Magalh=C3=A3es?= <88864312+JohnC0de@users.noreply.github.com> Date: Wed, 23 Sep 2026 02:57:47 -0300 Subject: [PATCH] feat(opencode2): bind the plugin to the 2.0.10+ event api with turn checkpoints OpenCode 2.0.10 replaced the lifecycle events the generated plugin listened for (session.idle and friends) with session.execution.{succeeded,failed, interrupted}, and it runs one long-lived service behind every CLI: closing a terminal is not a session end, so sessions captured through the old binding rarely produced a summary page or a baton. The plugin also re-ran two `git` processes per event, lost startup context after the first model request, dropped content-only tool results, and shared queue and spool state across the per-location instances OpenCode 2 loads. Plugin (install_hooks.rs, render_shared.rs, the node host fixture): - binds session.execution.*, session.text.ended, session.moved, session.deleted, session.compaction.started and the context/prompt/tool hooks; every name was checked against the OpenCode 2.0.14 binary; - claims startup context once per root session and re-injects it on every model request; child sessions never claim it or publish a baton; - keeps queue, spool and cleanup state per location instance; spool names can no longer collide within a millisecond, and a torn-down host cancels what is still in flight only after its final session-ends had the drain budget (all generated TypeScript integrations share this runtime now); - opt-in assistant capture (--capture-assistant) hands the last completed text to the native hook, which sanitizes and caps it before the spool or the wire, and falls back to the plain stop hook if that binary cannot run. Server (router.rs, ops.rs, reader.rs, writer.rs): - a completed root turn is a turn checkpoint: sessions/.md and the automatic baton are refreshed deterministically (no LLM) in the session row's own scope, keeping one open baton per live session, refreshed in place and audited; a checkpoint that lost the race with the session's end touches nothing; - an explicit, keyed session.moved rebinds the live session row to its new directory on first delivery (compare-and-set on the cwd it left), so the session's end and checkpoints follow it; - a SessionEnd whose resolved scope drifted under the same cwd (a .ai-memory.toml appeared mid-session) ends the session instead of stranding it open; scope-drifted ordinary events are recorded as upstream records any other drifted event; - the latest captured assistant excerpt rides in the automatic baton; it is not rendered into the git-tracked session page. Docs: install.md, support-matrix.md, auto-scope.md, SECURITY.md and DATA_HANDLING.md describe the checkpoint, routing and capture behavior. Verified: cargo fmt --all -- --check; cargo clippy --workspace --all-targets -- -D warnings; cargo test --workspace --all-targets (3569 passed, on release/2.5); the node host fixture runs under cargo test (node >= 22.6) and fails if unload aborts deliveries before its session-ends. Live on Windows 11 with OpenCode 2.0.14: after a 64 -> 66 store migration, the running OpenCode service loaded the regenerated plugin and one real turn through it (a Code Mode call to memory_status) recorded its prompt, tool events and stop, logged "turn checkpoint written; native session remains open", left the session open, wrote sessions/.md and exactly one open baton carrying the captured assistant excerpt, which the session page does not contain. --- CHANGELOG.md | 31 + DATA_HANDLING.md | 3 +- SECURITY.md | 3 + crates/ai-memory-cli/src/cli.rs | 11 +- .../src/commands/install_hooks.rs | 613 +++++++++---- .../src/commands/openclaw_plugin.rs | 10 +- .../src/commands/render_shared.rs | 99 ++- .../tests/fixtures/opencode2-plugin.mjs | 164 ++++ .../ai-memory-hooks/src/assistant_capture.rs | 34 +- crates/ai-memory-hooks/src/router.rs | 826 +++++++++++++++-- crates/ai-memory-store/src/ops.rs | 837 ++++++++++++++++-- crates/ai-memory-store/src/reader.rs | 26 + crates/ai-memory-store/src/writer.rs | 30 + docs/auto-scope.md | 23 +- docs/install.md | 24 +- docs/support-matrix.md | 2 +- 16 files changed, 2401 insertions(+), 335 deletions(-) create mode 100644 crates/ai-memory-cli/tests/fixtures/opencode2-plugin.mjs diff --git a/CHANGELOG.md b/CHANGELOG.md index e6685540..4466b631 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -86,6 +86,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 session in a directory. Concurrent OpenCode sessions in different projects no longer read each other's active project; explicit `workspace`/`project` arguments still win and `X-Memory-Actor-Session-Id` keeps precedence. (#864) +- OpenCode 2 turn checkpoints: every completed root turn + (`session.execution.succeeded`, `.failed` or `.interrupted`, OpenCode + 2.0.10+) now refreshes `sessions/.md` and the session's automatic + handoff without ending the native session, because OpenCode 2's shared + service outlives the CLI and closing a terminal is not a session end. The + checkpoint is deterministic (no LLM call), is written in the session's own + scope, keeps one open baton per live session (refreshed in place, audited as + `refresh_handoff`), and never touches a session that already ended. Child + sessions neither claim startup context nor publish batons. (#865) +- `install-hooks --agent opencode2 --capture-assistant` extends the + assistant/Stop capture double opt-in to OpenCode 2: the plugin forwards the + last completed assistant text through the native hook's sanitizer before it + reaches the spool or the wire. The captured excerpt continues the work in + the next session's automatic handoff; it is not rendered into the + git-tracked session page. (#865) ### Changed - Grok Build CLI shows a pending handoff, and an opted-in `[briefing]`, as @@ -103,6 +118,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 take another (the multi-session claim-once invariant is preserved). Delivery of this PostToolUse handoff is exempt from `AI_MEMORY_CAPTURE_OWNER` capture suppression, like the other context-delivery events. (#840) +- The generated OpenCode 2 plugin binds to the OpenCode 2.0.10+ event and hook + API (`session.execution.*`, `session.text.ended`, `session.moved`, the + `context` hook). Startup context is claimed once per root session and + retained on every later model request; content-only tool results are + captured; each location instance owns its own queue, spool state and + cleanup; an explicit `session.moved` rebinds the live session to its new + directory so its later end lands there. Checked against OpenCode 2.0.14. (#865) ### Fixed - `memory_query`'s vector stream called the generic `Embedder::embed` @@ -187,6 +209,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 event: their `git` lookups set `windowsHide`. The repo-root project lookup behind those spawns is memoized per cwd instead of running two synchronous `git` processes on every event. (#863) +- A `SessionEnd` whose resolved scope drifted from its session (a + `.ai-memory.toml` appeared under the running session) was refused as a + foreign-scope end and stranded the session open forever. It now ends the + session when owner and agent match and the event comes from the session's + own (normalized) cwd; a different cwd, operator or agent is still refused. (#865) +- Generated TypeScript integrations no longer overwrite each other's spooled + events written in the same millisecond, and a host that tears its capture + state down cancels still-pending deliveries after the drain budget instead + of waiting out each request's timeout. (#865) ## [2.4.0] - 2026-09-21 diff --git a/DATA_HANDLING.md b/DATA_HANDLING.md index 69dfcee8..96013700 100644 --- a/DATA_HANDLING.md +++ b/DATA_HANDLING.md @@ -52,7 +52,8 @@ and each requires a deliberate config change to turn on. opt-in and sanitized"). Persisting the coding assistant's final-turn text requires a double opt-in: `capture_assistant` on the server *and* `install-hooks --capture-assistant` on the client. Once enabled, captured - text flows into consolidation/reviewer prompts and — only if you have + text rides in the session's automatic handoff to the next session and + flows into consolidation/reviewer prompts and — only if you have separately configured a cloud LLM provider — is sent to that provider. The flag is global to the install; there is no per-project exclusion once it's on. diff --git a/SECURITY.md b/SECURITY.md index 40c95764..023097b4 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -101,6 +101,9 @@ what the project is and is not designed to defend against. `[REDACTED]`. - Captured assistant text flows into the consolidation and reviewer prompts, and — if you configure a cloud LLM provider — is sent to that provider. + The latest excerpt of a session also rides in its automatic handoff, so + the next session that claims the baton receives it as startup context. It + is not rendered into the git-tracked session page. - The opt-in is **global** to the install: there is no per-project marker to exclude a sensitive repository once the flag is on (assistant text is not path-attributable). Turn the server flag off to disable it everywhere. diff --git a/crates/ai-memory-cli/src/cli.rs b/crates/ai-memory-cli/src/cli.rs index 94eb1d9f..94c340e7 100644 --- a/crates/ai-memory-cli/src/cli.rs +++ b/crates/ai-memory-cli/src/cli.rs @@ -2313,7 +2313,7 @@ pub struct HookArgs { /// the local spool or the wire. #[arg(long, value_enum)] pub capture_mode: Option, - /// Opt in to assistant/Stop capture: on a Claude Code `stop` event, attach a + /// Opt in to assistant/Stop capture: on a supported agent's `stop` event, attach a /// sanitized, capped excerpt of the assistant's final turn as the Stop body. /// Baked onto the native `stop` command by /// `install-hooks --capture-assistant`; the server must also enable @@ -2391,11 +2391,10 @@ pub struct InstallHooksArgs { /// silently revert `repo-root` back to `basename`. #[arg(long, value_enum)] pub project_strategy: Option, - /// Bake `--capture-assistant` onto the installed native `stop` command so a - /// Claude Code `stop` event carries a sanitized excerpt of the assistant's - /// final turn (#196). Only valid for `--agent claude-code` on a native - /// platform; the server must also set `capture_assistant = true`. Re-running - /// without this flag removes it (idempotent). Default off. + /// Capture a sanitized excerpt of the assistant's final turn through the + /// native Stop hook. Supported for Claude Code, Codex and OpenCode 2 on a + /// native platform; the server must also set `capture_assistant = true`. + /// A bare re-apply preserves an existing opt-in. Default off. #[arg(long)] pub capture_assistant: bool, /// Persist the capture failure mode for this install (#446). Under diff --git a/crates/ai-memory-cli/src/commands/install_hooks.rs b/crates/ai-memory-cli/src/commands/install_hooks.rs index 513a3778..cec4fbe6 100644 --- a/crates/ai-memory-cli/src/commands/install_hooks.rs +++ b/crates/ai-memory-cli/src/commands/install_hooks.rs @@ -24,7 +24,7 @@ use crate::cli::{AgentChoice, CaptureModeArg, InstallHooksArgs, McpClient, Proje use crate::commands::apply_shared::{ApplyOutcome, apply_atomic, mutate_json, mutate_toml}; use crate::commands::install_mcp; use crate::commands::openclaw_plugin; -use crate::commands::path_util::home_dir; +use crate::commands::path_util::{home_dir, strip_windows_verbatim_prefix}; use crate::commands::render_shared::{ ANTIGRAVITY_LIFECYCLE_EVENTS, ANTIGRAVITY_TOOL_EVENTS, CODEX_PROFILE, COMMAND_CODE_PROFILE, CURSOR_PROFILE, GEMINI_PROFILE, KIMI_CODE_EVENTS, KIRO_CLI_V2_EVENTS, KIRO_CLI_V3_EVENTS, @@ -33,7 +33,7 @@ use crate::commands::render_shared::{ build_kiro_cli_v2_hooks_value, build_kiro_cli_v3_hooks_value, build_pool_settings_yaml_with_data_dir, build_profile_payload_for_agent, hook_script_for_claude_code, hook_script_for_current_platform, kimi_code_hook_commands, - local_hook_policy_v1_supported, ts_capture_policy_v1, ts_string_literal, + local_hook_policy_v1_supported, ts_capture_policy_v1, ts_string_literal, ts_timeout_signal, }; use crate::commands::uninstall::hook_command_is_ours; use crate::config::{Config, DEFAULT_SERVER_URL}; @@ -475,16 +475,16 @@ pub fn run(config: &Config, mut args: InstallHooksArgs) -> Result<()> { "[ai-memory] selected shell/PowerShell compatibility path does not enforce capture-policy v1; use a native platform selection or generated integration." ); } - // Assistant/Stop capture is Claude Code + native-platform only (#196). No + // Assistant capture requires a native client-side sanitizer. No // silent fallback: bail so an operator on a script-fallback platform or a // different agent is told the flag has no effect instead of installing a // command whose capture would be silently dropped. if args.capture_assistant && !capture_assistant_allowed(args.agent) { anyhow::bail!( - "--capture-assistant requires --agent claude-code on a native hook platform \ + "--capture-assistant requires --agent claude-code, codex, or opencode2 on a native hook platform \ (PosixNative/WindowsNative). The current selection uses the script fallback or a \ different agent, where the opt-in cannot take effect. Remove --capture-assistant or \ - switch to a native Claude Code install." + switch to a supported native install." ); } if (args.no_capture_prompts || args.capture_prompts) @@ -650,9 +650,13 @@ pub fn run(config: &Config, mut args: InstallHooksArgs) -> Result<()> { AgentChoice::OpenCode => { render_opencode_plugin(&server_url, auth, strategy, &preview_capture_mode) } - AgentChoice::OpenCode2 => { - render_opencode2_plugin(&server_url, auth, strategy, &preview_capture_mode) - } + AgentChoice::OpenCode2 => render_opencode2_plugin( + &server_url, + auth, + strategy, + &preview_capture_mode, + args.capture_assistant, + ), AgentChoice::Pi => render_pi_extension(&server_url, auth, strategy, &preview_capture_mode), AgentChoice::Omp => render_omp_extension( &server_url, @@ -888,8 +892,7 @@ fn baked_claude_prompt_capture(existing: &str) -> Option { /// Whether the currently-installed config for this agent already bakes the /// `--capture-assistant` opt-in. Used to preserve that opt-in across a bare /// re-apply (e.g. `ai-memory run` auto-wire) instead of dropping it. Only the -/// agents where assistant capture is allowed (Claude Code, Codex) reach here, -/// and both write the nested-hooks JSON shape. +/// agents where assistant capture is supported reach here. fn existing_capture_assistant_opt_in(args: &InstallHooksArgs) -> bool { existing_agent_config(args) .as_deref() @@ -901,6 +904,11 @@ fn existing_capture_assistant_opt_in(args: &InstallHooksArgs) -> bool { /// whether its Stop command carries `--capture-assistant`; `None` when it is not /// an ai-memory install at all (so there is nothing to preserve). fn baked_capture_assistant(existing: &str) -> Option { + if existing + .contains("// Auto-generated by `ai-memory install-hooks --agent opencode2 --apply`.") + { + return Some(existing.contains("const CAPTURE_ASSISTANT = true;")); + } let document: serde_json::Value = serde_json::from_str(existing).ok()?; let hooks = document.get("hooks")?.as_object()?; let mut saw_ai_memory = false; @@ -1616,14 +1624,16 @@ fn overlay_kiro_cli_event_hooks( /// ai-memory cares about (`CLAUDE_CODE_EVENTS`); preserve every other hook the /// user has wired up to other tools. /// Whether `--capture-assistant` may take effect for this agent + platform -/// (#196, #743): Claude Code and Codex on a native hook platform. Both carry -/// `last_assistant_message` on their `Stop` payload (see -/// `ai_memory_hooks::assistant_capture`). Any other agent or a script-fallback +/// (#196, #743): Claude Code, Codex and OpenCode 2 on a native hook platform. +/// OpenCode 2 forwards completed text through the same native Stop sanitizer. +/// Any other agent or a script-fallback /// platform cannot honor the opt-in, so the installer bails instead of enabling /// it silently. fn capture_assistant_allowed(agent: AgentChoice) -> bool { - matches!(agent, AgentChoice::ClaudeCode | AgentChoice::Codex) - && local_hook_policy_v1_supported() + matches!( + agent, + AgentChoice::ClaudeCode | AgentChoice::Codex | AgentChoice::OpenCode2 + ) && local_hook_policy_v1_supported() } fn prompt_capture_options_allowed(agent: AgentChoice) -> bool { @@ -3067,7 +3077,13 @@ fn apply_to_opencode2_plugin( None => opencode2_plugin_path()?, }; let strategy = args.project_strategy.and_then(ProjectStrategyArg::baked); - let body = build_opencode2_plugin(server_url, auth_token, strategy, capture_mode)?; + let body = build_opencode2_plugin( + server_url, + auth_token, + strategy, + capture_mode, + args.capture_assistant, + )?; let outcome = apply_atomic(&path, move |_existing| Ok(body.clone()))?; println!( @@ -3094,6 +3110,7 @@ fn render_opencode2_plugin( auth_token: Option<&str>, project_strategy: Option<&str>, capture_mode: &str, + capture_assistant: bool, ) -> Result<()> { println!( "// OpenCode 2.0 beta plugin — write to ~/.config/opencode/plugins/ai-memory-opencode2.ts" @@ -3103,7 +3120,13 @@ fn render_opencode2_plugin( println!(); println!( "{}", - build_opencode2_plugin(server_url, auth_token, project_strategy, capture_mode)? + build_opencode2_plugin( + server_url, + auth_token, + project_strategy, + capture_mode, + capture_assistant + )? ); Ok(()) } @@ -3119,6 +3142,7 @@ fn build_opencode2_plugin( auth_token: Option<&str>, project_strategy: Option<&str>, capture_mode: &str, + capture_assistant: bool, ) -> Result { let v1 = build_opencode_plugin(server_url, auth_token, project_strategy, capture_mode); const BANNER_V1: &str = @@ -3139,155 +3163,281 @@ fn build_opencode2_plugin( let Some(binding_at) = out.find(V1_BINDING_ANCHOR) else { anyhow::bail!("opencode2 template drifted: v1 host binding anchor not found"); }; - let mut rebuilt = out[..binding_at].to_string(); + let prelude = &out[..binding_at]; + let (imports, state) = prelude + .split_once("const SERVER =") + .context("missing OpenCode capture prelude")?; + let rebuilt = replace_opencode_anchor( + imports, + "import type { Plugin } from \"@opencode-ai/plugin\";", + "import type { Plugin } from \"@opencode/plugin\";", + "plugin import", + )?; + let mut rebuilt = replace_opencode_anchor( + &rebuilt, + "import { execFileSync } from \"node:child_process\";", + "import { execFileSync, spawn } from \"node:child_process\";", + "child_process import", + )?; + rebuilt.push_str("import { randomBytes } from \"node:crypto\";\n"); + rebuilt.push_str("const AiMemoryOpencode2: Plugin.Plugin = {\n id: \"ai-memory-opencode2\",\n setup: async (ctx) => {\n"); + // Only the assistant-capture opt-in bakes this host's executable path; + // without it the generated file stays portable between machines. Bake the + // path as invoked, like every other native hook config: a canonicalized + // one pins the target behind a versioned symlink or junction. + let native_hook = if capture_assistant { + let exe = + std::env::current_exe().context("cannot locate ai-memory native hook executable")?; + ts_string_literal(&strip_windows_verbatim_prefix(&exe.to_string_lossy())) + } else { + "undefined".to_string() + }; + rebuilt.push_str(&format!( + "const CAPTURE_ASSISTANT = {capture_assistant};\nconst NATIVE_HOOK: string | undefined = {native_hook};\n" + )); + rebuilt.push_str(&format!( + "const NATIVE_PROJECT_STRATEGY: string | undefined = {};\n", + project_strategy + .map(ts_string_literal) + .unwrap_or_else(|| "undefined".to_string()) + )); + // A module is shared by location graphs; queues and session ownership are not. + rebuilt.push_str("const SERVER ="); + // Cache an empty successful response, but retry an unavailable server on a + // later context hook. The v1 helper historically conflated both outcomes. + let state = replace_opencode_anchor( + state, + "return text.length > 0 ? text : undefined;", + "return text;", + "handoff result", + )?; + // Mint at event creation, not only on failure: move events must remain + // idempotent when an acknowledged delivery is later replayed from disk. + let state = replace_opencode_anchor( + &state, + "url.searchParams.set(\"event\", event);\n url.searchParams.set(\"agent\", AGENT);", + "url.searchParams.set(\"event\", event);\n url.searchParams.set(\"agent\", AGENT);\n url.searchParams.set(\"ingest_key\", randomBytes(16).toString(\"hex\"));", + "hook ingest key", + )?; + rebuilt.push_str(&state); rebuilt.push_str(OPENCODE2_BINDING); Ok(rebuilt) } fn replace_opencode_anchor(haystack: &str, from: &str, to: &str, what: &str) -> Result { - if !haystack.contains(from) { - anyhow::bail!("opencode2 template drifted: v1 {what} anchor not found"); + if haystack.matches(from).count() != 1 { + anyhow::bail!("opencode2 template drifted: v1 {what} anchor not found exactly once"); } Ok(haystack.replacen(from, to, 1)) } -/// V2 host binding for the shared capture prelude: the beta's `{ id, setup }` -/// plugin shape with its session/tool/event hooks (verified against -/// `@opencode-ai/plugin@beta`, including a `tsc --noEmit` pass over the -/// rendered file). Lifecycle arrives on `ctx.event.subscribe` whose -/// envelopes carry `{ type, data, location }` (v1's `event.properties` -/// is kept as a fallback because the beta schema is still changing). +/// V2 host binding for the shared capture prelude, written against +/// `@opencode/plugin@2.0.10`; every event it handles is still emitted by +/// OpenCode 2.0.14. Lifecycle envelopes carry `{ type, data, location }`. /// Handoff injection moved from v1's removed /// `experimental.chat.system.transform` to the `context` hook, which edits /// the outgoing model call without persisting into history — guarded to -/// inject once per session because it fires on every continuation. +/// claim once per session and retain the result on every model continuation. const OPENCODE2_BINDING: &str = r#" -// `Plugin.define` is an identity wrapper, so the binding exports the -// `{ id, setup }` shape directly and keeps the shared `import type` line: -// a runtime import of `@opencode-ai/plugin` does not resolve from the -// global plugins dir and fails the load. -const AiMemoryOpencode2: Plugin = { - id: "ai-memory-opencode2", - setup: async (ctx) => { - const ctxAny = ctx as any; - const directory = ctxAny?.location?.directory; - const controller = new AbortController(); - void (async () => { + const directory = ctx.location.directory; + type SessionID = Parameters[0]["sessionID"]; + type SessionInfo = Awaited>; + const sessions = new Map>(); + const assistantText = new Map(); + // Stops listening on unload. Deliveries keep `hookAbort` until the final + // session-ends have had their drain budget. + const unload = new AbortController(); + + async function sessionInfo(id: SessionID): Promise { + let pending = sessions.get(id); + if (!pending) { + pending = ctx.session.get({ sessionID: id }, { signal: unload.signal }); + sessions.set(id, pending); + pending.catch(() => sessions.delete(id)); + } + return pending; + } + + // `agent_id` marks a subagent session (#755); a root session has none. + function subagentMarker(session: SessionInfo | undefined): { agent_id?: SessionID } { + return session?.parentID ? { agent_id: session.parentID } : {}; + } + + function hookPayload(id: SessionID, session: SessionInfo) { + return { sessionID: id, cwd: session.location.directory, ...subagentMarker(session) }; + } + + function forgetSession(id: SessionID): void { + sessions.delete(id); + assistantText.delete(id); + handoffFetches.delete(id); + } + + async function ensureSession(id: SessionID): Promise { + const session = await sessionInfo(id); + // The public event stream can reach multiple loaded location plugins. + // Hooks and cleanup belong exclusively to the current location owner. + if (session.location.directory !== directory) { + sessions.delete(id); + return undefined; + } + startSession(id, session.location.directory, { title: session.title, ...subagentMarker(session) }); + return session; + } + + // Assistant text uses the native privacy boundary before it can reach disk + // or the wire. Only the explicit client/server double opt-in enables it. + async function postStop(id: SessionID): Promise { + const session = await ensureSession(id); + if (!session) return; + const text = assistantText.get(id); + assistantText.delete(id); + const payload = { ...hookPayload(id, session), turn_checkpoint: !session.parentID }; + if (!NATIVE_HOOK || !text) { + postHook("stop", payload); + return; + } + await requestHookDrain(); + if (unload.signal.aborted) return; + await new Promise((resolve) => { + const args = ["--data-dir", dirname(hookSpoolDir()), "hook", "--agent", AGENT, + "--event", "stop", "--server-url", SERVER, "--capture-assistant", + "--capture-mode", CAPTURE_MODE]; + if (NATIVE_PROJECT_STRATEGY) args.push("--project-strategy", NATIVE_PROJECT_STRATEGY); + const child = spawn(NATIVE_HOOK, args, { + stdio: ["pipe", "ignore", "pipe"], + windowsHide: true, + env: { ...process.env, AI_MEMORY_AUTH_TOKEN: resolveToken() ?? "" }, + timeout: 5000, + signal: unload.signal, + }); + child.stderr.resume(); + child.stdin.on("error", () => {}); + child.on("error", (error) => { + if (unload.signal.aborted) return resolve(); + // The hook never ran (moved or missing binary): keep the turn's + // Stop and checkpoint, without the assistant text. + console.warn("ai-memory assistant capture failed", error.message); + postHook("stop", payload); + resolve(); + }); + child.on("close", (code) => { + if (code !== 0 && !unload.signal.aborted) console.warn("ai-memory assistant capture exited", code); + resolve(); + }); + child.stdin.end(JSON.stringify({ ...payload, last_assistant_message: text })); + }); + } + + const eventTask = (async () => { try { - for await (const evt of ctx.event.subscribe({ signal: controller.signal })) { - const event = evt as any; - const type = event?.type; - const data = event?.data ?? event?.properties ?? {}; - const info = data?.info ?? {}; - const loc = data?.location?.directory ?? data?.directory - ?? event?.location?.directory ?? directory; - if (type === "session.created") { - const id = data?.sessionID ?? data?.id ?? info?.id; - const parentID = info?.parentID ?? data?.parentID; - startSession(id, data?.location?.directory ?? loc, { - title: data?.title ?? info?.title, - projectID: data?.projectID ?? info?.projectID, - // Subagent sessions carry a parentID; forward it as the `agent_id` - // marker so drop_subagent_captures works on OpenCode 2 too (#755). - // Root sessions have no parentID and must stay unmarked. - ...(typeof parentID === "string" && parentID - ? { agent_id: parentID } - : {}), - }); - } - if (type === "session.idle") { - const id = data?.sessionID ?? data?.id; - startSession(id, cwdFor(id, loc)); - postHook("stop", { sessionID: id, cwd: cwdFor(id, loc) }); - } - if (type === "session.deleted") { - const id = data?.sessionID ?? data?.id ?? info?.id; - endSession(id, loc, data?.directory ?? info?.directory); - } - if (type === "session.compaction.started") { - const id = data?.sessionID ?? data?.id; - postPreCompact(id, loc); - } - if (type === "session.compacted") { - const id = data?.sessionID ?? data?.id; - postPreCompact(id, loc); + for await (const event of ctx.event.subscribe({ signal: unload.signal })) { + try { + if (event.type === "session.created") { + await ensureSession(event.data.sessionID); + } else if (event.type === "session.text.ended" && CAPTURE_ASSISTANT) { + if (!(await ensureSession(event.data.sessionID))) continue; + // Bound in-memory buffering too; the native hook owns sanitizing + // and the smaller stored excerpt limit. + if (Buffer.byteLength(event.data.text, "utf8") <= 65536) { + assistantText.set(event.data.sessionID, event.data.text); + } else { + assistantText.delete(event.data.sessionID); + } + } else if (event.type === "session.execution.succeeded" || + event.type === "session.execution.failed" || + event.type === "session.execution.interrupted") { + await postStop(event.data.sessionID); + } else if (event.type === "session.deleted") { + const id = event.data.sessionID; + const session = await sessions.get(id)?.catch(() => undefined); + endSession(id, directory, undefined, subagentMarker(session)); + forgetSession(id); + } else if (event.type === "session.moved") { + const id = event.data.sessionID; + const previousCwd = sessionCwds.get(id); + const previous = await sessions.get(id)?.catch(() => undefined); + // Only the old owner emits the explicit transition. Replayed + // moves keep their ingest key, and stale ordinary events cannot + // silently rebind a live session in the store. + if (previousCwd === directory && event.data.location.directory !== previousCwd) { + postHook("session-start", { + sessionID: id, + cwd: event.data.location.directory, + session_moved: true, + session_moved_from_cwd: previousCwd, + ...subagentMarker(previous), + }); + } + forgetSession(id); + startedSessions.delete(id); + sessionCwds.delete(id); + } else if (event.type === "session.compaction.started") { + const session = await ensureSession(event.data.sessionID); + if (!session) continue; + postPreCompact(event.data.sessionID, session.location.directory); + } + } catch (error) { + if (!unload.signal.aborted) console.warn("ai-memory event capture failed", error); } } - } catch (_e) { - // The stream ends on unload. Capture is best-effort and must never - // break the host. + } catch (error) { + if (!unload.signal.aborted) console.warn("ai-memory event stream failed", error); } })(); - const promptRegistration = await ctx.session.hook("prompt", (event) => { - const e = event as any; - const id = e?.sessionID; - const prompt = e?.prompt ?? {}; - const cwd = cwdFor(id, directory); - startSession(id, cwd, { agent: prompt?.agent, model: prompt?.model }); + const promptRegistration = await ctx.session.hook("prompt", async (event) => { + const session = await ensureSession(event.sessionID); + if (!session) return; + assistantText.delete(event.sessionID); postHook("user-prompt", { - sessionID: id, - cwd, - agent: prompt?.agent, - model: prompt?.model, - messageID: e?.messageID, - prompt: typeof prompt?.text === "string" ? prompt.text : textFromParts(prompt?.parts), + ...hookPayload(event.sessionID, session), + messageID: event.messageID, + prompt: event.prompt.text, }); }); - const handoffInjected = new Set(); const contextRegistration = await ctx.session.hook("context", async (event) => { - const e = event as any; - const id = e?.sessionID; - if (!id || handoffInjected.has(id)) return; - startSession(id, cwdFor(id, directory)); + const id = event.sessionID; + const session = await ensureSession(id); + if (!session || session.parentID) return; let pending = handoffFetches.get(id); if (!pending) { - pending = fetchHandoff(cwdFor(id, directory), id); + // The handoff endpoint establishes scope directly. Do not wait for an + // offline capture backlog on the model's request path. + pending = fetchHandoff(session.location.directory, id); handoffFetches.set(id, pending); } const handoff = await pending; + if (handoff === undefined) handoffFetches.delete(id); if (handoff) { - // SystemPart is `{ type: "text", text }` — the `type` discriminator - // is required: without it the host fails the model call with a - // schema validation error (observed live on beta-18999). - e.system.push({ type: "text", text: handoff }); - handoffInjected.add(id); + event.system.push({ type: "text", text: handoff }); } }); - const beforeRegistration = await ctx.tool.hook("execute.before", (event) => { - const e = event as any; - const id = e?.sessionID; - startSession(id, cwdFor(id, directory)); + const beforeRegistration = await ctx.tool.hook("execute.before", async (event) => { + const session = await ensureSession(event.sessionID); + if (!session) return; postHook("pre-tool-use", { - sessionID: id, - cwd: cwdFor(id, directory), - tool: e?.tool, - callID: e?.id, - args: e?.input, + ...hookPayload(event.sessionID, session), + tool: event.tool, + callID: event.id, + args: event.input, }); }); - const afterRegistration = await ctx.tool.hook("execute.after", (event) => { - const e = event as any; - const id = e?.sessionID; - startSession(id, cwdFor(id, directory)); - // The beta reports failures on the same channel (`status: "error"`, - // no `result`); keep the message where v1 kept output so the - // failure reason survives consolidation. - const failed = e?.status === "error"; + const afterRegistration = await ctx.tool.hook("execute.after", async (event) => { + const session = await ensureSession(event.sessionID); + if (!session) return; postHook("post-tool-use", { - sessionID: id, - cwd: cwdFor(id, directory), - tool: e?.tool, - callID: e?.id, - args: e?.input, - title: failed ? undefined : e?.result?.title, - output: failed - ? String(e?.error?.message ?? e?.error ?? "tool failed") - : e?.result?.output, - metadata: failed ? undefined : e?.result?.metadata, + ...hookPayload(event.sessionID, session), + tool: event.tool, + callID: event.id, + args: event.input, + output: event.status === "error" ? event.error.message + : event.result.content ?? event.result.output, + metadata: event.status === "error" ? undefined : event.result.metadata, }); }); return async () => { - controller.abort(); + unload.abort(); + await Promise.race([eventTask, disposeDrainTimeout()]); for (const registration of [ promptRegistration, contextRegistration, @@ -3300,10 +3450,13 @@ const AiMemoryOpencode2: Plugin = { // Unload is best-effort; a dead host has nothing to unregister. } } - for (const id of Array.from(startedSessions)) { - endSession(id, directory); + for (const [id, pending] of sessions) { + const session = await pending.catch(() => undefined); + if (session) endSession(id, directory, undefined, subagentMarker(session)); } await drainHookQueueForDispose(); + // Whatever is still in flight after the budget fails over to the spool. + hookAbort.abort(); }; }, }; @@ -3611,6 +3764,7 @@ fn build_opencode_plugin( .unwrap_or_else(|| "const TOKEN: string | null = null;\n".to_string()); let apply_marker_params = ts_apply_marker_params(project_strategy); let capture_policy = ts_capture_policy_v1(capture_mode); + let timeout_signal = ts_timeout_signal(); let body = format!( r#"// Auto-generated by `ai-memory install-hooks --agent opencode --apply`. // Edit by re-running the command, not by hand — install-hooks @@ -3628,12 +3782,7 @@ const AGENT = "open-code"; {token_line} {capture_policy} -function timeoutSignal(ms: number): AbortSignal | undefined {{ - if (typeof AbortSignal === "undefined") return undefined; - const factory = (AbortSignal as unknown as {{ timeout?: (ms: number) => AbortSignal }}).timeout; - return factory ? factory(ms) : undefined; -}} - +{timeout_signal} function authHeaders(): Record {{ const token = resolveToken(); return token ? {{ Authorization: `Bearer ${{token}}` }} : {{}}; @@ -3718,7 +3867,7 @@ async function drainHookQueue(): Promise {{ }} catch (_e) {{ // Best-effort capture. Hooks must never block the agent. }} - if (hookQueue.length > 0) await sleep(HOOK_INTER_REQUEST_DELAY_MS); + if (hookQueue.length > 0 && !hookAbort.signal.aborted) await sleep(HOOK_INTER_REQUEST_DELAY_MS); }} }} finally {{ hookDraining = false; @@ -3809,10 +3958,10 @@ function startSession(id: string | undefined, cwd: string, extra: Record = {{}}): void {{ if (!id || !startedSessions.delete(id)) return; const resolvedCwd = cwd || cwdFor(id, directory); - postHook("session-end", {{ sessionID: id, cwd: resolvedCwd }}); + postHook("session-end", {{ sessionID: id, cwd: resolvedCwd, ...extra }}); sessionCwds.delete(id); handoffFetches.delete(id); preCompactLast.delete(id); @@ -4338,6 +4487,7 @@ fn build_omp_extension( .unwrap_or_else(|| "const TOKEN: string | null = null;\n".to_string()); let apply_marker_params = ts_apply_marker_params(project_strategy); let capture_policy = ts_capture_policy_v1(capture_mode); + let timeout_signal = ts_timeout_signal(); let body = format!( r#"// Auto-generated by `ai-memory install-hooks --agent omp --apply`. // Edit by re-running the command, not by hand — install-hooks @@ -4354,12 +4504,7 @@ const AGENT = "omp"; {token_line} {capture_policy} -function timeoutSignal(ms: number): AbortSignal | undefined {{ - if (typeof AbortSignal === "undefined") return undefined; - const factory = (AbortSignal as unknown as {{ timeout?: (ms: number) => AbortSignal }}).timeout; - return factory ? factory(ms) : undefined; -}} - +{timeout_signal} function authHeaders(): Record {{ const token = resolveToken(); return token ? {{ Authorization: `Bearer ${{token}}` }} : {{}}; @@ -4444,7 +4589,7 @@ async function drainHookQueue(): Promise {{ }} catch (_e) {{ // Best-effort capture. Hooks must never block the agent. }} - if (hookQueue.length > 0) await sleep(HOOK_INTER_REQUEST_DELAY_MS); + if (hookQueue.length > 0 && !hookAbort.signal.aborted) await sleep(HOOK_INTER_REQUEST_DELAY_MS); }} }} finally {{ hookDraining = false; @@ -5991,7 +6136,7 @@ mod tests { use tempfile::TempDir; #[test] - fn capture_assistant_allowed_only_for_claude_and_codex_native() { + fn capture_assistant_allowed_only_for_supported_native_agents() { use crate::cli::AgentChoice::*; // Every agent that cannot honor the opt-in is rejected regardless of // platform (#196): the installer must bail rather than enable it silently. @@ -6028,6 +6173,10 @@ mod tests { capture_assistant_allowed(Codex), local_hook_policy_v1_supported() ); + assert_eq!( + capture_assistant_allowed(OpenCode2), + local_hook_policy_v1_supported() + ); } #[test] @@ -8266,6 +8415,17 @@ model = "gpt-5" "pi", build_pi_extension("http://127.0.0.1:49374", Some("tok"), None, "denylist"), ), + ( + "opencode2", + build_opencode2_plugin( + "http://127.0.0.1:49374", + Some("tok"), + None, + "denylist", + false, + ) + .unwrap(), + ), ] { // The fire-and-forget delivery must be gone... assert!( @@ -8291,6 +8451,15 @@ model = "gpt-5" source.contains("if (hookDraining) return;\n requestSpoolDrain();"), "{name}: drainHookQueue must trigger the spool drain" ); + crate::commands::render_shared::assert_shared_ts_delivery_runtime(name, &source); + // Deliveries share one abort: a torn-down state spools the rest + // of its queue at once instead of pacing it out. + assert!( + source.contains( + "if (hookQueue.length > 0 && !hookAbort.signal.aborted) await sleep(HOOK_INTER_REQUEST_DELAY_MS);" + ), + "{name}: an aborted drain must not keep pacing" + ); // The runtime's fs needs made it into the import line. for f in [ "mkdirSync", @@ -8315,8 +8484,11 @@ model = "gpt-5" // timestamp is honoured, and the body survives round-tripping. let created_ms: u64 = 1_756_800_000_000; // Filename exactly as the TS builds it: - // `${String(createdMs).padStart(13, "0")}-${process.pid}-${seq}.json` - let name = format!("{created_ms:013}-4242-{:016x}.json", 7u64); + // `${String(createdMs).padStart(13, "0")}-${process.pid}-${spoolSeqPrefix}${seq}.json` + let name = format!( + "{created_ms:013}-4242-{:08x}{:08x}.json", + 0x9e37_79b9u32, 7u32 + ); // Entry JSON exactly as the TS `JSON.stringify` emits it (key order // irrelevant to serde, but the shapes and enum strings are load-bearing). let json = format!( @@ -8452,9 +8624,14 @@ model = "gpt-5" #[test] fn opencode2_plugin_binds_the_v2_api() { - let plugin = - build_opencode2_plugin("http://127.0.0.1:49374", Some("tok"), None, "denylist") - .unwrap(); + let plugin = build_opencode2_plugin( + "http://127.0.0.1:49374", + Some("tok"), + None, + "denylist", + false, + ) + .unwrap(); // Ownership markers the uninstall gate keys on. assert!(plugin.contains("install-hooks --agent opencode2 --apply")); @@ -8462,35 +8639,39 @@ model = "gpt-5" // Type-only import: `Plugin.define` is identity, and a runtime // import does not resolve from the global plugins dir (the beta // refuses the load). The host only needs the `{ id, setup }` shape. - assert!(plugin.contains("import type { Plugin } from \"@opencode-ai/plugin\";")); + assert!(plugin.contains("import type { Plugin } from \"@opencode/plugin\";")); assert!(!plugin.contains("export default Plugin.define")); assert!(!plugin.contains("Plugin.define({")); - assert!(plugin.contains("const AiMemoryOpencode2: Plugin = {")); + assert!(plugin.contains("const AiMemoryOpencode2: Plugin.Plugin = {")); assert!(plugin.contains("export default AiMemoryOpencode2;")); // #755: subagent sessions forward parentID as the `agent_id` marker, // only when it is a real string (root sessions stay unmarked). - assert!(plugin.contains("const parentID = info?.parentID ?? data?.parentID")); - assert!(plugin.contains("agent_id: parentID")); - assert!(plugin.contains("typeof parentID === \"string\"")); - // Beta hooks, verified against `@opencode-ai/plugin@beta`. + assert!(plugin.contains("agent_id: session.parentID")); + assert!(plugin.contains("if (!session || session.parentID) return")); + // V2 hooks, from `@opencode/plugin@2.0.10`. assert!(plugin.contains("ctx.event.subscribe")); assert!(plugin.contains("ctx.session.hook(\"prompt\"")); assert!(plugin.contains("ctx.session.hook(\"context\"")); assert!(plugin.contains("ctx.tool.hook(\"execute.before\"")); assert!(plugin.contains("ctx.tool.hook(\"execute.after\"")); - // V2 lifecycle envelopes carry `{ type, data, location }`. + // V2 lifecycle envelopes carry `{ type, data, location }`; each of + // these event names is present in the OpenCode 2.0.14 binary. for event in [ "session.created", - "session.idle", + "session.execution.succeeded", + "session.execution.failed", + "session.execution.interrupted", + "session.text.ended", + "session.moved", "session.deleted", - "session.compacted", + "session.compaction.started", ] { assert!(plugin.contains(event), "missing {event}"); } // Handoff injection moved off v1's removed experimental hook and // fires once per session (the context hook runs per model call). // System items are `{ type: "text", text }` objects on the V2 API, not strings. - assert!(plugin.contains("handoffInjected")); + assert!(!plugin.contains("handoffInjected")); assert!(plugin.contains("system.push({ type: \"text\", text: handoff })")); assert!(!plugin.contains("experimental.")); // Hook registrations are disposed on unload so a reload cannot @@ -8501,6 +8682,7 @@ model = "gpt-5" assert!(plugin.contains("session.compaction.started")); // Tool failures share the channel with the reason preserved. assert!(plugin.contains("status === \"error\"")); + assert!(plugin.contains("event.result.content ?? event.result.output")); // No v1 remnants. assert!(!plugin.contains("export default AiMemoryHooks")); assert!(!plugin.contains("chat.message")); @@ -8520,6 +8702,85 @@ model = "gpt-5" // the opencode2 plugin derives from v1 so it inherits that. assert!(plugin.contains("Bearer ${token}")); assert!(plugin.contains("tok")); + // The shared repo-root lookup (memoized, console-less) survives too. + assert!(plugin.contains(TS_REPO_ROOT_PROJECT)); + } + + #[test] + fn opencode2_assistant_capture_is_explicit_and_preserved() { + let off = build_opencode2_plugin("http://localhost:49374", None, None, "denylist", false) + .unwrap(); + let on = + build_opencode2_plugin("http://localhost:49374", None, None, "denylist", true).unwrap(); + assert_eq!(baked_capture_assistant(&off), Some(false)); + assert_eq!(baked_capture_assistant(&on), Some(true)); + // The executable path is baked only with the opt-in, so a plugin + // generated without it is portable between machines. + let exe = std::env::current_exe().unwrap(); + let exe = ts_string_literal(&strip_windows_verbatim_prefix(&exe.to_string_lossy())); + assert!(off.contains("const NATIVE_HOOK: string | undefined = undefined;")); + assert!(!off.contains(&exe)); + assert!(on.contains(&format!("const NATIVE_HOOK: string | undefined = {exe};"))); + assert!(!on.contains(r"\\\\?\\"), "no verbatim path is baked"); + // The spawn inherits no console from OpenCode 2's background service. + let spawn = &on[on.find("spawn(NATIVE_HOOK, args, {").unwrap()..]; + assert!(spawn[..spawn.find("});").unwrap()].contains("windowsHide: true,")); + // A hook that cannot run still records the turn, without its text. + let failed = &spawn[spawn.find("child.on(\"error\"").unwrap()..]; + assert!(failed[..failed.find("});").unwrap()].contains("postHook(\"stop\", payload);")); + assert!(on.contains("last_assistant_message: text")); + assert!(on.contains("child.stdin.end(JSON.stringify")); + assert!(!on.contains("postHook(\"stop\", { ...payload, last_assistant_message")); + assert!( + on.find("setup: async (ctx)").unwrap() < on.find("const hookQueue:").unwrap(), + "capture state belongs to each location instance" + ); + } + + /// Executes the generated plugin against a fake OpenCode 2 host + /// (`tests/fixtures/opencode2-plugin.mjs`) instead of only matching + /// template text: completion events, handoff claims, per-location + /// cleanup, same-millisecond spooling and bounded unload. + #[test] + fn opencode2_plugin_passes_the_node_host_fixture() { + let Some(()) = crate::commands::render_shared::node_strip_types_available() else { + return; + }; + let temp = tempfile::tempdir().unwrap(); + let plugin = temp.path().join("ai-memory-opencode2.ts"); + fs::write( + &plugin, + build_opencode2_plugin("http://127.0.0.1:49600", None, None, "denylist", false) + .unwrap(), + ) + .unwrap(); + // The fixture's sessions live in its cwd. Pin cwd, home and data dir + // inside the tempdir so no operator marker or spool is reachable. + let cwd = temp.path().join("location"); + fs::create_dir(&cwd).unwrap(); + let fixture = Path::new(env!("CARGO_MANIFEST_DIR")) + .join("tests") + .join("fixtures") + .join("opencode2-plugin.mjs"); + let output = std::process::Command::new("node") + .arg("--experimental-strip-types") + .arg(&fixture) + .arg(&plugin) + .current_dir(&cwd) + .env("HOME", temp.path()) + .env("USERPROFILE", temp.path()) + .env("AI_MEMORY_DATA_DIR", temp.path().join("data")) + .env_remove("AI_MEMORY_AUTH_TOKEN") + .env_remove("AI_MEMORY_RUN_ID") + .output() + .unwrap(); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + output.status.success() && stdout.matches("PASS:").count() == 3, + "opencode2 host fixture failed ({})\nstdout:\n{stdout}\nstderr:\n{}", + output.status, + String::from_utf8_lossy(&output.stderr), + ); } #[test] @@ -8531,6 +8792,7 @@ model = "gpt-5" None, Some("repo-root"), "denylist", + false, ) .unwrap(); assert!(plugin.contains("const TOKEN: string | null = null;")); @@ -8548,6 +8810,9 @@ model = "gpt-5" err.to_string().contains("banner"), "drift must name the moved anchor: {err:#}" ); + // A duplicated anchor is drift too: rewriting only the first copy + // would silently leave the other one on v1 behavior. + assert!(replace_opencode_anchor("a a", "a", "b", "banner").is_err()); } #[test] @@ -8555,9 +8820,14 @@ model = "gpt-5" // opencode2 reuses v1's capture prelude verbatim (see // `build_opencode2_plugin`'s doc comment), so the gate must survive // the anchor rewrite into the beta's `{ id, setup }` host binding. - let plugin = - build_opencode2_plugin("http://127.0.0.1:49374", Some("tok"), None, "allowlist") - .unwrap(); + let plugin = build_opencode2_plugin( + "http://127.0.0.1:49374", + Some("tok"), + None, + "allowlist", + false, + ) + .unwrap(); assert!( plugin.contains("const CAPTURE_MODE: \"allowlist\" | \"denylist\" = \"allowlist\";"), "{plugin}" @@ -8567,9 +8837,14 @@ model = "gpt-5" #[test] fn opencode2_plugin_denylist_bakes_inert_gate() { - let plugin = - build_opencode2_plugin("http://127.0.0.1:49374", Some("tok"), None, "denylist") - .unwrap(); + let plugin = build_opencode2_plugin( + "http://127.0.0.1:49374", + Some("tok"), + None, + "denylist", + false, + ) + .unwrap(); assert!( plugin.contains("const CAPTURE_MODE: \"allowlist\" | \"denylist\" = \"denylist\";"), "{plugin}" diff --git a/crates/ai-memory-cli/src/commands/openclaw_plugin.rs b/crates/ai-memory-cli/src/commands/openclaw_plugin.rs index 2eae9ff0..e9475c47 100644 --- a/crates/ai-memory-cli/src/commands/openclaw_plugin.rs +++ b/crates/ai-memory-cli/src/commands/openclaw_plugin.rs @@ -10,6 +10,7 @@ use crate::cli::InstallHooksArgs; use crate::commands::apply_shared::{ApplyOutcome, apply_atomic}; use crate::commands::render_shared::{ ts_capture_policy_v1, ts_resolve_token_fn, ts_spool_runtime, ts_string_literal, + ts_timeout_signal, }; pub(crate) const PLUGIN_ID: &str = "ai-memory"; @@ -340,12 +341,7 @@ const AGENT = "openclaw"; {token_line}{resolve_fn} {capture_policy} -function timeoutSignal(ms: number): AbortSignal | undefined {{ - if (typeof AbortSignal === "undefined") return undefined; - const factory = (AbortSignal as unknown as {{ timeout?: (ms: number) => AbortSignal }}).timeout; - return factory ? factory(ms) : undefined; -}} - +{timeout_signal} function authHeaders(): Record {{ const token = resolveToken(); return token ? {{ Authorization: `Bearer ${{token}}` }} : {{}}; @@ -585,6 +581,7 @@ export default definePluginEntry({{ server_literal = ts_string_literal(server_url), token_line = token_line, repo_root_project = super::install_hooks::TS_REPO_ROOT_PROJECT, + timeout_signal = ts_timeout_signal(), spool_runtime = ts_spool_runtime(), ) } @@ -606,6 +603,7 @@ mod tests { .contains("if (!resp || resp.status >= 500) spoolFailedHook(url, policy.payload);") ); assert!(plugin.contains("else requestSpoolDrain();")); + crate::commands::render_shared::assert_shared_ts_delivery_runtime("openclaw", &plugin); assert!(plugin.contains(r#"return join(env, "hook-spool");"#)); for f in [ "mkdirSync", diff --git a/crates/ai-memory-cli/src/commands/render_shared.rs b/crates/ai-memory-cli/src/commands/render_shared.rs index 5da3749c..8d851f0e 100644 --- a/crates/ai-memory-cli/src/commands/render_shared.rs +++ b/crates/ai-memory-cli/src/commands/render_shared.rs @@ -1883,6 +1883,25 @@ function resolveToken(): string | null { "# } +/// `timeoutSignal`, shared by every generated TypeScript integration: a +/// request deadline composed with `hookAbort`. A host that tears its capture +/// state down aborts `hookAbort` once its final deliveries had their drain +/// budget (OpenCode 2 on each location's unload), so whatever is still in +/// flight fails over to the spool at once instead of each waiting out its own +/// timeout. The module-level integrations never abort it. +pub(crate) fn ts_timeout_signal() -> &'static str { + r#"const hookAbort = new AbortController(); + +function timeoutSignal(ms: number): AbortSignal | undefined { + if (typeof AbortSignal === "undefined") return undefined; + const factory = (AbortSignal as unknown as { timeout?: (ms: number) => AbortSignal }).timeout; + const anyFactory = (AbortSignal as unknown as { any?: (signals: AbortSignal[]) => AbortSignal }).any; + if (!factory) return hookAbort.signal; + return anyFactory ? anyFactory([hookAbort.signal, factory(ms)]) : factory(ms); +} +"# +} + pub(crate) fn ts_spool_runtime() -> &'static str { r#" // ---- offline spool (#580): the same on-disk contract as `ai-memory hook` ---- @@ -1902,6 +1921,11 @@ function hookSpoolDir(): string { return join(base, "hook-spool"); } +// A random prefix per copy of this state keeps two copies in one process +// (OpenCode 2 location instances, OMP and Pi loading each other's extension) +// from renaming onto each other's same-millisecond entry; the counter keeps +// one copy's entries in order. +const spoolSeqPrefix = Math.floor(Math.random() * 0x100000000).toString(16).padStart(8, "0"); let spoolSeq = 0; function spoolFailedHook(url: URL | string, payload: Record): void { @@ -1924,7 +1948,7 @@ function spoolFailedHook(url: URL | string, payload: Record): v ...(token ? { token } : {}), attempts: 0, }; - const seq = (spoolSeq++ & 0xffffffff).toString(16).padStart(16, "0"); + const seq = spoolSeqPrefix + (spoolSeq++ & 0xffffffff).toString(16).padStart(8, "0"); const name = `${String(createdMs).padStart(13, "0")}-${process.pid}-${seq}.json`; const tmp = join(dir, `${name}.tmp`); writeFileSync(tmp, JSON.stringify(entry), { mode: 0o600 }); @@ -1986,6 +2010,60 @@ async function drainHookSpool(): Promise { "# } +/// `Some(())` when this machine can execute the emitted TypeScript. +/// Node is not a build dependency of this project, so a box without it +/// (or on a Node too old for type stripping) skips the runtime evidence +/// instead of failing. +#[cfg(test)] +pub(crate) fn node_strip_types_available() -> Option<()> { + let probe = std::process::Command::new("node") + .args(["--experimental-strip-types", "--version"]) + .output(); + match probe { + Ok(output) if output.status.success() => Some(()), + _ => { + eprintln!( + "skipping Node-required runtime evidence: node lacks --experimental-strip-types" + ); + None + } + } +} + +/// Two copies of the spool state in one process (OpenCode 2 location +/// instances, OMP and Pi sharing an extensions dir) must not write the +/// same file name in the same millisecond, and every request deadline +/// must compose with the state's abort signal. +#[cfg(test)] +pub(crate) fn assert_shared_ts_delivery_runtime(name: &str, source: &str) { + assert!( + source.contains("const spoolSeqPrefix = Math.floor(Math.random() * 0x100000000)"), + "{name}: spool names need a per-state random prefix" + ); + assert!( + source.contains( + "const seq = spoolSeqPrefix + (spoolSeq++ & 0xffffffff).toString(16).padStart(8, \"0\");" + ), + "{name}: spool seq must stay 16 hex digits" + ); + assert_eq!( + source + .matches("const hookAbort = new AbortController();") + .count(), + 1, + "{name}" + ); + assert_eq!( + source.matches("function timeoutSignal(").count(), + 1, + "{name}" + ); + assert!( + source.contains("anyFactory([hookAbort.signal, factory(ms)])"), + "{name}: request deadlines must also honour hookAbort" + ); +} + #[cfg(test)] mod tests { use super::*; @@ -2045,25 +2123,6 @@ mod tests { ) } - /// `Some(())` when this machine can execute the emitted TypeScript. - /// Node is not a build dependency of this project, so a box without it - /// (or on a Node too old for type stripping) skips the runtime evidence - /// instead of failing. - fn node_strip_types_available() -> Option<()> { - let probe = Command::new("node") - .args(["--experimental-strip-types", "--version"]) - .output(); - match probe { - Ok(output) if output.status.success() => Some(()), - _ => { - eprintln!( - "skipping Node-required runtime evidence: node lacks --experimental-strip-types" - ); - None - } - } - } - #[test] fn bearer_header_is_none_when_no_token() { assert!(bearer_header_value(None).is_none()); diff --git a/crates/ai-memory-cli/tests/fixtures/opencode2-plugin.mjs b/crates/ai-memory-cli/tests/fixtures/opencode2-plugin.mjs new file mode 100644 index 00000000..3ad373bf --- /dev/null +++ b/crates/ai-memory-cli/tests/fixtures/opencode2-plugin.mjs @@ -0,0 +1,164 @@ +// Execute the generated adapter, rather than only checking template substrings. +// Run by `cargo test` (opencode2_plugin_passes_the_node_host_fixture), or by hand: +// AI_MEMORY_DATA_DIR= node --experimental-strip-types opencode2-plugin.mjs /absolute/plugin.ts +import assert from "node:assert/strict"; +import { pathToFileURL } from "node:url"; +import { readdirSync, readFileSync } from "node:fs"; +import { join } from "node:path"; + +const requests = []; +globalThis.fetch = async (input, options = {}) => { + // Like a real fetch, a request cancelled before it is sent never arrives. + options.signal?.throwIfAborted(); + const url = new URL(input); + requests.push({ url, payload: options.body ? JSON.parse(options.body) : undefined }); + return new Response(url.pathname === "/handoff" ? "Remember the verified handoff." : "{}"); +}; +const { default: plugin } = await import(pathToFileURL(process.argv[2]).href); + +function host(directory, records) { + const hooks = new Map(); + const events = []; + let wake; + return { + hooks, + emit(event) { + events.push(event); + wake?.(); + }, + ctx: { + location: { directory }, + session: { + async get({ sessionID }) { + const record = records.get(sessionID); + assert.ok(record, `unknown session ${sessionID}`); + return record; + }, + async hook(name, callback) { + hooks.set(`session.${name}`, callback); + return { async dispose() {} }; + }, + }, + tool: { + async hook(name, callback) { + hooks.set(`tool.${name}`, callback); + return { async dispose() {} }; + }, + }, + event: { + async *subscribe({ signal }) { + signal.addEventListener("abort", () => wake?.(), { once: true }); + while (!signal.aborted) { + if (!events.length) await new Promise((resolve) => { wake = resolve; }); + while (events.length) yield events.shift(); + } + }, + }, + }, + }; +} + +async function until(predicate, message) { + const deadline = performance.now() + 5000; + while (!predicate()) { + assert.ok(performance.now() < deadline, message); + await new Promise((resolve) => setTimeout(resolve, 10)); + } +} + +const root = process.cwd(); +const info = (id, parentID) => ({ id, title: id, location: { directory: root }, parentID }); +const records = new Map([["root-a", info("root-a")], ["root-b", { ...info("root-b"), location: { directory: root + "/beta" } }], + ["child", info("child", "root-a")]]); +const a = host(root, records); +const b = host(root + "/beta", records); +const disposeA = await plugin.setup(a.ctx); +const disposeB = await plugin.setup(b.ctx); +try { + // Existing/resumed sessions may never emit session.created after plugin load. + await a.hooks.get("session.prompt")({ sessionID: "root-a", messageID: "m1", prompt: { text: "remember alpha" } }); + const first = { sessionID: "root-a", system: [] }; + const next = { sessionID: "root-a", system: [] }; + await a.hooks.get("session.context")(first); + await a.hooks.get("session.context")(next); + assert.deepEqual(first.system, next.system); + assert.equal(first.system.length, 1); + assert.equal(requests.filter((r) => r.url.pathname === "/handoff").length, 1, "claim once, inject repeatedly"); + + await a.hooks.get("session.prompt")({ sessionID: "child", messageID: "c1", prompt: { text: "child work" } }); + const child = { sessionID: "child", system: [] }; + await a.hooks.get("session.context")(child); + assert.equal(child.system.length, 0); + assert.equal(requests.filter((r) => r.url.pathname === "/handoff").length, 1, "child must not claim"); + + await a.hooks.get("tool.execute.after")({ sessionID: "root-a", tool: "example", id: "t1", input: {}, + status: "completed", result: { content: "content-only tool evidence" } }); + for (const type of ["session.execution.succeeded", "session.execution.failed", "session.execution.interrupted"]) { + a.emit({ type, data: { sessionID: "root-a" } }); + b.emit({ type, data: { sessionID: "root-a" } }); + } + a.emit({ type: "session.execution.succeeded", data: { sessionID: "child" } }); + await until(() => requests.filter((r) => r.url.searchParams.get("event") === "stop").length === 4, "completion events delivered"); + const stops = requests.filter((r) => r.url.searchParams.get("event") === "stop"); + assert.equal(stops.filter((r) => r.payload.turn_checkpoint).length, 3); + assert.equal(stops.find((r) => r.payload.sessionID === "child").payload.agent_id, "root-a"); + assert.ok(requests.some((r) => r.payload?.output === "content-only tool evidence")); + assert.ok(!requests.some((r) => r.url.searchParams.get("event") === "session-end")); + + await b.hooks.get("session.prompt")({ sessionID: "root-b", messageID: "b1", prompt: { text: "beta" } }); + await b.hooks.get("session.context")({ sessionID: "root-b", system: [] }); + await disposeA(); + assert.ok(requests.some((r) => r.url.searchParams.get("event") === "session-end" && r.payload.sessionID === "root-a")); + assert.ok(!requests.some((r) => r.url.searchParams.get("event") === "session-end" && r.payload.sessionID === "root-b"), "location cleanup must not close another instance"); + assert.equal(requests.find((r) => r.url.searchParams.get("event") === "session-end" && r.payload.sessionID === "child").payload.agent_id, "root-a", "child close must retain ancestry to suppress automatic handoffs"); + console.log("PASS: resumed sessions, retained handoff, child isolation, terminal events, content-only output, per-location cleanup"); +} finally { + await disposeB(); +} + +const spool = join(process.env.AI_MEMORY_DATA_DIR, "hook-spool"); +const realNow = Date.now; +globalThis.fetch = async () => { throw new Error("offline fixture"); }; +Date.now = () => 1800000000000; +try { + const offlineA = host(root, records); + const offlineB = host(root + "/beta", records); + const closeA = await plugin.setup(offlineA.ctx); + const closeB = await plugin.setup(offlineB.ctx); + await Promise.all([ + offlineA.hooks.get("session.prompt")({ sessionID: "root-a", messageID: "offline-a", prompt: { text: "offline alpha" } }), + offlineB.hooks.get("session.prompt")({ sessionID: "root-b", messageID: "offline-b", prompt: { text: "offline beta" } }), + ]); + await Promise.all([closeA(), closeB()]); + const entries = readdirSync(spool).filter((name) => name.endsWith(".json")).map((name) => JSON.parse(readFileSync(join(spool, name), "utf8"))); + for (const id of ["root-a", "root-b"]) { + assert.ok(entries.some((entry) => new URL(entry.url).searchParams.get("event") === "session-start" && JSON.parse(entry.body).sessionID === id), `same-millisecond spool must retain ${id}`); + } + assert.ok(entries.every((entry) => new URL(entry.url).searchParams.has("ingest_key")), "stable delivery keys survive spooling"); + console.log("PASS: concurrent location spools do not overwrite same-millisecond events"); +} finally { + Date.now = realNow; +} + +// A server which accepts requests but never answers must not hold plugin unload +// behind every queued request's individual timeout. +globalThis.fetch = (_url, { signal } = {}) => new Promise((_resolve, reject) => { + if (signal.aborted) reject(signal.reason); + else signal.addEventListener("abort", () => reject(signal.reason), { once: true }); +}); +const stalled = host(root, records); +const closeStalled = await plugin.setup(stalled.ctx); +await stalled.hooks.get("session.prompt")({ sessionID: "root-a", messageID: "stalled", prompt: { text: "bounded shutdown" } }); +for (let n = 0; n < 60; n++) { + await stalled.hooks.get("tool.execute.after")({ sessionID: "root-a", tool: "example", id: `pending-${n}`, input: {}, status: "completed", result: { content: "queued evidence" } }); +} +stalled.emit({ type: "session.text.ended", data: { sessionID: "root-a", text: "last response" } }); +stalled.emit({ type: "session.execution.succeeded", data: { sessionID: "root-a" } }); +await new Promise((resolve) => setImmediate(resolve)); +const shutdown = performance.now(); +// The plugin's drain timers are unref'd; a real host process stays alive. +const host_alive = setInterval(() => {}, 1000); +await closeStalled(); +clearInterval(host_alive); +assert.ok(performance.now() - shutdown < 5000, "unload must not wait for the entire network timeout backlog"); +console.log("PASS: unload cancels stalled capture and remains bounded"); diff --git a/crates/ai-memory-hooks/src/assistant_capture.rs b/crates/ai-memory-hooks/src/assistant_capture.rs index c6cf6ad4..200fd023 100644 --- a/crates/ai-memory-hooks/src/assistant_capture.rs +++ b/crates/ai-memory-hooks/src/assistant_capture.rs @@ -46,7 +46,7 @@ const ASSISTANT_MESSAGE_FIELDS: &[&str] = &["last_assistant_message"]; /// The raw field that carries the assistant's final message for `(agent, event)`, /// or `None` when the pair has no verified assistant-message field. /// -/// Closed table: `ClaudeCode + Stop` and `Codex + Stop` are supported today. +/// Closed table: Claude Code, Codex, and OpenCode support `Stop` capture. /// Extend deliberately — a new entry opts an agent/event into capture and MUST /// have its field name present in [`ASSISTANT_MESSAGE_FIELDS`] so the strip /// covers it (enforced by `closed_table_fields_are_all_stripped`). @@ -58,7 +58,7 @@ const ASSISTANT_MESSAGE_FIELDS: &[&str] = &["last_assistant_message"]; #[must_use] pub fn assistant_message_field(agent: AgentKind, event: HookEvent) -> Option<&'static str> { match (agent, event) { - (AgentKind::ClaudeCode, HookEvent::Stop) | (AgentKind::Codex, HookEvent::Stop) => { + (AgentKind::ClaudeCode | AgentKind::Codex | AgentKind::OpenCode, HookEvent::Stop) => { Some("last_assistant_message") } _ => None, @@ -223,10 +223,10 @@ pub fn strip_assistant_message_raw(raw: &mut serde_json::Value) -> bool { mod tests { use super::*; - /// Only `ClaudeCode + Stop` and `Codex + Stop` are capture candidates; every + /// Only Claude Code, Codex, and OpenCode Stop are capture candidates; every /// other agent/event pair across the full agent surface must return `None`. #[test] - fn only_claude_and_codex_stop_are_capture_candidates() { + fn only_supported_stop_events_are_capture_candidates() { let events = [ HookEvent::SessionStart, HookEvent::UserPrompt, @@ -243,8 +243,10 @@ mod tests { ]; for agent in AgentKind::ALL { for event in events { - let expected = matches!(agent, AgentKind::ClaudeCode | AgentKind::Codex) - && event == HookEvent::Stop; + let expected = matches!( + agent, + AgentKind::ClaudeCode | AgentKind::Codex | AgentKind::OpenCode + ) && event == HookEvent::Stop; assert_eq!( assistant_message_field(agent, event).is_some(), expected, @@ -403,6 +405,26 @@ mod tests { assert!(raw.get(ASSISTANT_MARKER_KEY).is_some()); } + #[test] + fn client_transform_captures_opencode_stop_before_spooling() { + let secret = "AKIA".to_string() + &"A".repeat(16); + let mut raw = serde_json::json!({ + "last_assistant_message": format!("Completed work with {secret}"), + "turn_checkpoint": true, + }); + let out = transform_for_client(&mut raw, AgentKind::OpenCode, HookEvent::Stop); + assert!(out.captured && out.changed); + assert!(raw.get("last_assistant_message").is_none()); + assert!(!raw.to_string().contains(&secret)); + assert_eq!(raw["turn_checkpoint"], true); + assert!( + raw[ASSISTANT_MARKER_KEY]["excerpt"] + .as_str() + .unwrap() + .contains("Completed work") + ); + } + #[test] fn backstop_populates_body_when_all_gates_pass() { let mut raw = serde_json::json!({ "last_assistant_message": "done" }); diff --git a/crates/ai-memory-hooks/src/router.rs b/crates/ai-memory-hooks/src/router.rs index 9db28f4a..ad3fcdb2 100644 --- a/crates/ai-memory-hooks/src/router.rs +++ b/crates/ai-memory-hooks/src/router.rs @@ -844,6 +844,7 @@ async fn handle_hook_batch( } let _permit = permit; state.ingest_metrics.record_accepted(); + let (session, agent, event) = (resolve_session_id(&env).ok(), env.agent, env.event); if let Err(e) = process_authorized( &state, env, @@ -857,7 +858,13 @@ async fn handle_hook_batch( e.downcast_ref::(), Some(StoreError::SessionCollision) ) { - warn!("hook batch session collision/recovery rejection dropped"); + warn!( + session = ?session, + agent = %agent.as_str(), + event = ?event, + reason = SESSION_COLLISION_REASON, + "hook batch session collision/recovery rejection dropped" + ); accepted_indices.push(idx); continue; } @@ -2383,6 +2390,13 @@ fn sticky_cwd_admits( && meaningful_session_anchor(session_cwd, home_dir).is_some()) } +/// Why a `SessionCollision` is refused. The store error deliberately carries +/// no row details, so the log names the rule instead: scope is not identity, +/// only the owner and agent of an existing session id are (plus the +/// all-owners recovery authorization checked in `process_authorized`). +const SESSION_COLLISION_REASON: &str = + "session id belongs to another owner or agent, or all-owners recovery was refused"; + /// Returns `true` when the event cleared the writer, so the caller can stamp /// the ingest "last write" metric. A rejected or failed event returns `false`: /// nothing was persisted, and pretending otherwise hides exactly the outage @@ -2394,12 +2408,19 @@ async fn process_envelope( level: ai_memory_core::AuthLevel, skip_webhooks: Vec, ) -> bool { + let (session, agent, event) = (resolve_session_id(&env).ok(), env.agent, env.event); if let Err(e) = process_authorized(&state, env, actor, level, skip_webhooks).await { if matches!( e.downcast_ref::(), Some(StoreError::SessionCollision) ) { - warn!("hook session collision dropped"); + warn!( + session = ?session, + agent = %agent.as_str(), + event = ?event, + reason = SESSION_COLLISION_REASON, + "hook session collision dropped" + ); } else { warn!(error = %e, "hook processing failed"); } @@ -2507,6 +2528,27 @@ async fn process_authorized( skip_webhooks: Vec, ) -> anyhow::Result<()> { let session_id = resolve_session_id(&env)?; + // An OpenCode `session.moved` relocation, forwarded by the plugin as a + // SessionStart naming the directory the session left. Admission rebinds + // the live session row to where it now runs, so its later SessionEnd and + // sticky routing follow it; a child session moves like its root. Managed + // runs are pinned to their own scope and never move. + let moved_from_cwd = (env.agent == AgentKind::OpenCode + && env.event == HookEvent::SessionStart + && env.managed_run.is_none() + && env + .raw + .get("session_moved") + .and_then(serde_json::Value::as_bool) + == Some(true)) + .then(|| { + env.raw + .get("session_moved_from_cwd") + .and_then(serde_json::Value::as_str) + .filter(|cwd| !cwd.is_empty()) + .map(str::to_owned) + }) + .flatten(); // Build the actor key used to scope the in-process `ActiveProject` // pointer. `user` is the qualified storage key of whatever identity the // auth middleware extracted from this request; `session_id` is the RAW @@ -2544,12 +2586,13 @@ async fn process_authorized( // - Under `[routing] mid_session = "sticky"` the session also overrules a // host-derived `repo-root` override, closing the cross-repo `cd` case; // marker-declared scopes still win. See `overrides_permit_sticky`. - let sticky_scope = if overrides_permit_sticky( - env.workspace_override.as_deref(), - env.project_override.as_deref(), - env.project_source, - state.mid_session_routing, - ) { + let sticky_scope = if moved_from_cwd.is_none() + && overrides_permit_sticky( + env.workspace_override.as_deref(), + env.project_override.as_deref(), + env.project_source, + state.mid_session_routing, + ) { state .reader .find_session_scope(session_id) @@ -2675,6 +2718,7 @@ async fn process_authorized( sanitized, owner_filter.clone(), ingest_key.clone(), + moved_from_cwd.clone(), ) .await; match result { @@ -2840,8 +2884,34 @@ async fn process_authorized( // On SessionEnd, close boundary-only sessions without generated artifacts. // Substantive sessions synthesize the summary page and auto-handoff below. - if matches!(env.event, HookEvent::SessionEnd) { + // OpenCode's service outlives its CLI: a completed root turn publishes a + // continuation checkpoint without claiming the native session has ended. + let turn_checkpoint = env.agent == AgentKind::OpenCode + && env.event == HookEvent::Stop + && env + .raw + .get("turn_checkpoint") + .and_then(serde_json::Value::as_bool) + == Some(true) + && !managed + && !body_is_subagent(&env.raw); + if matches!(env.event, HookEvent::SessionEnd) || turn_checkpoint { let mut observations = state.reader.observations_for_session(session_id).await?; + // A checkpoint writes where the session's end will: the session row's + // scope, not wherever this Stop resolved (a marker may have appeared + // mid-session). One that lost the race with the end writes nothing. + let checkpoint_scope = if turn_checkpoint && !is_ephemeral_session(&observations) { + state.reader.open_session_scope(session_id).await? + } else { + None + }; + if turn_checkpoint && checkpoint_scope.is_none() { + if let Some(key) = ingest_key { + state.writer.complete_observation_ingest(proj, key).await?; + } + return Ok(()); + } + let (page_ws, page_proj) = checkpoint_scope.unwrap_or((ws, proj)); if is_ephemeral_session(&observations) { let outcome = state .writer @@ -2880,8 +2950,13 @@ async fn process_authorized( } } } - let new_page = - synthesize_session_page(ws, proj, session_id, admitted.agent_kind(), &observations); + let new_page = synthesize_session_page( + page_ws, + page_proj, + session_id, + admitted.agent_kind(), + &observations, + ); let page_id = state .wiki .write_page(ai_memory_wiki::WritePageRequest { @@ -2912,15 +2987,21 @@ async fn process_authorized( // baton lands in a bucket the operator's actorless transport cannot // read. let handoff_owner = owner_stamp_for_event(state, session_owner.as_ref()).await; - let handoff = (!managed).then(|| { + // An OpenCode child session has its own id and forwards its parent as + // `agent_id`; its end must not hand the next session a sub-task baton. + // Other harnesses are not gated: Claude Code also sets `agent_type` on + // a top-level `--agent` session, which still owns its baton. + let child_session = env.agent == AgentKind::OpenCode && body_is_subagent(&env.raw); + let handoff = (!managed && !child_session).then(|| { build_auto_handoff( - ws, - proj, + page_ws, + page_proj, env.agent, session_id, env.cwd.clone(), &observations, handoff_owner, + turn_checkpoint, ) }); // Automatic SessionEnd handoffs are the bulk of handoff traffic; @@ -2945,8 +3026,8 @@ async fn process_authorized( match state .wiki .authorize_operation( - ws, - proj, + page_ws, + page_proj, ai_memory_wiki::AdmissionOp::HandoffBegin, session_actor.clone(), skip_webhooks, @@ -2958,7 +3039,7 @@ async fn process_authorized( warn!( session = %session_id, error = %e, - "auto handoff refused by admission chain; session ends without a baton", + "auto handoff refused by admission chain; continuing without a baton", ); (None, None) } @@ -2970,19 +3051,26 @@ async fn process_authorized( // between them would leave an ended session whose successor has // nothing to pick up. A managed run and an admission refusal both take // the second arm, ending the session with no handoff at all. - let handoff_id = match handoff { - Some(handoff) => Some( - state - .writer - .end_admitted_session_with_handoff(admitted.clone(), Some(page_id), handoff) - .await?, - ), - None => { - state - .writer - .end_admitted_session(admitted.clone(), Some(page_id)) - .await?; - None + let handoff_id = if turn_checkpoint { + match handoff { + Some(handoff) => state.writer.checkpoint_session_handoff(handoff).await?, + None => None, + } + } else { + match handoff { + Some(handoff) => Some( + state + .writer + .end_admitted_session_with_handoff(admitted.clone(), Some(page_id), handoff) + .await?, + ), + None => { + state + .writer + .end_admitted_session(admitted.clone(), Some(page_id)) + .await?; + None + } } }; if handoff_id.is_some() @@ -3006,30 +3094,35 @@ async fn process_authorized( // deterministic wiki writes are committed so the worker cannot race // their git snapshot. Stale redelivery above repairs cancellation in // the narrow window after `end_session`. - enqueue_session_end_consolidation(state, session_id, ws, proj).await?; - if let Some(handoff_id) = handoff_id { - info!( - session = %session_id, - page = %new_page.path, - handoff = %handoff_id, - "session ended; summary page + open handoff created", - ); - } else if managed { - info!( - session = %session_id, - page = %new_page.path, - managed_run = ?managed_run, - "managed session ended; summary page written without duplicate legacy handoff", - ); + if turn_checkpoint { + info!(session = %session_id, page = %new_page.path, "turn checkpoint written; native session remains open"); } else { - // Only reachable through the admission refusal above, which already - // warned with the reason; without this arm the refusal would be - // logged as a managed session end and hide why the baton is gone. - info!( - session = %session_id, - page = %new_page.path, - "session ended; summary page written without a handoff (admission refused)", - ); + enqueue_session_end_consolidation(state, session_id, ws, proj).await?; + if let Some(handoff_id) = handoff_id { + info!( + session = %session_id, + page = %new_page.path, + handoff = %handoff_id, + "session ended; summary page + open handoff created", + ); + } else if managed || child_session { + info!( + session = %session_id, + page = %new_page.path, + managed_run = ?managed_run, + "managed or child session ended; summary page written without legacy handoff", + ); + } else { + // Only reachable through the admission refusal above, which + // already warned with the reason; without this arm the refusal + // would be logged as a managed session end and hide why the + // baton is gone. + info!( + session = %session_id, + page = %new_page.path, + "session ended; summary page written without a handoff (admission refused)", + ); + } } } @@ -3087,6 +3180,7 @@ fn is_ephemeral_session(observations: &[ai_memory_core::Observation]) -> bool { }) } +#[allow(clippy::too_many_arguments)] fn build_auto_handoff( workspace_id: WorkspaceId, project_id: ProjectId, @@ -3095,6 +3189,7 @@ fn build_auto_handoff( cwd: Option, observations: &[ai_memory_core::Observation], owner_user: Option, + turn_checkpoint: bool, ) -> NewHandoff { // Prefer obs.body (the full prompt) over obs.title (first-line + // truncated to 80 chars for log/list display). When body is @@ -3136,16 +3231,42 @@ fn build_auto_handoff( } let first_prompt = prompts.first().cloned(); let last_prompt = prompts.last().cloned(); - let summary = match (&first_prompt, &last_prompt) { + let mut summary = match (&first_prompt, &last_prompt) { (Some(first), Some(last)) if first == last => format!("Session focused on: {}", cap(first)), (Some(first), Some(last)) => format!("Started: {}\n\nLast: {}", cap(first), cap(last),), (Some(first), None) => format!("Started: {}", cap(first)), _ => format!( - "Session ended; {} observations recorded.", + "{}; {} observations recorded.", + if turn_checkpoint { + "Turn checkpoint" + } else { + "Session ended" + }, observations.len() ), }; - let open_questions = derive_open_questions(observations, &last_prompt); + // Stop bodies exist only after the assistant-capture double opt-in. The + // excerpt continues the work in the next session's baton and stays out of + // the git-tracked session page. + if let Some(assistant) = observations.iter().rev().find(|observation| { + observation.kind == ObservationKind::Stop && !observation.body.trim().is_empty() + }) { + summary.push_str("\n\nLatest assistant response: "); + summary.push_str(&assistant.body); + } + let open_questions = if turn_checkpoint { + // A completed turn is neither a native-session exit nor evidence that + // the user's last question remains unanswered. + last_prompt + .as_deref() + .map(str::trim) + .filter(|prompt| !prompt.is_empty() && !is_acknowledgment(prompt)) + .map(|prompt| format!("Continue from last request: {}", cap_handoff_text(prompt))) + .into_iter() + .collect() + } else { + derive_open_questions(observations, &last_prompt) + }; let next_steps = if tools.is_empty() { Vec::new() } else { @@ -8347,6 +8468,598 @@ mod tests { ); } + fn opencode_turn_event(session: &str, event: &str, text: &str) -> HookEnvelope { + HookEnvelope::from_query_and_body( + HookQuery { + event: event.into(), + agent: Some("opencode2".into()), + capture_assistant: Some("1".into()), + ..Default::default() + }, + serde_json::json!({ + "session_id": session, + "prompt": text, + "turn_checkpoint": true, + "_ai_memory_assistant": { "version": 1, "excerpt": text }, + }), + ) + } + + #[tokio::test] + async fn opencode_turn_checkpoints_refresh_page_and_baton_without_ending_session() { + let tmp = TempDir::new().unwrap(); + let mut state = make_state(&tmp).await; + state.consolidate_on_session_end = true; + let llm = Arc::new(RecordingLlm(Mutex::new(None))); + state.consolidator = Some(Arc::new(Consolidator::new( + state.reader.clone(), + state.writer.clone(), + state.wiki.clone(), + llm.clone(), + state.workspace_id, + state.project_id, + ))); + state.session_consolidation_notify = Some(Arc::new(tokio::sync::Notify::new())); + let session = SessionId::new(); + let mut previous_page = None; + let mut previous_handoff = None; + for text in ["Implemented the first turn", "Verified the second turn"] { + process( + &state, + opencode_turn_event(&session.to_string(), "user-prompt", text), + None, + Vec::new(), + ) + .await + .unwrap(); + let mut stop = opencode_turn_event(&session.to_string(), "stop", text); + crate::assistant_capture::apply_assistant_backstop(&mut stop, true); + process(&state, stop, None, Vec::new()).await.unwrap(); + let pages = state + .reader + .recent_pages_for_project(state.workspace_id, state.project_id, 20) + .await + .unwrap(); + let page = pages + .iter() + .find(|page| page.path.as_str().starts_with("sessions/")) + .unwrap(); + let body = state + .wiki + .read_page(state.workspace_id, state.project_id, &page.path) + .unwrap() + .body; + assert!(body.contains(text)); + assert_ne!(previous_page, Some(page.id)); + previous_page = Some(page.id); + let handoff = state + .reader + .latest_open_handoff( + state.workspace_id, + state.project_id, + None, + ai_memory_core::OwnerFilter::Any, + ) + .await + .unwrap() + .unwrap(); + assert!( + handoff + .content + .summary + .contains(&format!("Latest assistant response: {text}")) + ); + if let Some(previous) = previous_handoff { + assert_eq!( + previous, handoff.scope.id, + "a live session keeps one baton, refreshed in place" + ); + } + assert!( + handoff + .content + .open_questions + .iter() + .all(|question| !question.contains("exit")) + ); + previous_handoff = Some(handoff.scope.id); + assert!( + state + .reader + .latest_completed_session_for_project(state.workspace_id, state.project_id) + .await + .unwrap() + .is_none() + ); + assert_eq!( + state + .reader + .session_end_disposition( + session, + state.workspace_id, + state.project_id, + AgentKind::OpenCode + ) + .await + .unwrap(), + ai_memory_store::SessionEndDisposition::Open + ); + } + let now = Timestamp::now().as_microsecond(); + assert!( + state + .writer + .claim_session_consolidation(now, now - 1) + .await + .unwrap() + .is_none() + ); + assert!(llm.0.lock().unwrap().is_none()); + // Each turn replaces the session's own unclaimed baton rather than + // leaving an expired row behind per turn. + let batons = state + .reader + .list_handoffs( + state.workspace_id, + state.project_id, + None, + ai_memory_core::OwnerFilter::Any, + 10, + ) + .await + .unwrap(); + assert_eq!(batons.len(), 1, "one baton row per live session"); + + // A real end still closes the same resumable native session. + process( + &state, + opencode_turn_event(&session.to_string(), "session-end", ""), + None, + Vec::new(), + ) + .await + .unwrap(); + assert_eq!( + state + .reader + .latest_completed_session_for_project(state.workspace_id, state.project_id) + .await + .unwrap(), + Some(session) + ); + } + + #[tokio::test] + async fn opencode_turn_checkpoint_requires_root_substantive_unmanaged_marked_stop() { + for case in [ + "noop", + "subagent", + "managed", + "unmarked", + "wrong-agent", + "wrong-type", + ] { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + let session = SessionId::new().to_string(); + if case != "noop" { + let mut prompt = opencode_turn_event(&session, "user-prompt", "Real work"); + if case == "wrong-agent" { + prompt.agent = AgentKind::Codex; + } + process(&state, prompt, None, Vec::new()).await.unwrap(); + } + let mut stop = opencode_turn_event(&session, "stop", "Completed work"); + match case { + "subagent" => stop.raw["agent_id"] = serde_json::json!("child"), + "managed" => stop.managed_run = Some(ManagedRunId::new().to_string()), + "unmarked" => stop.raw["turn_checkpoint"] = serde_json::json!(false), + "wrong-type" => stop.raw["turn_checkpoint"] = serde_json::json!("true"), + "wrong-agent" => stop.agent = AgentKind::Codex, + _ => {} + } + process(&state, stop, None, Vec::new()).await.unwrap(); + assert!(session_pages(&state).await.is_empty(), "{case}"); + assert!(!open_handoff_exists(&state).await, "{case}"); + } + } + + #[tokio::test] + async fn opencode_turn_checkpoint_admission_refusal_keeps_page_and_live_session() { + let tmp = TempDir::new().unwrap(); + let state = make_state_with_admission( + &tmp, + refusing_admission_chain("guard", vec![ai_memory_wiki::AdmissionOp::HandoffBegin]), + ) + .await; + let session = SessionId::new().to_string(); + for event in ["user-prompt", "stop"] { + process( + &state, + opencode_turn_event(&session, event, "Real work"), + None, + Vec::new(), + ) + .await + .unwrap(); + } + assert!(!session_pages(&state).await.is_empty()); + assert!(!open_handoff_exists(&state).await); + assert!( + state + .reader + .latest_completed_session_for_project(state.workspace_id, state.project_id) + .await + .unwrap() + .is_none() + ); + } + + #[tokio::test] + async fn opencode_child_session_end_closes_and_summarizes_without_baton() { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + let session = SessionId::new(); + for event in ["user-prompt", "stop", "session-end"] { + let mut env = opencode_turn_event(&session.to_string(), event, "Child captured work"); + env.raw["agent_id"] = serde_json::json!("child"); + process(&state, env, None, Vec::new()).await.unwrap(); + } + assert!(!session_pages(&state).await.is_empty()); + assert!(!open_handoff_exists(&state).await); + assert_eq!( + state + .reader + .latest_completed_session_for_project(state.workspace_id, state.project_id) + .await + .unwrap(), + Some(session) + ); + } + + // Claude Code stamps `agent_type` on a top-level `--agent` session, so + // that marker must not cost it its baton; and a captured Stop body is + // surfaced for any agent in the capture table, not only OpenCode. + #[tokio::test] + async fn claude_code_agent_session_end_keeps_baton_and_latest_response() { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + let session = SessionId::new().to_string(); + for event in ["user-prompt", "stop", "session-end"] { + let mut env = HookEnvelope::from_query_and_body( + HookQuery { + event: event.into(), + agent: Some("claude-code".into()), + capture_assistant: Some("1".into()), + ..Default::default() + }, + serde_json::json!({ + "session_id": session, + "agent_type": "reviewer", + "prompt": "Review the patch", + "_ai_memory_assistant": { "version": 1, "excerpt": "Found two bugs" }, + }), + ); + crate::assistant_capture::apply_assistant_backstop(&mut env, true); + process(&state, env, None, Vec::new()).await.unwrap(); + } + let handoff = state + .reader + .latest_open_handoff( + state.workspace_id, + state.project_id, + None, + ai_memory_core::OwnerFilter::Any, + ) + .await + .unwrap() + .expect("a top-level --agent session keeps its baton"); + assert!( + handoff + .content + .summary + .contains("Latest assistant response: Found two bugs") + ); + let pages = session_pages(&state).await; + let body = state + .wiki + .read_page( + state.workspace_id, + state.project_id, + &ai_memory_core::PagePath::new(pages[0].clone()).unwrap(), + ) + .unwrap() + .body; + assert!( + !body.contains("Found two bugs"), + "the assistant excerpt stays out of the git-tracked session page" + ); + } + + #[tokio::test] + async fn opencode_turn_checkpoint_assistant_capture_requires_both_opt_ins() { + for (server, client) in [(false, false), (false, true), (true, false), (true, true)] { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + let session = SessionId::new().to_string(); + process( + &state, + opencode_turn_event(&session, "user-prompt", "Fix the bug"), + None, + Vec::new(), + ) + .await + .unwrap(); + let secret = "AKIA".to_string() + &"A".repeat(16); + let mut stop = + opencode_turn_event(&session, "stop", &format!("Assistant-only result {secret}")); + stop.capture_assistant_requested = client; + crate::assistant_capture::apply_assistant_backstop(&mut stop, server); + process(&state, stop, None, Vec::new()).await.unwrap(); + let handoff = state + .reader + .latest_open_handoff( + state.workspace_id, + state.project_id, + None, + ai_memory_core::OwnerFilter::Any, + ) + .await + .unwrap() + .unwrap(); + assert_eq!( + handoff.content.summary.contains("Assistant-only result"), + server && client + ); + assert!(!handoff.content.summary.contains(&secret)); + let pages = state + .reader + .recent_pages_for_project(state.workspace_id, state.project_id, 20) + .await + .unwrap(); + let page = pages + .iter() + .find(|page| page.path.as_str().starts_with("sessions/")) + .unwrap(); + let body = state + .wiki + .read_page(state.workspace_id, state.project_id, &page.path) + .unwrap() + .body; + assert!(!body.contains("Assistant-only result")); + } + } + + #[tokio::test] + async fn opencode_turn_checkpoint_handoffs_are_project_scoped_and_claimed_once() { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + for project in ["alpha", "beta"] { + let session = SessionId::new().to_string(); + for turn in ["first", "latest"] { + for event in ["user-prompt", "stop"] { + let mut env = + opencode_turn_event(&session, event, &format!("{project}-{turn}")); + env.workspace_override = Some("checkpoint-workspace".into()); + env.project_override = Some(project.into()); + crate::assistant_capture::apply_assistant_backstop(&mut env, true); + process(&state, env, None, Vec::new()).await.unwrap(); + } + } + } + for (project, other) in [("alpha", "beta"), ("beta", "alpha")] { + let query = HandoffQuery { + workspace: Some("checkpoint-workspace".into()), + project: Some(project.into()), + agent: Some("codex".into()), + session_id: Some(SessionId::new().to_string()), + ..Default::default() + }; + let first = fetch_and_accept_handoff(&state, query.clone(), None, Vec::new()) + .await + .unwrap() + .unwrap(); + assert!(first.contains(&format!("Latest assistant response: {project}-latest"))); + assert!(!first.contains(&format!("{other}-"))); + let second = fetch_and_accept_handoff(&state, query, None, Vec::new()) + .await + .unwrap(); + assert!( + second.is_none(), + "a consumed or superseded baton must not reappear" + ); + } + } + + // Live incident: a `.ai-memory.toml` naming a workspace appeared under + // running OpenCode sessions, so their later events (same cwd) resolved to + // a new `(workspace, project)`. Scope is not identity: those events are + // recorded in the scope they name instead of being dropped as a session + // collision, the session row keeps the scope it began in, and the + // drifted SessionEnd still ends it. + #[tokio::test] + async fn marker_added_mid_session_records_later_events_and_still_ends() { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + let session = SessionId::new(); + let event = |name: &str, workspace: Option<&str>| { + HookEnvelope::from_query_and_body( + HookQuery { + event: name.into(), + agent: Some("opencode2".into()), + cwd: Some("/work/projects".into()), + workspace: workspace.map(str::to_owned), + ..Default::default() + }, + serde_json::json!({ + "session_id": session.to_string(), + "prompt": "keep working", + "tool_name": "bash", + "turn_checkpoint": name == "stop", + }), + ) + }; + let admit = |env| { + process_authorized( + &state, + env, + None, + ai_memory_core::AuthLevel::Anonymous, + Vec::new(), + ) + }; + for name in ["session-start", "user-prompt"] { + admit(event(name, None)).await.unwrap(); + } + let (began_ws, began_proj, _) = state + .reader + .find_session_scope(session) + .await + .unwrap() + .unwrap(); + for name in ["user-prompt", "pre-tool-use", "stop"] { + admit(event(name, Some("windows"))).await.unwrap(); + } + + let observations = state + .reader + .observations_for_session(session) + .await + .unwrap(); + assert_eq!(observations.len(), 5, "no event may be dropped"); + let marker_scoped = observations + .iter() + .filter(|obs| obs.workspace_id != began_ws && obs.project_id != began_proj) + .count(); + assert_eq!( + marker_scoped, 3, + "post-marker events land where they resolve" + ); + let (row_ws, row_proj, _) = state + .reader + .find_session_scope(session) + .await + .unwrap() + .unwrap(); + assert_eq!((row_ws, row_proj), (began_ws, began_proj)); + + // The drifted Stop was a turn checkpoint. Its page and baton belong to + // the session, so they land where its end will write them, and a purge + // of the session's scope can find them. + let marker = observations.last().unwrap(); + let (marker_ws, marker_proj) = (marker.workspace_id, marker.project_id); + let session_page = |ws, proj| { + let reader = state.reader.clone(); + async move { + reader + .recent_pages_for_project(ws, proj, 20) + .await + .unwrap() + .into_iter() + .any(|page| page.path.as_str().starts_with("sessions/")) + } + }; + let open_baton = |ws, proj| { + let reader = state.reader.clone(); + async move { + reader + .latest_open_handoff(ws, proj, None, ai_memory_core::OwnerFilter::Any) + .await + .unwrap() + .is_some() + } + }; + assert!(session_page(began_ws, began_proj).await); + assert!(open_baton(began_ws, began_proj).await); + assert!(!session_page(marker_ws, marker_proj).await); + assert!(!open_baton(marker_ws, marker_proj).await); + + // The plugin's end also resolves to the marker scope. It comes from + // the session's own cwd, so it ends the session where it began + // instead of stranding it open. + admit(event("session-end", Some("windows"))).await.unwrap(); + assert_eq!( + state + .reader + .latest_completed_session_for_project(began_ws, began_proj) + .await + .unwrap(), + Some(session) + ); + } + + // OpenCode `session.moved`: the plugin forwards the relocation as a keyed + // SessionStart naming the directory the session left. The live row + // follows it, so the SessionEnd sent from the new directory ends the + // session there instead of being ignored as a foreign-scope end. A child + // session (it carries its parent as `agent_id`) moves the same way. + #[tokio::test] + async fn opencode_native_move_lets_the_session_end_where_it_now_runs() { + for child in [false, true] { + let tmp = TempDir::new().unwrap(); + let state = make_state(&tmp).await; + let session = SessionId::new(); + let event = |project: &str, name: &str| { + let mut body = serde_json::json!({ + "session_id": session.to_string(), + "prompt": format!("{project} work"), + }); + if child { + body["agent_id"] = serde_json::json!("parent"); + } + HookEnvelope::from_query_and_body( + HookQuery { + event: name.into(), + agent: Some("opencode2".into()), + cwd: Some(format!("/repo/{project}")), + workspace: Some("move-workspace".into()), + project: Some(project.into()), + ..Default::default() + }, + body, + ) + }; + for name in ["session-start", "user-prompt"] { + process(&state, event("alpha", name), None, Vec::new()) + .await + .unwrap(); + } + let mut moved = event("beta", "session-start"); + moved.raw["session_moved"] = serde_json::json!(true); + moved.raw["session_moved_from_cwd"] = serde_json::json!("/repo/alpha"); + moved.ingest_key = Some("alpha-to-beta".into()); + process(&state, moved, None, Vec::new()).await.unwrap(); + for name in ["user-prompt", "session-end"] { + process(&state, event("beta", name), None, Vec::new()) + .await + .unwrap(); + } + + let (ws, project, cwd) = state + .reader + .find_session_scope(session) + .await + .unwrap() + .unwrap(); + assert_eq!(cwd.as_deref(), Some("/repo/beta"), "child: {child}"); + let beta = state + .writer + .get_or_create_project(ws, "beta", None) + .await + .unwrap(); + assert_eq!(project, beta, "child: {child}"); + assert_eq!( + state + .reader + .latest_completed_session_for_project(ws, beta) + .await + .unwrap(), + Some(session), + "child: {child}" + ); + } + } + #[tokio::test] async fn session_end_closes_only_matching_scoped_session() { let tmp = TempDir::new().unwrap(); @@ -9816,6 +10529,7 @@ mod tests { None, &observations, None, + false, ); state .writer diff --git a/crates/ai-memory-store/src/ops.rs b/crates/ai-memory-store/src/ops.rs index 4307c582..91efdc7a 100644 --- a/crates/ai-memory-store/src/ops.rs +++ b/crates/ai-memory-store/src/ops.rs @@ -99,9 +99,10 @@ pub enum HookSessionAdmission { }, /// A terminal event named no persisted session and created nothing. InvalidMissingEnd, - /// A terminal event named a persisted session in a different scope, so it - /// is not that session's end. Mirrors the pre-guard - /// `SessionEndDisposition::DropInvalid` arm. + /// A terminal event named a persisted session in a different scope and + /// from a different or missing cwd, so it is not that session's end. A + /// same-cwd end is admitted in the session's own scope instead. Mirrors + /// the pre-guard `SessionEndDisposition::DropInvalid` arm. InvalidScopedEnd, } /// Result of conditionally ending a session whose persisted observations are @@ -1658,12 +1659,17 @@ pub fn insert_observation_keyed( /// Find or create the hook session, validate its immutable tuple and owner, /// optionally claim an ingest key, and append the observation in one writer /// transaction. Validation always precedes key mutation. +/// +/// `moved_from_cwd` marks an explicit native relocation (OpenCode +/// `session.moved`): when the stored cwd still matches it, the live session +/// row is rebound to this event's scope and cwd in the same transaction. pub fn admit_hook_session_event( conn: &mut Connection, session: &NewSession, obs: &NewObservation, owner_filter: &OwnerFilter, ingest_key: Option<&str>, + moved_from_cwd: Option<&str>, ) -> StoreResult { if obs.session_id != session.id || obs.workspace_id != session.workspace_id @@ -1676,14 +1682,22 @@ pub fn admit_hook_session_event( let session_end = obs.kind == ObservationKind::SessionEnd; let now = Timestamp::now().as_microsecond(); let tx = conn.transaction()?; - type Row = (Vec, Vec, String, Option, Option, u64); + type Row = ( + Vec, + Vec, + String, + Option, + Option, + u64, + Option, + ); let existing: Option = tx.query_row( - "SELECT workspace_id, project_id, agent_kind, actor_user, ended_at, ended_observation_count FROM sessions WHERE id = ?1", + "SELECT workspace_id, project_id, agent_kind, actor_user, ended_at, ended_observation_count, cwd FROM sessions WHERE id = ?1", params![session.id.as_bytes()], - |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?, r.get(5)?)), + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?, r.get(5)?, r.get(6)?)), ).optional()?; - let (owner, ended_at, ended_count) = match existing { - Some((ws, project, agent, owner, ended_at, ended_count)) => { + let (owner, ended_at, ended_count, stored_cwd, scope) = match existing { + Some((ws, project, agent, owner, ended_at, ended_count, cwd)) => { // Corrupt owners fail closed, including for Any recovery. Owner // and agent identify WHO the session belongs to, so a mismatch // there is a genuine UUID collision and is terminal. @@ -1703,17 +1717,35 @@ pub fn admit_hook_session_event( // exists to opt out of. Treating the difference as a collision // silently DROPPED those events instead of recording them. // - // A terminal event is the one exception: an end naming a different - // scope is not this session's end, so it is dropped rather than - // ending someone else's session (the pre-guard - // `SessionEndDisposition::DropInvalid` arm). + // A terminal event is the exception: a session ends in the scope + // it was recorded in. An end naming a different scope still ends + // it when it comes from the session's own cwd — the scope drifted + // under the same directory, e.g. a `.ai-memory.toml` appeared + // mid-session. Otherwise it is not this session's end and is + // dropped (the pre-guard `SessionEndDisposition::DropInvalid` + // arm). Owner and agent were already checked above, so a foreign + // operator or agent never gets this far. let scoped_to_session = ws.as_slice() == session.workspace_id.as_bytes() && project.as_slice() == session.project_id.as_bytes(); - if session_end && !scoped_to_session { + let same_cwd = matches!( + (cwd.as_deref(), session.cwd.as_deref()), + (Some(stored), Some(event)) + if crate::reader::normalize_cwd(stored) + == crate::reader::normalize_cwd(&event.to_string_lossy()) + ); + if session_end && !scoped_to_session && !same_cwd { tx.commit()?; return Ok(HookSessionAdmission::InvalidScopedEnd); } - (owner, ended_at, ended_count) + let scope = if session_end { + ( + WorkspaceId::from_slice(&ws)?, + ProjectId::from_slice(&project)?, + ) + } else { + (session.workspace_id, session.project_id) + }; + (owner, ended_at, ended_count, cwd, scope) } None if session_end => { tx.commit()?; @@ -1728,13 +1760,31 @@ pub fn admit_hook_session_event( "INSERT INTO sessions (id, workspace_id, project_id, agent_kind, cwd, started_at, actor_user) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)", params![session.id.as_bytes(), session.workspace_id.as_bytes(), session.project_id.as_bytes(), session.agent_kind.as_str(), session.cwd.as_ref().map(|p| p.to_string_lossy().into_owned()), now, session.actor_user.as_deref()], )?; - (session.actor_user.clone(), None, 0) + ( + session.actor_user.clone(), + None, + 0, + None, + (session.workspace_id, session.project_id), + ) } }; + let (workspace_id, project_id) = scope; + let rescoped; + let obs = if (obs.workspace_id, obs.project_id) == scope { + obs + } else { + rescoped = NewObservation { + workspace_id, + project_id, + ..obs.clone() + }; + &rescoped + }; let guard = AdmittedSession { session_id: session.id, - workspace_id: session.workspace_id, - project_id: session.project_id, + workspace_id, + project_id, agent_kind: session.agent_kind, owner, }; @@ -1777,6 +1827,29 @@ pub fn admit_hook_session_event( } else { IngestObservationOutcome::Inserted(insert_observation_row(&tx, obs)?) }; + // Compare-and-set on the source cwd so a stale or out-of-order move never + // rebinds a newer location. Only the first admission of a keyed move may + // rebind: a redelivery (complete or resumed) already had its chance in the + // transaction that claimed the key, and replaying it after A -> B -> A + // would drag the session back. An unkeyed move cannot prove it is not + // such a replay. Earlier observations keep their scope. + if let (Some(from), Some(stored), Some(target)) = + (moved_from_cwd, stored_cwd.as_deref(), session.cwd.as_ref()) + && ingest_key.is_some() + && ended_at.is_none() + && matches!(ingest, IngestObservationOutcome::Inserted(_)) + && crate::reader::normalize_cwd(stored) == crate::reader::normalize_cwd(from) + { + tx.execute( + "UPDATE sessions SET workspace_id = ?1, project_id = ?2, cwd = ?3 WHERE id = ?4", + params![ + session.workspace_id.as_bytes(), + session.project_id.as_bytes(), + target.to_string_lossy(), + session.id.as_bytes() + ], + )?; + } tx.commit()?; if !session_end { Ok(HookSessionAdmission::Observation { @@ -2700,6 +2773,72 @@ pub fn insert_handoff(conn: &mut Connection, h: &NewHandoff) -> StoreResult StoreResult> { + let Some(session_id) = handoff.from_session_id.as_ref() else { + return Err(StoreError::InvalidState( + "checkpoint handoff has no source session".into(), + )); + }; + let tx = conn.transaction()?; + let live: Option<(Vec, Vec, Option)> = tx + .query_row( + "SELECT workspace_id, project_id, actor_user FROM sessions \ + WHERE id = ?1 AND ended_at IS NULL", + params![session_id.as_bytes()], + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), + ) + .optional()?; + let Some((workspace_id, project_id, actor_user)) = live else { + return Ok(None); + }; + if workspace_id.as_slice() != handoff.workspace_id.as_bytes() + || project_id.as_slice() != handoff.project_id.as_bytes() + || actor_user != handoff.owner_user + { + return Err(StoreError::InvalidState( + "checkpoint handoff scope or owner does not match its session".into(), + )); + } + let existing: Option> = tx + .query_row( + "SELECT id FROM handoffs \ + WHERE from_session_id = ?1 AND workspace_id = ?2 AND project_id = ?3 \ + AND owner_user IS ?4 AND state = 'open' \ + ORDER BY created_at DESC LIMIT 1", + params![ + session_id.as_bytes(), + workspace_id, + project_id, + actor_user.as_deref() + ], + |r| r.get(0), + ) + .optional()?; + let id = match existing { + Some(id) => { + let id = HandoffId::from_slice(&id)?; + refresh_handoff_row(&tx, &id, handoff)?; + id + } + None => insert_handoff_row(&tx, handoff)?, + }; + tx.commit()?; + Ok(Some(id)) +} + /// Atomically stamp a session ended and insert its automatic handoff. /// /// A failed handoff insert rolls the end stamp back, so a keyed retry can run @@ -2765,68 +2904,133 @@ fn bound_handoff_list(items: &[String]) -> Vec { .collect() } -fn insert_handoff_row(conn: &Transaction<'_>, h: &NewHandoff) -> StoreResult { - validate_identity_storage_key(h.owner_user.as_deref(), "handoff owner")?; - let id = HandoffId::new(); - let now = Timestamp::now().as_microsecond(); +/// A handoff's prose and cwd in their stored form. +struct HandoffFields { + summary: String, + open_questions: String, + next_steps: String, + files_touched: String, + cwd: Option, +} + +fn handoff_fields(h: &NewHandoff) -> StoreResult { // Store-boundary bound (defense in depth): the MCP/hook callers already // scrub and cap handoff prose, but the store is the last gate before // durable persistence — mirror the observation body's bound so a caller // that ever forgets cannot write unbounded content to the DB. Generous // enough never to fire under the callers' tighter caps. - let summary = bound_handoff_field(&h.summary); - let open_q = serde_json::to_string(&bound_handoff_list(&h.open_questions))?; - let next_s = serde_json::to_string(&bound_handoff_list(&h.next_steps))?; - let files = serde_json::to_string(&bound_handoff_list(&h.files_touched))?; - let from_session: Option<&[u8]> = h.from_session_id.as_ref().map(|s| &s.as_bytes()[..]); + // // Normalize the stored cwd: strip trailing path separators (keep a bare root // as "/"). The hook extractor preserves whatever the agent payload sent, // so this single write point guarantees a consistent stored form for both // manual and auto (SessionEnd) handoffs, keeping the next session's // path-boundary match robust to trailing slash/backslash drift. - let cwd: Option = h.cwd.as_ref().map(|p| { - let s = p.to_string_lossy(); - let trimmed = s.trim_end_matches(['/', '\\']); - if trimmed.is_empty() { - "/".to_string() - } else { - trimmed.to_string() - } - }); + Ok(HandoffFields { + summary: bound_handoff_field(&h.summary), + open_questions: serde_json::to_string(&bound_handoff_list(&h.open_questions))?, + next_steps: serde_json::to_string(&bound_handoff_list(&h.next_steps))?, + files_touched: serde_json::to_string(&bound_handoff_list(&h.files_touched))?, + cwd: h.cwd.as_ref().map(|p| { + let s = p.to_string_lossy(); + let trimmed = s.trim_end_matches(['/', '\\']); + if trimmed.is_empty() { + "/".to_string() + } else { + trimmed.to_string() + } + }), + }) +} + +/// A newer automatic handoff from the exact same cwd is the only one that +/// can ever win there, even before a SessionStart occurs. Bound abandoned +/// same-directory sessions without touching deliberate manual handoffs or +/// independent parent/sibling cwd scopes — and without crossing an operator +/// boundary: the same directory inside a shared container is the norm, so +/// owner equality is the only thing keeping one operator's SessionEnd from +/// retiring another's pending baton. `keep` spares the baton being refreshed. +fn expire_same_cwd_auto_handoffs( + conn: &Transaction<'_>, + h: &NewHandoff, + cwd: Option<&str>, + keep: Option<&HandoffId>, + now: i64, +) -> StoreResult<()> { + let expired = conn.execute( + "UPDATE handoffs SET state = 'expired' \ + WHERE workspace_id = ?1 AND project_id = ?2 \ + AND state = 'open' AND from_session_id IS NOT NULL \ + AND (cwd = ?3 OR (cwd IS NULL AND ?3 IS NULL)) \ + AND owner_user IS ?4 AND id IS NOT ?5", + params![ + h.workspace_id.as_bytes(), + h.project_id.as_bytes(), + cwd, + h.owner_user.as_deref(), + keep.map(HandoffId::as_bytes) + ], + )?; + if expired > 0 { + audit( + conn, + "expire_superseded_handoffs", + Some(h.workspace_id.as_bytes()), + Some(h.project_id.as_bytes()), + None, + None, + now, + )?; + } + Ok(()) +} + +/// Rewrite an open automatic handoff with a newer checkpoint of the same +/// session, as if it had just been inserted. +fn refresh_handoff_row(conn: &Transaction<'_>, id: &HandoffId, h: &NewHandoff) -> StoreResult<()> { + let now = Timestamp::now().as_microsecond(); + let fields = handoff_fields(h)?; + expire_same_cwd_auto_handoffs(conn, h, fields.cwd.as_deref(), Some(id), now)?; + conn.execute( + "UPDATE handoffs SET cwd = ?2, summary = ?3, open_questions = ?4, next_steps = ?5, \ + files_touched = ?6, created_at = ?7 WHERE id = ?1 AND state = 'open'", + params![ + id.as_bytes(), + fields.cwd, + fields.summary, + fields.open_questions, + fields.next_steps, + fields.files_touched, + now + ], + )?; + audit( + conn, + "refresh_handoff", + Some(h.workspace_id.as_bytes()), + Some(h.project_id.as_bytes()), + None, + None, + now, + )?; + Ok(()) +} + +fn insert_handoff_row(conn: &Transaction<'_>, h: &NewHandoff) -> StoreResult { + validate_identity_storage_key(h.owner_user.as_deref(), "handoff owner")?; + let id = HandoffId::new(); + let now = Timestamp::now().as_microsecond(); + let HandoffFields { + summary, + open_questions: open_q, + next_steps: next_s, + files_touched: files, + cwd, + } = handoff_fields(h)?; + let from_session: Option<&[u8]> = h.from_session_id.as_ref().map(|s| &s.as_bytes()[..]); let from_agent = h.from_agent.as_str(); let to_agent = h.to_agent.map(AgentKind::as_str); - // A newer automatic handoff from the exact same cwd is the only one that - // can ever win there, even before a SessionStart occurs. Bound abandoned - // same-directory sessions without touching deliberate manual handoffs or - // independent parent/sibling cwd scopes — and without crossing an operator - // boundary: the same directory inside a shared container is the norm, so - // owner equality is the only thing keeping one operator's SessionEnd from - // retiring another's pending baton. if from_session.is_some() { - let expired = conn.execute( - "UPDATE handoffs SET state = 'expired' \ - WHERE workspace_id = ?1 AND project_id = ?2 \ - AND state = 'open' AND from_session_id IS NOT NULL \ - AND (cwd = ?3 OR (cwd IS NULL AND ?3 IS NULL)) \ - AND owner_user IS ?4", - params![ - h.workspace_id.as_bytes(), - h.project_id.as_bytes(), - cwd, - h.owner_user.as_deref() - ], - )?; - if expired > 0 { - audit( - conn, - "expire_superseded_handoffs", - Some(h.workspace_id.as_bytes()), - Some(h.project_id.as_bytes()), - None, - None, - now, - )?; - } + expire_same_cwd_auto_handoffs(conn, h, cwd.as_deref(), None, now)?; } // Insert + audit atomically. Handoffs are keyed by agent/session, not a DB // user, so the audit author is NULL — the row records the lifecycle event @@ -6384,6 +6588,125 @@ pub(crate) mod tests { assert_eq!(state, "accepted", "an accepted handoff is not backlog"); } + // A turn checkpoint refreshes only its own session's unclaimed baton in + // the same scope and owner bucket (invariant #16): a claimed baton, + // another session's, a manual one, another operator's bucket and another + // project all survive. The refreshed baton keeps its id, and a checkpoint + // that arrives after the session ended touches nothing. + #[test] + fn checkpoint_session_handoff_refreshes_only_the_live_sessions_own_baton() { + let (_tmp, mut conn, ws, proj) = fresh_db(); + let other_proj = get_or_create_project(&mut conn, &ws, "other", None).unwrap(); + let session = hook_session(SessionId::new(), ws, proj, Some("user:alice")); + begin_session(&mut conn, &session).unwrap(); + let other_session = hook_session(SessionId::new(), ws, proj, Some("user:alice")); + begin_session(&mut conn, &other_session).unwrap(); + let baton = + |from: Option, project: ProjectId, owner: &str, cwd: &str| NewHandoff { + workspace_id: ws, + project_id: project, + from_session_id: from, + from_agent: AgentKind::OpenCode, + to_agent: None, + cwd: Some(cwd.into()), + summary: format!("{cwd} baton"), + open_questions: Vec::new(), + next_steps: Vec::new(), + files_touched: Vec::new(), + owner_user: Some(owner.into()), + }; + let first = checkpoint_session_handoff( + &mut conn, + &baton(Some(session.id), proj, "user:alice", "/turn/1"), + ) + .unwrap() + .expect("a live session publishes its baton"); + let claimed = insert_handoff( + &mut conn, + &baton(Some(session.id), proj, "user:alice", "/claimed"), + ) + .unwrap(); + conn.execute( + "UPDATE handoffs SET state = 'accepted' WHERE id = ?1", + params![claimed.as_bytes()], + ) + .unwrap(); + let survivors = [ + claimed, + insert_handoff( + &mut conn, + &baton(Some(other_session.id), proj, "user:alice", "/sibling"), + ) + .unwrap(), + insert_handoff(&mut conn, &baton(None, proj, "user:alice", "/manual")).unwrap(), + insert_handoff( + &mut conn, + &baton(Some(session.id), proj, "user:bob", "/bob"), + ) + .unwrap(), + insert_handoff( + &mut conn, + &baton(Some(session.id), other_proj, "user:alice", "/elsewhere"), + ) + .unwrap(), + ]; + + let second = checkpoint_session_handoff( + &mut conn, + &baton(Some(session.id), proj, "user:alice", "/turn/2"), + ) + .unwrap(); + + assert_eq!(second, Some(first), "the live baton is refreshed in place"); + let (state, summary): (String, String) = conn + .query_row( + "SELECT state, summary FROM handoffs WHERE id = ?1", + params![first.as_bytes()], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert_eq!( + (state.as_str(), summary.as_str()), + ("open", "/turn/2 baton") + ); + for survivor in survivors { + let untouched: bool = conn + .query_row( + "SELECT summary NOT LIKE '/turn/%' FROM handoffs WHERE id = ?1", + params![survivor.as_bytes()], + |r| r.get(0), + ) + .unwrap(); + assert!(untouched); + } + let rows: i64 = conn + .query_row("SELECT COUNT(*) FROM handoffs", [], |r| r.get(0)) + .unwrap(); + assert_eq!(rows, 6, "no expired row per turn"); + + let wrong_scope = checkpoint_session_handoff( + &mut conn, + &baton(Some(session.id), other_proj, "user:alice", "/turn/3"), + ); + assert!(matches!(wrong_scope, Err(StoreError::InvalidState(_)))); + + end_session(&mut conn, &session.id, None).unwrap(); + let late = checkpoint_session_handoff( + &mut conn, + &baton(Some(session.id), proj, "user:alice", "/late"), + ) + .unwrap(); + assert_eq!(late, None, "a checkpoint after the end touches nothing"); + let summary: String = conn + .query_row( + "SELECT summary FROM handoffs WHERE id = ?1", + params![first.as_bytes()], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(summary, "/turn/2 baton"); + } + fn embed_failure_rows(conn: &Connection) -> i64 { conn.query_row("SELECT COUNT(*) FROM page_embed_failures", [], |r| r.get(0)) .unwrap() @@ -11091,6 +11414,7 @@ pub(crate) mod tests { &bob_event, &OwnerFilter::User(bob.into()), Some("fresh-key"), + None, ), Err(StoreError::SessionCollision) )); @@ -11109,6 +11433,7 @@ pub(crate) mod tests { &bob_event, &OwnerFilter::User(alice.into()), Some("fresh-key"), + None, ) .unwrap(), HookSessionAdmission::Observation { @@ -11143,6 +11468,7 @@ pub(crate) mod tests { &observation, &OwnerFilter::User("user:alice".into()), None, + None, ) .unwrap(), HookSessionAdmission::Observation { @@ -11166,6 +11492,196 @@ pub(crate) mod tests { assert_eq!(session_project.as_slice(), proj.as_bytes()); } + // An explicit native move rebinds the live row only when the stored cwd + // still matches the move's source (compared through `normalize_cwd`), the + // move is keyed, and the session is open. Every other shape is recorded + // like any scope-drifted event; only owner/agent mismatch is refused. + #[test] + fn native_move_rebinds_live_session_only_from_its_current_cwd() { + for (case, stored_cwd, from, key, ended, owner, rebinds) in [ + ( + "moves", + "/repo/alpha", + "/repo/alpha", + Some("move"), + false, + "user:alice", + true, + ), + ( + "case-folded windows source", + r"C:\Repo\Alpha", + "c:/repo/alpha/", + Some("move"), + false, + "user:alice", + true, + ), + ( + "stale source", + "/repo/alpha", + "/repo/elsewhere", + Some("move"), + false, + "user:alice", + false, + ), + ( + "unix case differs", + "/Repo/alpha", + "/repo/alpha", + Some("move"), + false, + "user:alice", + false, + ), + ( + "unkeyed", + "/repo/alpha", + "/repo/alpha", + None, + false, + "user:alice", + false, + ), + ( + "ended", + "/repo/alpha", + "/repo/alpha", + Some("move"), + true, + "user:alice", + false, + ), + ] { + let (_tmp, mut conn, ws, proj) = fresh_db(); + let mut source = hook_session(SessionId::new(), ws, proj, Some("user:alice")); + source.agent_kind = AgentKind::OpenCode; + source.cwd = Some(stored_cwd.into()); + begin_session(&mut conn, &source).unwrap(); + if ended { + end_session(&mut conn, &source.id, None).unwrap(); + } + let target = get_or_create_project(&mut conn, &ws, "beta", None).unwrap(); + let mut moved = source.clone(); + moved.project_id = target; + moved.cwd = Some("/repo/beta".into()); + let mut start = hook_observation(&moved); + start.kind = ObservationKind::SessionStart; + let owner = OwnerFilter::User(owner.into()); + assert!( + matches!( + admit_hook_session_event(&mut conn, &moved, &start, &owner, key, Some(from)), + Ok(HookSessionAdmission::Observation { + ingest: IngestObservationOutcome::Inserted(_), + .. + }) + ), + "{case}: the move event itself is always recorded" + ); + let (project, cwd): (Vec, String) = conn + .query_row( + "SELECT project_id, cwd FROM sessions WHERE id = ?1", + params![source.id.as_bytes()], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .unwrap(); + let expected = if rebinds { + (target, "/repo/beta") + } else { + (proj, stored_cwd) + }; + assert_eq!( + (project.as_slice(), cwd.as_str()), + (expected.0.as_bytes().as_slice(), expected.1), + "{case}" + ); + } + } + + // Identity still wins over a move: a foreign operator's move is refused + // before anything is written, and a completed move replayed after the + // session moved back never rebinds it to the stale target. + #[test] + fn native_move_rejects_foreign_owner_and_ignores_completed_replay() { + let (_tmp, mut conn, ws, proj) = fresh_db(); + let mut source = hook_session(SessionId::new(), ws, proj, Some("user:alice")); + source.agent_kind = AgentKind::OpenCode; + source.cwd = Some("/repo/alpha".into()); + begin_session(&mut conn, &source).unwrap(); + let target = get_or_create_project(&mut conn, &ws, "beta", None).unwrap(); + let mut moved = source.clone(); + moved.project_id = target; + moved.cwd = Some("/repo/beta".into()); + let mut start = hook_observation(&moved); + start.kind = ObservationKind::SessionStart; + assert!(matches!( + admit_hook_session_event( + &mut conn, + &moved, + &start, + &OwnerFilter::User("user:bob".into()), + Some("move"), + Some("/repo/alpha") + ), + Err(StoreError::SessionCollision) + )); + let keys: i64 = conn + .query_row("SELECT COUNT(*) FROM ingest_keys", [], |row| row.get(0)) + .unwrap(); + assert_eq!(keys, 0, "a refused move must not claim its key"); + + let owner = OwnerFilter::User("user:alice".into()); + admit_hook_session_event( + &mut conn, + &moved, + &start, + &owner, + Some("move"), + Some("/repo/alpha"), + ) + .unwrap(); + complete_observation_ingest(&mut conn, &target, "move").unwrap(); + let mut back = hook_observation(&source); + back.kind = ObservationKind::SessionStart; + admit_hook_session_event( + &mut conn, + &source, + &back, + &owner, + Some("return"), + Some("/repo/beta"), + ) + .unwrap(); + complete_observation_ingest(&mut conn, &proj, "return").unwrap(); + assert!(matches!( + admit_hook_session_event( + &mut conn, + &moved, + &start, + &owner, + Some("move"), + Some("/repo/alpha") + ) + .unwrap(), + HookSessionAdmission::Observation { + ingest: IngestObservationOutcome::AlreadyComplete, + .. + } + )); + let (project, cwd): (Vec, String) = conn + .query_row( + "SELECT project_id, cwd FROM sessions WHERE id = ?1", + params![source.id.as_bytes()], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .unwrap(); + assert_eq!( + (project.as_slice(), cwd.as_str()), + (proj.as_bytes().as_slice(), "/repo/alpha") + ); + } + // Identity is still identity: a different operator or a different agent // reusing the UUID stays terminal, in the same scope-moved shape as above. #[test] @@ -11185,6 +11701,7 @@ pub(crate) mod tests { &observation, &OwnerFilter::User("user:bob".into()), None, + None, ), Err(StoreError::SessionCollision) )); @@ -11199,6 +11716,7 @@ pub(crate) mod tests { &observation, &OwnerFilter::User("user:alice".into()), None, + None, ), Err(StoreError::SessionCollision) )); @@ -11229,6 +11747,7 @@ pub(crate) mod tests { &end, &OwnerFilter::User("user:alice".into()), None, + None, ) .unwrap(), HookSessionAdmission::InvalidScopedEnd @@ -11248,6 +11767,180 @@ pub(crate) mod tests { assert_eq!(observations, 0); } + // A scope-drifted SessionEnd (a marker appeared under a running session) + // still ends that session, in the scope it was recorded in — but only + // from the session's own cwd, and never for another operator or agent + // (invariant #16). Every refused shape leaves the session open and writes + // nothing. + #[test] + fn drifted_session_end_ends_only_the_same_actor_agent_and_cwd() { + enum Expect { + Ended, + Ignored, + Collision, + } + let alice = OwnerFilter::User("user:alice".into()); + let bob = OwnerFilter::User("user:bob".into()); + for (case, stored_cwd, event_cwd, owner, agent, expect) in [ + ( + "same actor, agent and cwd", + Some("/repo/app"), + Some("/repo/app/"), + &alice, + AgentKind::OpenCode, + Expect::Ended, + ), + ( + "windows case-only cwd drift", + Some(r"C:\Repo\App"), + Some("c:/repo/app"), + &alice, + AgentKind::OpenCode, + Expect::Ended, + ), + ( + "other agent kinds too", + Some("/repo/app"), + Some("/repo/app"), + &alice, + AgentKind::Codex, + Expect::Ended, + ), + ( + "different cwd", + Some("/repo/app"), + Some("/repo/other"), + &alice, + AgentKind::OpenCode, + Expect::Ignored, + ), + ( + "unix case differs", + Some("/repo/App"), + Some("/repo/app"), + &alice, + AgentKind::OpenCode, + Expect::Ignored, + ), + ( + "event without cwd", + Some("/repo/app"), + None, + &alice, + AgentKind::OpenCode, + Expect::Ignored, + ), + ( + "session without cwd", + None, + Some("/repo/app"), + &alice, + AgentKind::OpenCode, + Expect::Ignored, + ), + ( + "other operator", + Some("/repo/app"), + Some("/repo/app"), + &bob, + AgentKind::OpenCode, + Expect::Collision, + ), + ( + "unattributed caller", + Some("/repo/app"), + Some("/repo/app"), + &OwnerFilter::Unattributed, + AgentKind::OpenCode, + Expect::Collision, + ), + ] { + let (_tmp, mut conn, ws, proj) = fresh_db(); + let mut session = hook_session(SessionId::new(), ws, proj, Some("user:alice")); + session.agent_kind = agent; + session.cwd = stored_cwd.map(Into::into); + begin_session(&mut conn, &session).unwrap(); + let marker_ws = get_or_create_workspace(&mut conn, "windows").unwrap(); + let marker_proj = get_or_create_project(&mut conn, &marker_ws, "app", None).unwrap(); + let mut drifted = session.clone(); + drifted.workspace_id = marker_ws; + drifted.project_id = marker_proj; + drifted.cwd = event_cwd.map(Into::into); + let end = session_end_observation(&drifted); + + let result = + admit_hook_session_event(&mut conn, &drifted, &end, owner, Some("end"), None); + let ended: Option = conn + .query_row( + "SELECT ended_at FROM sessions WHERE id = ?1", + params![session.id.as_bytes()], + |row| row.get(0), + ) + .unwrap(); + let landed: Vec> = conn + .prepare("SELECT project_id FROM observations") + .unwrap() + .query_map([], |row| row.get(0)) + .unwrap() + .collect::>() + .unwrap(); + match expect { + Expect::Ended => { + let Ok(HookSessionAdmission::EndOpen { + session: admitted, .. + }) = result + else { + panic!("{case}: expected the end to be admitted, got {result:?}"); + }; + assert_eq!( + (admitted.workspace_id, admitted.project_id), + (ws, proj), + "{case}" + ); + assert_eq!(landed, vec![proj.as_bytes().to_vec()], "{case}"); + end_admitted_session(&mut conn, &admitted, None).unwrap(); + let ended: Option = conn + .query_row( + "SELECT ended_at FROM sessions WHERE id = ?1", + params![session.id.as_bytes()], + |row| row.get(0), + ) + .unwrap(); + assert!(ended.is_some(), "{case}"); + } + Expect::Ignored => { + assert!( + matches!(result, Ok(HookSessionAdmission::InvalidScopedEnd)), + "{case}: {result:?}" + ); + assert!(ended.is_none() && landed.is_empty(), "{case}"); + } + Expect::Collision => { + assert!( + matches!(result, Err(StoreError::SessionCollision)), + "{case}: {result:?}" + ); + assert!(ended.is_none() && landed.is_empty(), "{case}"); + } + } + } + + // Same operator and cwd, but another agent reusing the id. + let (_tmp, mut conn, ws, proj) = fresh_db(); + let mut session = hook_session(SessionId::new(), ws, proj, Some("user:alice")); + session.agent_kind = AgentKind::OpenCode; + session.cwd = Some("/repo/app".into()); + begin_session(&mut conn, &session).unwrap(); + let mut drifted = session.clone(); + drifted.project_id = get_or_create_project(&mut conn, &ws, "marker", None).unwrap(); + drifted.agent_kind = AgentKind::ClaudeCode; + let end = session_end_observation(&drifted); + assert!(matches!( + admit_hook_session_event(&mut conn, &drifted, &end, &alice, None, None), + Err(StoreError::SessionCollision) + )); + } + #[test] fn concurrent_alice_and_bob_events_only_admit_alice_observation() { let (tmp, mut conn, ws, proj) = fresh_db(); @@ -11269,6 +11962,7 @@ pub(crate) mod tests { &hook_observation(&session), &OwnerFilter::User(owner.into()), Some(key), + None, ) }) }; @@ -11305,6 +11999,7 @@ pub(crate) mod tests { &observation, &OwnerFilter::User("user:alice".into()), Some("tuple-key"), + None, ), Err(StoreError::InvalidState(_)) )); @@ -11329,6 +12024,7 @@ pub(crate) mod tests { &observation, &OwnerFilter::User("user:alice".into()), Some("agent-key"), + None, ), Err(StoreError::SessionCollision) )); @@ -11353,6 +12049,7 @@ pub(crate) mod tests { &hook_observation(&owned), &OwnerFilter::User("user:bob".into()), Some("denied-owner"), + None, ), Err(StoreError::SessionCollision) )); @@ -11368,6 +12065,7 @@ pub(crate) mod tests { &hook_observation(&shared), &OwnerFilter::User("user:alice".into()), None, + None, ) .unwrap(); assert!(admitted_session(admission).owner().is_none()); @@ -11387,6 +12085,7 @@ pub(crate) mod tests { &hook_observation(&shared), &OwnerFilter::User("user:bob".into()), Some("shared-bob"), + None, ) .unwrap(), HookSessionAdmission::Observation { .. } @@ -11420,6 +12119,7 @@ pub(crate) mod tests { &session_end_observation(&session), &OwnerFilter::User("user:alice".into()), Some(key), + None, ) .unwrap(); match admission { @@ -11458,6 +12158,7 @@ pub(crate) mod tests { &session_end_observation(&session), &OwnerFilter::User("user:alice".into()), Some(key), + None, ) .unwrap(); match admission { @@ -11486,6 +12187,7 @@ pub(crate) mod tests { &session_end_observation(&session), &OwnerFilter::User("user:alice".into()), Some("missing-end"), + None, ) .unwrap(), HookSessionAdmission::InvalidMissingEnd @@ -11513,6 +12215,7 @@ pub(crate) mod tests { &hook_observation(&session), &OwnerFilter::User("user:alice".into()), Some("ordinary-complete"), + None, ) .unwrap(), HookSessionAdmission::Observation { @@ -11533,6 +12236,7 @@ pub(crate) mod tests { &hook_observation(&session), &OwnerFilter::User("user:alice".into()), None, + None, ) .unwrap(), ); @@ -11583,6 +12287,7 @@ pub(crate) mod tests { &lifecycle_observation, &OwnerFilter::User("user:alice".into()), None, + None, ) .unwrap(), ); diff --git a/crates/ai-memory-store/src/reader.rs b/crates/ai-memory-store/src/reader.rs index 56e45ae9..acb08772 100644 --- a/crates/ai-memory-store/src/reader.rs +++ b/crates/ai-memory-store/src/reader.rs @@ -3339,6 +3339,32 @@ impl ReaderPool { .await } + /// The recorded `(workspace, project)` of a session that has not ended. + /// A mid-session checkpoint writes its artifacts there, next to where the + /// session's eventual end writes them, whatever scope the event resolved + /// to. + /// + /// # Errors + /// Propagates any SQL or pool error. + pub async fn open_session_scope( + &self, + session_id: SessionId, + ) -> StoreResult> { + self.with_conn(move |conn| { + let row = conn + .query_row( + "SELECT workspace_id, project_id FROM sessions \ + WHERE id = ?1 AND ended_at IS NULL", + params![session_id.as_bytes()], + |row| Ok((row.get::<_, Vec>(0)?, row.get::<_, Vec>(1)?)), + ) + .optional()?; + row.map(|(ws, proj)| Ok((WorkspaceId::from_slice(&ws)?, ProjectId::from_slice(&proj)?))) + .transpose() + }) + .await + } + /// Ids of every session that touches `(workspace_id, project_id)`: a /// `sessions` row in the scope OR at least one observation stamped into /// it. The second leg catches the phantom projects that mid-session diff --git a/crates/ai-memory-store/src/writer.rs b/crates/ai-memory-store/src/writer.rs index 3993eb56..4289f3bf 100644 --- a/crates/ai-memory-store/src/writer.rs +++ b/crates/ai-memory-store/src/writer.rs @@ -164,6 +164,7 @@ pub(crate) enum WriteCmd { obs: NewObservation, owner_filter: OwnerFilter, ingest_key: Option, + moved_from_cwd: Option, reply: oneshot::Sender>, }, EndAdmittedSession { @@ -220,6 +221,10 @@ pub(crate) enum WriteCmd { handoff: NewHandoff, reply: oneshot::Sender>, }, + CheckpointSessionHandoff { + handoff: NewHandoff, + reply: oneshot::Sender>>, + }, AcceptHandoff { acceptance: HandoffAcceptance, reply: oneshot::Sender>, @@ -1021,12 +1026,15 @@ impl WriterHandle { } /// Atomically validate/admit a hook session and insert its observation. + /// `moved_from_cwd` marks an explicit native relocation; see + /// `ops::admit_hook_session_event`. pub async fn admit_hook_session_event( &self, session: NewSession, obs: Sanitized, owner_filter: OwnerFilter, ingest_key: Option, + moved_from_cwd: Option, ) -> StoreResult { let (tx, rx) = oneshot::channel(); self.send(WriteCmd::AdmitHookSessionEvent { @@ -1034,6 +1042,7 @@ impl WriterHandle { obs: obs.into_inner(), owner_filter, ingest_key, + moved_from_cwd, reply: tx, }) .await?; @@ -1211,6 +1220,21 @@ impl WriterHandle { rx.await.map_err(|_| StoreError::WriterClosed)? } + /// Publish a live session's turn-checkpoint baton; `None` when the session + /// already ended or is gone. See `ops::checkpoint_session_handoff`. + /// + /// # Errors + /// Returns [`StoreError::WriterClosed`] or propagates SQL errors. + pub async fn checkpoint_session_handoff( + &self, + handoff: NewHandoff, + ) -> StoreResult> { + let (tx, rx) = oneshot::channel(); + self.send(WriteCmd::CheckpointSessionHandoff { handoff, reply: tx }) + .await?; + rx.await.map_err(|_| StoreError::WriterClosed)? + } + /// Mark a handoff accepted by the given agent / session. /// /// Returns whether this call is the one that claimed it; `false` means the @@ -2857,6 +2881,7 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver) { obs, owner_filter, ingest_key, + moved_from_cwd, reply, } => { let result = ops::admit_hook_session_event( @@ -2865,6 +2890,7 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver) { &obs, &owner_filter, ingest_key.as_deref(), + moved_from_cwd.as_deref(), ); send_or_warn(reply, result, "admit_hook_session_event"); } @@ -2958,6 +2984,10 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver) { let result = ops::insert_handoff(&mut conn, &handoff); send_or_warn(reply, result, "insert_handoff"); } + WriteCmd::CheckpointSessionHandoff { handoff, reply } => { + let result = ops::checkpoint_session_handoff(&mut conn, &handoff); + send_or_warn(reply, result, "checkpoint_session_handoff"); + } WriteCmd::AcceptHandoff { acceptance, reply } => { let result = ops::accept_handoff(&mut conn, &acceptance); send_or_warn(reply, result, "accept_handoff"); diff --git a/docs/auto-scope.md b/docs/auto-scope.md index 0bf71852..f2a8641c 100644 --- a/docs/auto-scope.md +++ b/docs/auto-scope.md @@ -93,6 +93,7 @@ AI_MEMORY_AUTO_SCOPE__MAX_ENTRIES=8192 | Auth middleware (rung 1b trusted proxy) | username or OIDC `(issuer, subject)` pair | | Auth middleware (rung 2 DB user) | `user` ← `users.username` | | MCP request header `X-Memory-Actor-Session-Id` | `session_id` for tool calls | +| MCP request `_meta["ai.opencode/sessionID"]` (OpenCode 2) | native lifecycle `session_id` for tool calls, before transport-ID fallback | | MCP request header `Mcp-Session-Id` | fallback `session_id` for tool calls | | Anonymous / no token | empty actor → single slot | @@ -120,7 +121,17 @@ a differing `(workspace, project)` is recorded rather than rejected — see [`[routing] mid_session`](marker-file.md#mid-session-navigation-routing-mid_session) for how those events are attributed. The one exception is a terminal event: a `SessionEnd` naming a different scope than its session is not that session's -end, so it is dropped rather than ending someone else's session. +end, so it is dropped rather than ending someone else's session — unless it +comes from the session's own cwd. Then the scope drifted under the same +directory (a `.ai-memory.toml` appeared mid-session), and the end closes the +session in the scope it was recorded in. + +The OpenCode 2 adapter reports a native `session.moved` explicitly, with the +directory the session left and a stable ingest key. Its first delivery rebinds +the live session row to the new scope and cwd, provided the stored cwd still +matches the one it left, so the session's later end and checkpoints land where +it now runs; earlier observations keep their scope. A redelivered move never +rebinds again. ## Client requirements @@ -130,6 +141,14 @@ MCP client config files can only declare static URL/auth headers. Static configs cannot inject the current agent-run session id into every tool call. +OpenCode 2 (2.0.4+) sends the native session id on every tool call as +`CallToolRequest.params._meta["ai.opencode/sessionID"]`, including Code Mode +and subagent calls; `initialize` carries none. Its `Mcp-Session-Id` is shared +by every session in a directory, so the metadata key is what tells concurrent +sessions apart. It is a routing coordinate, not authentication, and takes +precedence over the transport header. With the generated OpenCode 2 lifecycle +adapter, no separate bridge is needed. + Claude Code can opt into ai-memory's session-aware stdio bridge: ```bash @@ -152,7 +171,7 @@ the shared single slot. Use `per_session` only when your client or bridge can send the same opaque session id from the hook payload on each MCP request as -`X-Memory-Actor-Session-Id` (preferred) or `Mcp-Session-Id`. Otherwise +`X-Memory-Actor-Session-Id`, native MCP metadata as above, or `Mcp-Session-Id`. Otherwise requests that carry a different MCP session id fail closed to the baked default, while requests with no usable actor identity still degrade to the legacy single slot. diff --git a/docs/install.md b/docs/install.md index 24ed2c76..c0d3d4c9 100644 --- a/docs/install.md +++ b/docs/install.md @@ -589,8 +589,8 @@ The client sanitizes (built-in patterns) and truncates the excerpt before it touches the spool or wire; the server re-scrubs with its `[sanitize]` patterns before storing. If either side is off — or the marker is malformed — the Stop stays empty. Re-running `install-hooks` without `--capture-assistant` removes -the flag (idempotent). `--capture-assistant` is Claude Code and Codex on a -native hook platform only; on any other agent or the script fallback the +the flag (idempotent). `--capture-assistant` is Claude Code, Codex and +OpenCode 2 (`--agent opencode2`) on a native hook platform only; on any other agent or the script fallback the installer refuses it rather than enabling something that cannot take effect. Assistant text is privacy-sensitive — read the `SECURITY.md` notes on what it can contain and where it flows (consolidation/reviewer prompts, and out to a cloud LLM provider if one @@ -1213,6 +1213,26 @@ loaded at startup. ### OpenCode 2 (beta) +The generated plugin targets the OpenCode 2.0.10+ event API (checked against +2.0.14). OpenCode 2 runs one long-lived service behind every CLI, so closing a +terminal is not a session end: each completed root turn instead writes a +deterministic checkpoint (no LLM call) of `sessions/.md` and refreshes the +session's automatic handoff, keeping one open baton per live session. The next +session claims the latest checkpoint. Startup context is claimed once per root +session and retained on every later model request; child sessions never claim +it or publish a baton. + +Assistant text stays opt-in, as for Claude Code and Codex: set +`capture_assistant = true` on the server and install with +`ai-memory install-hooks --agent opencode2 --capture-assistant --apply`. The +plugin hands the last completed text to the native hook, which sanitizes and +caps it before it reaches the spool or the wire; the excerpt then rides in the +next session's automatic handoff. A bare re-apply preserves the opt-in. + +For a commented `opencode.jsonc`, preview `install-mcp --client opencode2` and +merge the entry into the existing `mcp.servers` object by hand: the apply path +writes strict JSON. + The 2.0 beta installs side by side as `opencode2` and shares v1's config dir and session store, but its MCP schema and plugin API changed. Wire it with the `opencode2` client/agent names: diff --git a/docs/support-matrix.md b/docs/support-matrix.md index 62e23acd..addbee89 100644 --- a/docs/support-matrix.md +++ b/docs/support-matrix.md @@ -15,7 +15,7 @@ | Command Code | Supported | MCP config (`~/.commandcode/mcp.json`) + its four stable lifecycle-hook events (`~/.commandcode/settings.json`); native commands enforce capture exclusions and `SessionStart` injects handoffs. `Stop` is only a turn boundary, so use `ai-memory finalize-session --agent command-code` after the final turn. `ai-memory run command-code` adds exact v3 native-session resume and visible-event import; experimental unsandboxed Mods remain excluded. | | Devin CLI | Supported | MCP config + lifecycle hooks. Hooks use Devin's `PostCompaction` event, inject handoffs via `hookSpecificOutput.additionalContext`, and omit subagent events because Devin does not expose them. | | OpenCode | Supported | Remote MCP config + generated TypeScript plugin; generated plugin enforces capture exclusions. | -| OpenCode 2 (`opencode2` beta) | Supported | V2 `mcp.servers` remote MCP config + generated `Plugin.define` TypeScript plugin (`ai-memory-opencode2.ts`); shares v1's config dir, session store, and agent kind. Managed runs launch, resume, and import; ledger-delta acknowledgement is limited by the shared background service (see managed-workstreams notes). Beta plugin API and schema — expect churn. | +| OpenCode 2 (`opencode2` beta) | Supported | V2 `mcp.servers` remote MCP config + generated TypeScript plugin (`ai-memory-opencode2.ts`) for the OpenCode 2.0.10+ event API, checked against 2.0.14; shares v1's config dir, session store, and agent kind. Completed root turns checkpoint the session page and baton without ending the shared service's session; MCP calls route by the native `_meta` session id, so concurrent sessions in different projects stay isolated. `--capture-assistant` is supported (double opt-in). Managed runs launch, resume, and import; ledger-delta acknowledgement is limited by the shared background service (see managed-workstreams notes). | | Cursor | Supported | MCP config + lifecycle hooks. | | Gemini CLI | Supported | MCP config + lifecycle hooks. | | Oh My Pi / OMP | Supported | Use `--client omp` / `--agent omp` (or `oh-my-pi`) for native `.omp` MCP config + TypeScript extension; generated extension enforces capture exclusions. |