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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 33 additions & 3 deletions bt-daemon/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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"
),
Comment thread
AbhiPrasad marked this conversation as resolved.
_ => false,
}
}

/// Capture one hook event from `stdin` and forward it to the daemon.
Expand Down Expand Up @@ -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 {
Expand Down
54 changes: 54 additions & 0 deletions bt-daemon/tests/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
92 changes: 92 additions & 0 deletions src/plugins/opencode/content/src/tracing/daemon.test.ts
Original file line number Diff line number Diff line change
@@ -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<Record<string, unknown>>,
flushes: [] as string[],
closed: 0,
logGate: undefined as Promise<void> | undefined,
}));

vi.mock("../runtime/daemon-client", () => ({
DaemonClient: class {
async log(envelope: Record<string, unknown>): Promise<boolean> {
mockState.logs.push(envelope);
await mockState.logGate;
return true;
}
async flush(sessionId: string): Promise<boolean> {
mockState.flushes.push(sessionId);
return true;
}
async close(): Promise<void> {
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<void>;

let acknowledge!: () => void;
mockState.logGate = new Promise<void>((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);
});
});
7 changes: 3 additions & 4 deletions src/plugins/opencode/content/src/tracing/daemon.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
Loading