From b40149084c60f3c81c63ff668954bcc48025c640 Mon Sep 17 00:00:00 2001 From: Abhijeet Prasad Date: Mon, 14 Sep 2026 18:27:39 -0400 Subject: [PATCH] fix(opencode): move delivery flushing into the daemon OpenCode waited for backend delivery at lifecycle boundaries. A slow flush blocked idle handling and shutdown for up to 10 seconds. Before: OpenCode -> journal ack -> wait for backend flush -> next event / shutdown After: OpenCode -> journal ack -> next event / shutdown | +-> daemon flushes to Braintrust in the background Keep ordered capture acknowledgement, but let the daemon flush idle, deletion, error, and disposal events. Delivery continues after OpenCode disconnects. Requires both the updated OpenCode plugin and bt daemon. --- bt-daemon/src/lib.rs | 36 +++++++- bt-daemon/tests/pipeline.rs | 54 +++++++++++ .../content/src/tracing/daemon.test.ts | 92 +++++++++++++++++++ .../opencode/content/src/tracing/daemon.ts | 7 +- 4 files changed, 182 insertions(+), 7 deletions(-) create mode 100644 src/plugins/opencode/content/src/tracing/daemon.test.ts 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; }