diff --git a/bt-daemon/src/lib.rs b/bt-daemon/src/lib.rs index 8ae2ef0..d3124a1 100644 --- a/bt-daemon/src/lib.rs +++ b/bt-daemon/src/lib.rs @@ -349,8 +349,8 @@ pub(crate) fn should_flush_ingress_event(env: &wire::Envelope) -> bool { Some(wire::FlushMode::FlushOnTurnEnd) ); should_flush_hook_event(&env.event, flush_on_turn_end) - || (env.source == "pi" - && match env.event.as_str() { + || match env.source.as_str() { + "pi" => match env.event.as_str() { // Preserve Pi's explicit lifecycle flushes even in batched mode. "session_shutdown" | "session_compact" | "session_tree" => true, // OMP can end an attempt while keeping the same turn open for retry. @@ -360,7 +360,13 @@ pub(crate) fn should_flush_ingress_event(env: &wire::Envelope) -> bool { != Some(true) } _ => false, - }) + }, + "opencode" => matches!( + env.event.as_str(), + "session.idle" | "session.deleted" | "session.error" | "server.instance.disposed" + ), + _ => false, + } } /// Capture one hook event from `stdin` and forward it to the daemon. @@ -1514,6 +1520,30 @@ mod tests { } } + #[test] + fn opencode_lifecycle_events_flush_from_ingress() { + let mut env: wire::Envelope = serde_json::from_value(serde_json::json!({ + "source": "opencode", "session_id": "opencode-session", "event": "session.idle", + "ts_ms": 1, "payload": {}, + "route": {"flush_mode": "fire_and_forget"} + })) + .unwrap(); + for event in [ + "session.idle", + "session.deleted", + "session.error", + "server.instance.disposed", + ] { + env.event = event.into(); + assert!(should_flush_ingress_event(&env)); + } + env.event = "session.compacted".into(); + assert!(!should_flush_ingress_event(&env)); + env.event = "session.idle".into(); + env.source = "pi".into(); + assert!(!should_flush_ingress_event(&env)); + } + #[test] fn additional_metadata_overrides_a_route_only_with_a_json_object() { let mut route = SessionRoute { diff --git a/bt-daemon/tests/pipeline.rs b/bt-daemon/tests/pipeline.rs index 5cd08e2..7be859c 100644 --- a/bt-daemon/tests/pipeline.rs +++ b/bt-daemon/tests/pipeline.rs @@ -1304,6 +1304,60 @@ async fn pi_lifecycle_flushes_without_an_explicit_client_flush() { handle.await.unwrap(); } +#[tokio::test] +async fn opencode_lifecycle_flushes_without_an_explicit_client_flush() { + let (socket, handle, flushes, _tmp) = start_tracking_daemon("test").await; + let host = dummy_host(); + for event in [ + "session.idle", + "session.deleted", + "session.error", + "server.instance.disposed", + ] { + let session_id = format!("opencode-background-{event}"); + let native_session_id = format!("native-{event}"); + let mut start = envelope(&session_id, "session.created", 1); + start.source = "opencode".into(); + start.payload = serde_json::json!({ + "properties": {"info": {"id": native_session_id}} + }); + forward_envelope(&start, &socket, &host, false) + .await + .unwrap(); + + let mut terminal = envelope(&session_id, event, 2); + terminal.source = "opencode".into(); + terminal.payload = serde_json::json!({ + "properties": {"sessionID": native_session_id} + }); + forward_envelope(&terminal, &socket, &host, false) + .await + .unwrap(); + + // The forwarding client has disconnected; backend delivery belongs to + // the daemon and must continue independently. + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if flushes + .lock() + .unwrap() + .get(&session_id) + .copied() + .unwrap_or_default() + > 0 + { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap_or_else(|_| panic!("OpenCode {event} was not flushed by the daemon")); + } + shutdown(&socket).await; + handle.await.unwrap(); +} + #[tokio::test] async fn hook_capture_stops_at_the_durable_journal_boundary() { let (socket, handle, tmp) = start_slow_daemon().await; diff --git a/src/plugins/opencode/content/src/tracing/daemon.test.ts b/src/plugins/opencode/content/src/tracing/daemon.test.ts new file mode 100644 index 0000000..43eac52 --- /dev/null +++ b/src/plugins/opencode/content/src/tracing/daemon.test.ts @@ -0,0 +1,92 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { PluginInput } from "@opencode-ai/plugin"; +import type { Event } from "@opencode-ai/sdk"; + +const mockState = vi.hoisted(() => ({ + logs: [] as Array>, + flushes: [] as string[], + closed: 0, + logGate: undefined as Promise | undefined, +})); + +vi.mock("../runtime/daemon-client", () => ({ + DaemonClient: class { + async log(envelope: Record): Promise { + mockState.logs.push(envelope); + await mockState.logGate; + return true; + } + async flush(sessionId: string): Promise { + mockState.flushes.push(sessionId); + return true; + } + async close(): Promise { + mockState.closed += 1; + } + }, +})); + +import { createDaemonTracingHooks } from "./daemon"; + +describe("OpenCode daemon adapter", () => { + beforeEach(() => { + mockState.logs.length = 0; + mockState.flushes.length = 0; + mockState.closed = 0; + mockState.logGate = undefined; + }); + + it("awaits journal acknowledgement but leaves lifecycle delivery flushing to the daemon", async () => { + const input = { + directory: "/tmp/project", + worktree: "/tmp/project", + } as PluginInput; + const hooks = createDaemonTracingHooks( + input, + { + projectName: "agents", + route: { + destination: { type: "project_logs", project_name: "agents" }, + flush_mode: "fire_and_forget", + }, + }, + () => {}, + ); + const handleEvent = hooks.event as (input: { event: Event }) => Promise; + + let acknowledge!: () => void; + mockState.logGate = new Promise((resolve) => { + acknowledge = resolve; + }); + let idleForwarded = false; + const idle = handleEvent({ + event: { type: "session.idle", properties: { sessionID: "native-session" } } as Event, + }).then(() => { + idleForwarded = true; + }); + await Promise.resolve(); + expect(idleForwarded).toBe(false); + acknowledge(); + await idle; + + mockState.logGate = undefined; + for (const type of ["session.deleted", "session.error"] as const) { + await handleEvent({ + event: { type, properties: { sessionID: "native-session" } } as Event, + }); + } + await handleEvent({ + event: { type: "server.instance.disposed", properties: {} } as Event, + }); + + expect(mockState.logs.map((log) => log.event)).toEqual([ + "session.idle", + "session.deleted", + "session.error", + "server.instance.disposed", + ]); + expect(mockState.logs.every((log) => log.source === "opencode")).toBe(true); + expect(mockState.flushes).toHaveLength(0); + expect(mockState.closed).toBe(1); + }); +}); diff --git a/src/plugins/opencode/content/src/tracing/daemon.ts b/src/plugins/opencode/content/src/tracing/daemon.ts index ec51029..a5e3792 100644 --- a/src/plugins/opencode/content/src/tracing/daemon.ts +++ b/src/plugins/opencode/content/src/tracing/daemon.ts @@ -55,15 +55,14 @@ export function createDaemonTracingHooks( }, route: config.route, }); - if (event === "session.idle" || event === "session.deleted" || event === "session.error") { - await daemon.flush(daemonSessionId); - } }; return { event: async ({ event }: { event: Event }) => { if (event.type === "server.instance.disposed") { - await daemon.flush(daemonSessionId); + // Give the daemon a durable terminal event before disconnecting. It owns + // the potentially slow backend flush, so OpenCode shutdown is not blocked. + await forward(event.type, { properties: event.properties }); await daemon.close(); return; }