diff --git a/bt-daemon/src/dispatch.rs b/bt-daemon/src/dispatch.rs index a0cb73a..9ecdc3c 100644 --- a/bt-daemon/src/dispatch.rs +++ b/bt-daemon/src/dispatch.rs @@ -34,6 +34,7 @@ enum SessionMsg { Event(Box, oneshot::Sender<()>), Configure(Box, oneshot::Sender<()>), Flush(oneshot::Sender), + Finalize(oneshot::Sender), Shutdown(oneshot::Sender<()>), } @@ -157,6 +158,20 @@ impl Session { } } + /// Finalize an invocation-local session and flush its sink. Unlike a + /// delivery checkpoint, a managed-run completion is a terminal boundary. + pub async fn finalize(&self, timeout: std::time::Duration) -> (bool, u64) { + self.touch(); + let (reply_tx, reply_rx) = oneshot::channel(); + if self.tx.send(SessionMsg::Finalize(reply_tx)).await.is_err() { + return (false, self.counters.queued.load(Ordering::Relaxed)); + } + match tokio::time::timeout(timeout, reply_rx).await { + Ok(Ok(pending)) => (pending == 0, pending), + _ => (false, self.counters.queued.load(Ordering::Relaxed)), + } + } + /// Reconfigure the sink before a refresh-triggered flush. Queue ordering /// guarantees that all earlier events are processed first. pub async fn configure(&self, config: crate::wire::SessionConfig) -> anyhow::Result<()> { @@ -249,7 +264,9 @@ struct SessionActor { enum BatchMode { Live, Replay, - Flush, + Checkpoint, + TerminalFinalize, + ShutdownFinalize, } impl BatchMode { @@ -257,12 +274,19 @@ impl BatchMode { match self { Self::Live => ("translate failed", "sink emit failed"), Self::Replay => ("journal replay failed", "sink replay emit failed"), - Self::Flush => ("translate flush failed", "sink emit (flush) failed"), + Self::Checkpoint => ( + "translate checkpoint failed", + "sink emit (checkpoint) failed", + ), + Self::TerminalFinalize | Self::ShutdownFinalize => ( + "translate finalization failed", + "sink emit (finalization) failed", + ), } } fn observes_correlation(self) -> bool { - !matches!(self, Self::Flush) + !matches!(self, Self::ShutdownFinalize) } } @@ -291,6 +315,9 @@ impl SessionActor { SessionMsg::Flush(r) => { let _ = r.send(0); } + SessionMsg::Finalize(r) => { + let _ = r.send(0); + } SessionMsg::Shutdown(r) => { let _ = r.send(()); break; @@ -357,11 +384,28 @@ impl SessionActor { let _ = reply.send(()); } SessionMsg::Flush(reply) => { - self.drain_flush(&mut translator, &mut sink, &ctx).await; + self.checkpoint_and_flush(&mut translator, &mut sink, &ctx) + .await; + let _ = reply.send(self.counters.queued.load(Ordering::Relaxed)); + } + SessionMsg::Finalize(reply) => { + self.finalize_and_flush( + &mut translator, + &mut sink, + &ctx, + BatchMode::TerminalFinalize, + ) + .await; let _ = reply.send(self.counters.queued.load(Ordering::Relaxed)); } SessionMsg::Shutdown(reply) => { - self.drain_flush(&mut translator, &mut sink, &ctx).await; + self.finalize_and_flush( + &mut translator, + &mut sink, + &ctx, + BatchMode::ShutdownFinalize, + ) + .await; let _ = reply.send(()); break; } @@ -464,16 +508,53 @@ impl SessionActor { } } - async fn drain_flush( + async fn checkpoint_and_flush( &self, translator: &mut Box, sink: &mut Box, ctx: &SessionCtx, ) { - let translated = translator.flush(ctx); - let _ = self - .emit_translator_batches(translator, sink, ctx, translated, BatchMode::Flush) + let translated = translator.checkpoint(ctx); + let correlation_changed = self + .emit_translator_batches(translator, sink, ctx, translated, BatchMode::Checkpoint) + .await; + self.persist_correlation_if_changed(correlation_changed) .await; + self.flush_sink(sink).await; + } + + async fn finalize_and_flush( + &self, + translator: &mut Box, + sink: &mut Box, + ctx: &SessionCtx, + mode: BatchMode, + ) { + let translated = translator.finalize(ctx); + let correlation_changed = self + .emit_translator_batches(translator, sink, ctx, translated, mode) + .await; + self.persist_correlation_if_changed(correlation_changed) + .await; + self.flush_sink(sink).await; + } + + async fn persist_correlation_if_changed(&self, changed: bool) { + if !changed { + return; + } + if let Err(error) = crate::server::persist_active_parent_snapshot( + &self.data_dir, + &self.correlation_key, + &self.correlation, + ) + .await + { + self.set_error(error); + } + } + + async fn flush_sink(&self, sink: &mut Box) { if let Err(e) = sink.flush().await { self.set_error(format!("sink flush failed: {e}")); } diff --git a/bt-daemon/src/server.rs b/bt-daemon/src/server.rs index 0359648..2f301a1 100644 --- a/bt-daemon/src/server.rs +++ b/bt-daemon/src/server.rs @@ -478,7 +478,7 @@ impl Daemon { let session = { self.sessions.lock().unwrap().get(&key).cloned() }; if let Some(session) = session { let (flushed, pending) = session - .flush(Duration::from_millis(params.timeout_ms)) + .finalize(Duration::from_millis(params.timeout_ms)) .await; result.flushed &= flushed; result.pending = result.pending.saturating_add(pending); diff --git a/bt-daemon/src/translate/antigravity.rs b/bt-daemon/src/translate/antigravity.rs index c5ec613..86752b5 100644 --- a/bt-daemon/src/translate/antigravity.rs +++ b/bt-daemon/src/translate/antigravity.rs @@ -516,7 +516,7 @@ impl AgentTranslator for AntigravityTranslator { Ok(ops) } - fn flush(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { + fn finalize(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { let mut ops = Vec::new(); self.close_pending(self.last_ts_ms, None, &mut ops); self.close_turn(self.last_ts_ms, None, &mut ops); diff --git a/bt-daemon/src/translate/claude.rs b/bt-daemon/src/translate/claude.rs index 9e5bf63..6d6d973 100644 --- a/bt-daemon/src/translate/claude.rs +++ b/bt-daemon/src/translate/claude.rs @@ -97,6 +97,7 @@ struct ClaudeTranslator { git: Arc, current_cwd: Option, last_turn_cwd: Option, + last_ts_ms: i64, } impl ClaudeTranslator { @@ -126,6 +127,7 @@ impl ClaudeTranslator { git, current_cwd: None, last_turn_cwd: None, + last_ts_ms: 0, } } @@ -700,6 +702,7 @@ impl ClaudeTranslator { impl AgentTranslator for ClaudeTranslator { fn handle(&mut self, event: &Envelope, ctx: &SessionCtx) -> anyhow::Result> { + self.last_ts_ms = self.last_ts_ms.max(event.ts_ms); anyhow::ensure!( self.pending_emission.is_none(), "Claude translator has pending catch-up work; drain it before handling another event" @@ -778,8 +781,47 @@ impl AgentTranslator for ClaudeTranslator { Ok(Some(ops)) } - fn flush(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { - Ok(Vec::new()) + fn finalize(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { + let end_ms = self.last_ts_ms; + let mut ops = Vec::new(); + for (_, tool) in self.pending_tools.drain() { + ops.push(SpanOp::Merge(SpanRow { + span_id: tool.span_id, + root_span_id: self.root_span_id.clone(), + end_ms: Some(end_ms), + error: Some("Session ended before tool completion".into()), + ..Default::default() + })); + } + if let Some(turn) = self.turn.take() { + ops.push(SpanOp::Merge(SpanRow { + span_id: turn.id, + root_span_id: self.root_span_id.clone(), + end_ms: Some(end_ms), + error: Some("Session ended before turn completion".into()), + ..Default::default() + })); + } + for (_, subagent) in self.subagents.drain() { + ops.push(SpanOp::Merge(SpanRow { + span_id: subagent.span_id, + root_span_id: self.root_span_id.clone(), + end_ms: Some(end_ms), + error: Some("Session ended before subagent completion".into()), + ..Default::default() + })); + } + if self.root_open && !self.root_ended { + self.root_ended = true; + ops.push(SpanOp::Merge(SpanRow { + span_id: self.session_span_id.clone(), + root_span_id: self.root_span_id.clone(), + end_ms: Some(end_ms), + ..Default::default() + })); + } + self.release_terminal_state(); + Ok(ops) } } diff --git a/bt-daemon/src/translate/codex.rs b/bt-daemon/src/translate/codex.rs index 1584770..44671a9 100644 --- a/bt-daemon/src/translate/codex.rs +++ b/bt-daemon/src/translate/codex.rs @@ -145,9 +145,10 @@ enum PendingWork { through_ms: Option, after: DeferredHook, }, - Flush { + CatchUp { paths: Vec, next_path: usize, + finalize: bool, }, } @@ -259,24 +260,30 @@ impl AgentTranslator for CodexTranslator { }); } } - PendingWork::Flush { + PendingWork::CatchUp { paths, mut next_path, + finalize, } => { while next_path < paths.len() { let path = &paths[next_path]; if self.catch_up_chunk(path, 0, None, &mut ops) { - if let Some(mut scope) = self.scopes.remove(path) { - self.close_dangling(&mut scope, None, &mut ops); - self.scopes.insert(path.clone(), scope); + if finalize { + if let Some(mut scope) = self.scopes.remove(path) { + self.close_dangling(&mut scope, None, &mut ops); + self.scopes.insert(path.clone(), scope); + } } next_path += 1; } // Return after any completed scope or a bounded partial read. - // This keeps a flush over many scopes bounded as well. + // This keeps catch-up over many scopes bounded as well. if !ops.is_empty() || next_path < paths.len() { - self.pending = (next_path < paths.len()) - .then_some(PendingWork::Flush { paths, next_path }); + self.pending = (next_path < paths.len()).then_some(PendingWork::CatchUp { + paths, + next_path, + finalize, + }); break; } } @@ -285,25 +292,35 @@ impl AgentTranslator for CodexTranslator { Ok(Some(ops)) } - fn flush(&mut self, ctx: &SessionCtx) -> anyhow::Result> { + fn checkpoint(&mut self, ctx: &SessionCtx) -> anyhow::Result> { + self.start_catch_up(ctx, false) + } + + fn finalize(&mut self, ctx: &SessionCtx) -> anyhow::Result> { + self.start_catch_up(ctx, true) + } +} + +impl CodexTranslator { + fn start_catch_up(&mut self, ctx: &SessionCtx, finalize: bool) -> anyhow::Result> { anyhow::ensure!( self.pending.is_none(), - "Codex translator has pending catch-up work; drain it before flushing" + "Codex translator has pending catch-up work; drain it before checkpointing" ); - // Re-read each scope to catch a late task_complete, then close dangling. + // Re-read each scope to catch a late task_complete. Finalization also + // closes dangling work whose terminal native event never arrived. let paths: Vec = self.scopes.keys().cloned().collect(); if paths.is_empty() { return Ok(Vec::new()); } - self.pending = Some(PendingWork::Flush { + self.pending = Some(PendingWork::CatchUp { paths, next_path: 0, + finalize, }); Ok(self.drain_pending(ctx)?.unwrap_or_default()) } -} -impl CodexTranslator { fn ensure_main_scope(&mut self, path: &str) { if self.scopes.contains_key(path) { return; diff --git a/bt-daemon/src/translate/debug.rs b/bt-daemon/src/translate/debug.rs index 2bbaf5c..199d428 100644 --- a/bt-daemon/src/translate/debug.rs +++ b/bt-daemon/src/translate/debug.rs @@ -74,7 +74,7 @@ impl AgentTranslator for DebugTranslator { Ok(ops) } - fn flush(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { + fn finalize(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { Ok(Vec::new()) } } diff --git a/bt-daemon/src/translate/mod.rs b/bt-daemon/src/translate/mod.rs index d97f496..620f5f2 100644 --- a/bt-daemon/src/translate/mod.rs +++ b/bt-daemon/src/translate/mod.rs @@ -92,18 +92,34 @@ pub trait AgentTranslator: Send { /// Handle one event, returning span ops to emit. fn handle(&mut self, event: &Envelope, ctx: &SessionCtx) -> anyhow::Result>; - /// Continue bounded work started by [`Self::handle`] or [`Self::flush`]. + /// Continue bounded work started by [`Self::handle`], [`Self::checkpoint`], + /// or [`Self::finalize`]. /// `Some` means the caller must emit this batch and call again; `None` /// means the translator is fully caught up. fn drain_pending(&mut self, _ctx: &SessionCtx) -> anyhow::Result>> { Ok(None) } - /// Emit any pending spans (e.g. close dangling turns) at flush/shutdown. - fn flush(&mut self, ctx: &SessionCtx) -> anyhow::Result> { + /// Catch up externally buffered observations without ending the logical + /// agent session. Delivery barriers call this before flushing the sink. + fn checkpoint(&mut self, ctx: &SessionCtx) -> anyhow::Result> { + let _ = ctx; + Ok(Vec::new()) + } + + /// Finish the logical agent session and defensively close work whose + /// terminal native event never arrived. Called only when the actor itself + /// is shutting down or being retired. + fn finalize(&mut self, ctx: &SessionCtx) -> anyhow::Result> { let _ = ctx; Ok(Vec::new()) } + + /// Backward-compatible terminal flush used by transcript import callers. + /// Live delivery barriers use [`Self::checkpoint`] instead. + fn flush(&mut self, ctx: &SessionCtx) -> anyhow::Result> { + self.finalize(ctx) + } } /// Builds translator instances for a given `source`. diff --git a/bt-daemon/src/translate/opencode.rs b/bt-daemon/src/translate/opencode.rs index e8c9ad5..9532abe 100644 --- a/bt-daemon/src/translate/opencode.rs +++ b/bt-daemon/src/translate/opencode.rs @@ -50,6 +50,7 @@ struct NativeSession { reasoning_parts: HashMap, tool_calls: HashMap>, tool_starts: HashMap, + tool_names: HashMap, tool_args: HashMap, tool_outputs: HashMap, tool_errors: HashMap, @@ -99,7 +100,7 @@ impl AgentTranslator for OpenCodeTranslator { Ok(ops) } - fn flush(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { + fn finalize(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { let now = self.last_ts_ms; let ids: Vec = self.sessions.keys().cloned().collect(); let mut ops = Vec::new(); @@ -465,6 +466,7 @@ impl OpenCodeTranslator { .or_else(|| event.payload.get("tool")) .and_then(Value::as_str) .unwrap_or("tool"); + s.tool_names.insert(call.into(), tool.into()); return vec![SpanOp::Insert(SpanRow { span_id: ids::span_id(&self.daemon_session_id, &format!("tool:{sid}:{call}")), root_span_id: s.effective_root_span_id.clone(), @@ -500,6 +502,7 @@ impl OpenCodeTranslator { return vec![]; }; if s.denied_tools.remove(call) { + s.tool_names.remove(call); return vec![]; } let Some(turn) = s.current_turn_span_id.clone() else { @@ -552,6 +555,7 @@ impl OpenCodeTranslator { ..Default::default() }; s.tool_starts.remove(call); + s.tool_names.remove(call); vec![if had_start { SpanOp::Merge(row) } else { @@ -582,6 +586,7 @@ impl OpenCodeTranslator { return vec![]; }; s.denied_tools.insert(call.into()); + s.tool_names.remove(call); let tool = props.get("tool").and_then(Value::as_str).unwrap_or("tool"); return vec![SpanOp::Insert(SpanRow { span_id: ids::span_id(&self.daemon_session_id, &format!("tool:{sid}:{call}")), @@ -624,6 +629,24 @@ impl OpenCodeTranslator { return vec![]; }; let mut ops = vec![]; + for (call, _start_ms) in s.tool_starts.drain() { + if s.denied_tools.remove(&call) { + continue; + } + let tool_name = s.tool_names.remove(&call).unwrap_or_else(|| "tool".into()); + ops.push(SpanOp::Merge(SpanRow { + span_id: ids::span_id(&self.daemon_session_id, &format!("tool:{sid}:{call}")), + root_span_id: s.effective_root_span_id.clone(), + end_ms: Some(ts), + metadata: Some(json!({ + "tool_name": tool_name, + "call_id": call, + "tool_outcome": "error", + })), + error: Some("Interrupted before tool completion".into()), + ..Default::default() + })); + } if let Some(turn) = s.current_turn_span_id.take() { ops.push(SpanOp::Merge(SpanRow { span_id: turn, diff --git a/bt-daemon/src/translate/pi.rs b/bt-daemon/src/translate/pi.rs index 8360e39..99473ef 100644 --- a/bt-daemon/src/translate/pi.rs +++ b/bt-daemon/src/translate/pi.rs @@ -137,8 +137,10 @@ impl AgentTranslator for PiTranslator { Ok(ops) } - fn flush(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { - let mut ops = self.close_turn(self.last_ts, Some("Interrupted before completion".into())); + fn finalize(&mut self, _ctx: &SessionCtx) -> anyhow::Result> { + let error = "Interrupted before completion"; + let mut ops = self.close_dangling(self.last_ts, error); + ops.extend(self.close_turn(self.last_ts, Some(error.into()))); if self.opened { ops.push(self.close_root(self.last_ts)); } @@ -147,6 +149,56 @@ impl AgentTranslator for PiTranslator { } impl PiTranslator { + fn close_dangling(&mut self, ts: i64, error: &str) -> Vec { + let Some((turn, _)) = &self.turn else { + self.pending_llms.clear(); + self.tools.clear(); + return Vec::new(); + }; + let mut ops = Vec::new(); + for pending in self.pending_llms.drain(..) { + self.llm_seq += 1; + ops.push(SpanOp::Insert(SpanRow { + span_id: ids::span_id( + &self.session_id, + &format!("llm:{}:{}", self.turn_seq, self.llm_seq), + ), + root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![turn.clone()], + name: "llm".into(), + span_type: SpanType::Llm, + start_ms: Some(pending.start_ms), + end_ms: Some(ts), + input: Some(pending.input), + metadata: pending.provider, + error: Some(error.into()), + ..Default::default() + })); + } + for (call, tool) in self.tools.drain() { + self.total_tools += 1; + ops.push(SpanOp::Insert(SpanRow { + span_id: ids::span_id(&self.session_id, &format!("tool:{}:{call}", self.turn_seq)), + root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![turn.clone()], + name: tool.name.clone(), + span_type: SpanType::Tool, + start_ms: Some(tool.start_ms), + end_ms: Some(ts), + input: Some(tool.args), + metadata: Some(json!({ + "tool_name": tool.name, + "tool_call_id": call, + "tool_approval": "approved", + "tool_outcome": "error", + })), + error: Some(error.into()), + ..Default::default() + })); + } + ops + } + fn ensure_root(&mut self, envelope: &Envelope, ctx: &SessionCtx) -> Vec { if self.opened { return Vec::new(); diff --git a/bt-daemon/tests/opencode_translator.rs b/bt-daemon/tests/opencode_translator.rs index 4960075..210e718 100644 --- a/bt-daemon/tests/opencode_translator.rs +++ b/bt-daemon/tests/opencode_translator.rs @@ -186,3 +186,68 @@ fn opencode_additional_metadata_reaches_roots_without_overriding_session_fields( assert_eq!(root.metadata.as_ref().unwrap()["team"], "platform"); assert_eq!(root.metadata.as_ref().unwrap()["source"], "opencode"); } + +#[test] +fn opencode_checkpoint_preserves_session_state_for_later_turns() { + let registry = Registry::default_agents(); + let mut translator = registry.create("opencode", "root-session"); + let ctx = SessionCtx { + session_id: "root-session".into(), + config: None, + }; + translator + .handle( + &event( + "session.created", + 1, + json!({"properties":{"info":{"id":"native"}}}), + ), + &ctx, + ) + .unwrap(); + translator + .handle( + &event("chat.message", 2, json!({"input":{"sessionID":"native"}})), + &ctx, + ) + .unwrap(); + + assert!(translator.checkpoint(&ctx).unwrap().is_empty()); + let ops = translator + .handle( + &event("chat.message", 3, json!({"input":{"sessionID":"native"}})), + &ctx, + ) + .unwrap(); + assert!(ops + .iter() + .any(|op| matches!(op, SpanOp::Insert(row) if row.name == "Turn 2"))); + assert!(ops + .iter() + .all(|op| !matches!(op, SpanOp::Insert(row) if row.name == "OpenCode"))); +} + +#[test] +fn opencode_finalization_closes_a_missing_tool_completion() { + let registry = Registry::default_agents(); + let mut translator = registry.create("opencode", "root-session"); + let ctx = SessionCtx { + session_id: "root-session".into(), + config: None, + }; + for envelope in [ + event("chat.message", 1, json!({"input":{"sessionID":"native"}})), + event( + "tool.execute.before", + 2, + json!({"input":{"sessionID":"native","callID":"call","tool":"read"},"output":{"args":{"path":"x"}}}), + ), + ] { + translator.handle(&envelope, &ctx).unwrap(); + } + let ops = translator.finalize(&ctx).unwrap(); + assert!(ops.iter().any(|op| matches!( + op, + SpanOp::Merge(row) if row.end_ms == Some(2) && row.error.is_some() + ))); +} diff --git a/bt-daemon/tests/pi_translator.rs b/bt-daemon/tests/pi_translator.rs index 46705e4..fa4d5a6 100644 --- a/bt-daemon/tests/pi_translator.rs +++ b/bt-daemon/tests/pi_translator.rs @@ -166,3 +166,68 @@ fn pi_additional_metadata_reaches_roots_without_overriding_session_fields() { assert!(root.parent_span_ids.is_empty()); assert_eq!(root.root_span_id, root.span_id); } + +#[test] +fn pi_checkpoint_preserves_the_open_session_and_turn() { + let registry = Registry::default_agents(); + let mut translator = registry.create("pi", "pi-session"); + let ctx = SessionCtx { + session_id: "pi-session".into(), + config: None, + }; + translator + .handle(&event("session_start", 1, json!({"reason":"new"})), &ctx) + .unwrap(); + translator + .handle( + &event("before_agent_start", 2, json!({"prompt":"first"})), + &ctx, + ) + .unwrap(); + + assert!(translator.checkpoint(&ctx).unwrap().is_empty()); + let ops = translator + .handle(&event("agent_end", 3, json!({"messages":[]})), &ctx) + .unwrap(); + assert!(ops.iter().any( + |op| matches!(op, SpanOp::Merge(row) if row.name.is_empty() && row.end_ms == Some(3)) + )); + assert!(ops.iter().all(|op| !matches!( + op, + SpanOp::Merge(row) + if row + .metadata + .as_ref() + .is_some_and(|metadata| metadata.get("total_turns").is_some()) + ))); +} + +#[test] +fn pi_finalization_closes_missing_llm_and_tool_events() { + let registry = Registry::default_agents(); + let mut translator = registry.create("pi", "pi-session"); + let ctx = SessionCtx { + session_id: "pi-session".into(), + config: None, + }; + for envelope in [ + event("before_agent_start", 1, json!({"prompt":"work"})), + event("context", 2, json!({"messages":[]})), + event( + "tool_execution_start", + 3, + json!({"toolCallId":"call","toolName":"read","args":{"path":"x"}}), + ), + ] { + translator.handle(&envelope, &ctx).unwrap(); + } + let ops = translator.finalize(&ctx).unwrap(); + assert!(ops.iter().any(|op| matches!( + op, + SpanOp::Insert(row) if row.span_type == SpanType::Llm && row.error.is_some() + ))); + assert!(ops.iter().any(|op| matches!( + op, + SpanOp::Insert(row) if row.span_type == SpanType::Tool && row.error.is_some() + ))); +} diff --git a/bt-daemon/tests/pipeline.rs b/bt-daemon/tests/pipeline.rs index 16a1063..1c178f9 100644 --- a/bt-daemon/tests/pipeline.rs +++ b/bt-daemon/tests/pipeline.rs @@ -70,6 +70,7 @@ struct RouteSinkRecord { destination: Mutex>, emitted: std::sync::atomic::AtomicU64, flushes: std::sync::atomic::AtomicU64, + ops: Mutex>, } #[derive(Default)] @@ -102,6 +103,7 @@ impl Sink for RouteRecordingSink { } async fn emit(&mut self, ops: &[SpanOp]) -> anyhow::Result { + self.record.ops.lock().unwrap().extend_from_slice(ops); self.record .emitted .fetch_add(ops.len() as u64, std::sync::atomic::Ordering::Relaxed); @@ -1291,6 +1293,50 @@ async fn managed_run_flush_is_scoped_to_its_accepted_sessions() { handle.await.unwrap(); } +#[tokio::test] +async fn managed_run_completion_finalizes_an_open_agent_root() { + let provider = Arc::new(TestAuthProvider { + calls: Mutex::new(Vec::new()), + fail: false, + first_lease_expired: false, + }); + let recording = Arc::new(RouteRecordingSinkFactory::default()); + let (_data_dir, socket, handle, _tmp) = start_daemon_with(provider, recording.clone()).await; + let mut env = envelope("managed-opencode", "session.created", 7); + env.source = "opencode".into(); + env.managed_run_id = Some("managed-run".into()); + env.payload = serde_json::json!({"properties":{"info":{"id":"native"}}}); + forward_envelope(&env, &socket, &dummy_host(), false) + .await + .unwrap(); + + let checkpoint = flush_session("managed-opencode", &socket, 5_000) + .await + .unwrap(); + assert!(checkpoint.flushed); + let record = recording.sinks.lock().unwrap()[0].clone(); + assert!(record + .ops + .lock() + .unwrap() + .iter() + .all(|op| !matches!(op, SpanOp::Merge(row) if row.end_ms.is_some()))); + + let finalized = flush_managed_run("managed-run", &socket, 5_000) + .await + .unwrap(); + assert!(finalized.flushed); + assert!(record + .ops + .lock() + .unwrap() + .iter() + .any(|op| matches!(op, SpanOp::Merge(row) if row.end_ms == Some(7)))); + + shutdown(&socket).await; + handle.await.unwrap(); +} + #[cfg(feature = "cli")] struct EnvVarGuard { key: &'static str,