diff --git a/CHANGELOG.md b/CHANGELOG.md index 0885538..096514c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,8 +5,33 @@ Versioning and Keep a Changelog conventions. ## [Unreleased] +### Security + +- `opencode serve` now requires HTTP Basic credentials + (`OPENCODE_SERVER_USERNAME`/`OPENCODE_SERVER_PASSWORD`, a per-process + password derived from a random SDK secret and the port), and the SDK's HTTP + bridge sends them. Another local process can no longer drive a served + session through its loopback port, including a port forwarded to the SDK + host for an SSH target. + ### Added +- Manual compaction for Codex and OpenCode. `RuntimeHandle::compact` now runs + `thread/compact/start` on the Codex app server and `/session/{id}/summarize` + on `opencode serve` (which needs a pinned model), reporting the same manual + compaction lifecycle as Claude and ending the invocation when the provider + confirms it. `TurnCapabilities::manual_compaction` states which adapters + support it; the retained driver's `manual_compaction` capability now follows + it instead of naming Claude alone. +- `TurnCapabilities::compaction_instructions` and + `RuntimeDriverCapabilities::compaction_instructions`. Only Claude's native + compaction takes summary instructions; `RuntimeHandle::compact` with + instructions on Codex or OpenCode fails with `CapabilityUnavailable` instead + of compacting without them. +- `CommandSpec::loopback_ports`: loopback ports a provider process serves that + the SDK host must reach. `SshTransport` forwards each one (`ssh -L`, with + `ExitOnForwardFailure`), so `OpenCode::serve()` works on SSH targets; local + execution ignores it. - A pi adapter. `providers::Pi`, behind the new default `pi` feature, drives `pi --mode rpc` (tested with pi 1.0.0): text and reasoning deltas, tool lifecycle events, usage with cost and context-window occupancy, the diff --git a/Cargo.lock b/Cargo.lock index dba1b5e..31bb47f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -967,6 +967,7 @@ version = "0.1.0" dependencies = [ "async-trait", "base64 0.23.1", + "getrandom 0.4.3", "nix", "reqwest", "serde", diff --git a/Cargo.toml b/Cargo.toml index 1010f86..71159fc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -42,6 +42,7 @@ temps-sandbox = ["dep:base64", "dep:reqwest"] [dependencies] async-trait = "0.1" +getrandom = "0.4" base64 = { version = "0.23", optional = true } reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"], optional = true } serde = { version = "1", features = ["derive"] } diff --git a/examples/fullstack-chat/Cargo.lock b/examples/fullstack-chat/Cargo.lock index 561723b..28028e0 100644 --- a/examples/fullstack-chat/Cargo.lock +++ b/examples/fullstack-chat/Cargo.lock @@ -1294,6 +1294,7 @@ version = "0.1.0" dependencies = [ "async-trait", "base64 0.23.1", + "getrandom 0.4.3", "nix", "reqwest", "serde", diff --git a/src/adapter.rs b/src/adapter.rs index 67616eb..a02d094 100644 --- a/src/adapter.rs +++ b/src/adapter.rs @@ -32,6 +32,10 @@ pub struct CommandSpec { pub initial_stdin: Option>, /// Whether stdin must remain available for interaction responses. pub interactive_stdin: bool, + /// Loopback ports the process listens on that the SDK host must reach at + /// the same port, such as `opencode serve`'s HTTP server. Local + /// execution needs nothing; a remote transport forwards each one. + pub loopback_ports: Vec, } impl fmt::Debug for CommandSpec { @@ -50,6 +54,7 @@ impl fmt::Debug for CommandSpec { &self.initial_stdin.as_ref().map(Vec::len), ) .field("interactive_stdin", &self.interactive_stdin) + .field("loopback_ports", &self.loopback_ports) .finish() } } @@ -91,6 +96,7 @@ impl CommandSpec { clear_environment: true, initial_stdin: None, interactive_stdin: false, + loopback_ports: Vec::new(), } } @@ -117,6 +123,7 @@ impl CommandSpec { clear_environment: self.clear_environment, initial_stdin: self.initial_stdin, interactive_stdin: self.interactive_stdin, + loopback_ports: self.loopback_ports, } } } diff --git a/src/providers/claude.rs b/src/providers/claude.rs index 9da0346..7dca5e4 100644 --- a/src/providers/claude.rs +++ b/src/providers/claude.rs @@ -1206,6 +1206,9 @@ impl AgentAdapter for Claude { // when it ends, for automatic and manual compaction alike. compaction_lifecycle: true, live_messages: true, + // `/compact [instructions]` is Claude Code's own command. + manual_compaction: true, + compaction_instructions: true, ..TurnCapabilities::default() } } @@ -1991,12 +1994,7 @@ fn claude_account_usage(value: &Value) -> Option { } /// Whether a prompt is Claude Code's native `/compact` command. -fn is_manual_compaction_prompt(prompt: &str) -> bool { - let prompt = prompt.trim_start(); - prompt - .strip_prefix("/compact") - .is_some_and(|rest| rest.is_empty() || rest.starts_with(char::is_whitespace)) -} +use super::is_manual_compaction_prompt; /// Maximum characters of a provider compaction diagnostic carried in events. const MAX_COMPACTION_ERROR_CHARS: usize = 512; diff --git a/src/providers/codex.rs b/src/providers/codex.rs index c5c5ec5..a8c1c03 100644 --- a/src/providers/codex.rs +++ b/src/providers/codex.rs @@ -680,6 +680,10 @@ impl AgentAdapter for Codex { // The app server reports `contextCompaction` items as they start // and complete; `exec --json` does not surface compaction. compaction_lifecycle: self.app_server_mode(), + // `thread/compact/start` is an app-server request. + manual_compaction: self.app_server_mode(), + // `thread/compact/start` takes no summary instructions. + compaction_instructions: false, live_messages: false, } } diff --git a/src/providers/codex_app_server.rs b/src/providers/codex_app_server.rs index b5d5652..ae596b0 100644 --- a/src/providers/codex_app_server.rs +++ b/src/providers/codex_app_server.rs @@ -83,6 +83,9 @@ struct TurnState { /// Latest active-context occupancy, the pre-compaction size when a /// compaction starts. last_context_tokens: Option, + /// This invocation compacts the thread (`thread/compact/start`) instead of + /// starting a turn. + manual_compaction: bool, } pub(super) fn mark_retained(state: &mut AdapterState) { @@ -154,6 +157,7 @@ pub(super) fn handshake() -> Vec { /// Build the thread and turn parameters this transport sends once the app /// server accepts the handshake. pub(super) fn prepare_turn(request: &TurnRequest, state: &mut AdapterState) -> Result<()> { + super::refuse_compaction_instructions(Provider::Codex, &request.prompt)?; let controls = super::codex::resolve_controls(request)?; let (approval, collaboration, sandbox) = controls.unwrap_or(("never", "default", "danger-full-access")); @@ -218,6 +222,7 @@ pub(super) fn prepare_turn(request: &TurnRequest, state: &mut AdapterState) -> R turn_params, resume: request.session_id.is_some(), model: request.model.clone(), + manual_compaction: super::is_manual_compaction_prompt(&request.prompt), ..TurnState::default() }, ); @@ -296,6 +301,31 @@ pub(super) fn parse_line(line: &str, state: &mut AdapterState) -> Result { + // A manual compaction is complete once the thread says so. + if turn.open_compaction.take().is_some() { + output.events.push(TurnEvent::CompactionCompleted { + compaction: completed_compaction_with( + CompactionTrigger::Manual, + turn.last_context_tokens, + ), + }); + } + output.terminal = true; + } "thread/compacted" => { // Deprecated in favor of the `contextCompaction` item, but still // the only signal from older app servers. Report it only when no @@ -672,7 +726,11 @@ fn notification( // A turn cannot end with its compaction still running; the // application must see it close instead of spinning forever. output.events.push(TurnEvent::CompactionFailed { - trigger: CompactionTrigger::Automatic, + trigger: if turn.manual_compaction { + CompactionTrigger::Manual + } else { + CompactionTrigger::Automatic + }, message: Some(turn.error_message.clone().unwrap_or_else(|| { "Codex ended the turn before the compaction finished.".to_string() })), @@ -740,16 +798,38 @@ fn compaction_item( } return; } - turn.open_compaction = None; + // A manual compaction opened its own lifecycle when it was requested. + let trigger = if turn.manual_compaction { + CompactionTrigger::Manual + } else { + CompactionTrigger::Automatic + }; + let was_open = turn.open_compaction.take().is_some(); + if turn.manual_compaction { + // The compaction is the whole invocation: its completed item ends it, + // whether or not `turn/completed` or `thread/compacted` follows. + output.terminal = true; + if !was_open { + // `thread/compacted` already reported this manual compaction. + return; + } + } turn.saw_compaction_item = true; output.events.push(TurnEvent::CompactionCompleted { - compaction: completed_compaction(turn.last_context_tokens), + compaction: completed_compaction_with(trigger, turn.last_context_tokens), }); } fn completed_compaction(pre_tokens: Option) -> ContextCompaction { + completed_compaction_with(CompactionTrigger::Automatic, pre_tokens) +} + +fn completed_compaction_with( + trigger: CompactionTrigger, + pre_tokens: Option, +) -> ContextCompaction { ContextCompaction { - trigger: CompactionTrigger::Automatic, + trigger, pre_tokens, post_tokens: None, dropped_tokens: None, @@ -1521,6 +1601,118 @@ mod tests { ); } + #[test] + fn a_manual_compaction_compacts_the_opened_thread_instead_of_starting_a_turn() { + let mut request = TurnRequest::new(Provider::Codex, ".", "/compact"); + request.session_id = Some("thread-9".into()); + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state).unwrap(); + parse_line(r#"{"id":1,"result":{}}"#, &mut state).unwrap(); + let opened = parse_line( + r#"{"id":2,"result":{"thread":{"id":"thread-9"}}}"#, + &mut state, + ) + .unwrap(); + let compact = decode(opened.writes[0].clone()); + assert_eq!(compact["method"], "thread/compact/start"); + assert_eq!(compact["params"]["threadId"], "thread-9"); + assert!(opened.events.iter().any(|event| matches!( + event, + TurnEvent::CompactionStarted { + trigger: CompactionTrigger::Manual + } + ))); + + let done = parse_line( + r#"{"jsonrpc":"2.0","method":"thread/compacted","params":{"threadId":"thread-9"}}"#, + &mut state, + ) + .unwrap(); + assert!(done.terminal); + assert!(done.events.iter().any(|event| matches!( + event, + TurnEvent::CompactionCompleted { compaction } if compaction.trigger == CompactionTrigger::Manual + ))); + } + + #[test] + fn a_retained_manual_compaction_follows_only_its_own_threads_turn() { + let mut request = TurnRequest::new(Provider::Codex, ".", "/compact"); + request.session_id = Some("thread-9".into()); + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state).unwrap(); + mark_retained(&mut state); + parse_line( + r#"{"id":2,"result":{"thread":{"id":"thread-9"}}}"#, + &mut state, + ) + .unwrap(); + parse_line(r#"{"id":4,"result":{}}"#, &mut state).unwrap(); + + // A late turn from another thread on the shared connection is not + // this compaction's turn. + parse_line( + r#"{"method":"turn/started","params":{"threadId":"thread-other","turn":{"id":"turn-x"}}}"#, + &mut state, + ) + .unwrap(); + let stray = parse_line( + r#"{"method":"item/completed","params":{"threadId":"thread-other","turnId":"turn-x","item":{"type":"contextCompaction","id":"c-x"}}}"#, + &mut state, + ) + .unwrap(); + assert!(!stray.terminal); + assert_eq!(stray.events, [] as [crate::TurnEvent; 0]); + + parse_line( + r#"{"method":"turn/started","params":{"threadId":"thread-9","turn":{"id":"turn-c"}}}"#, + &mut state, + ) + .unwrap(); + let done = parse_line( + r#"{"method":"item/completed","params":{"threadId":"thread-9","turnId":"turn-c","item":{"type":"contextCompaction","id":"c-1"}}}"#, + &mut state, + ) + .unwrap(); + // The completed compaction item ends the invocation by itself. + assert!(done.terminal); + assert!(done.events.iter().any(|event| matches!( + event, + TurnEvent::CompactionCompleted { compaction } if compaction.trigger == CompactionTrigger::Manual + ))); + } + + #[test] + fn a_retained_manual_compaction_also_ends_on_the_older_thread_compacted() { + let mut request = TurnRequest::new(Provider::Codex, ".", "/compact"); + request.session_id = Some("thread-9".into()); + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state).unwrap(); + mark_retained(&mut state); + parse_line( + r#"{"id":2,"result":{"thread":{"id":"thread-9"}}}"#, + &mut state, + ) + .unwrap(); + parse_line(r#"{"id":4,"result":{}}"#, &mut state).unwrap(); + let other = parse_line( + r#"{"method":"thread/compacted","params":{"threadId":"thread-other"}}"#, + &mut state, + ) + .unwrap(); + assert!(!other.terminal); + let done = parse_line( + r#"{"method":"thread/compacted","params":{"threadId":"thread-9"}}"#, + &mut state, + ) + .unwrap(); + assert!(done.terminal); + assert!(done + .events + .iter() + .any(|event| matches!(event, TurnEvent::CompactionCompleted { .. }))); + } + #[test] fn resume_and_writer_conflict_fork_exclude_historical_turns() { let mut request = TurnRequest::new(Provider::Codex, ".", "continue"); diff --git a/src/providers/mod.rs b/src/providers/mod.rs index df5dd7b..5921fd3 100644 --- a/src/providers/mod.rs +++ b/src/providers/mod.rs @@ -73,3 +73,54 @@ fn merge_usage(target: &mut Usage, incoming: &Usage) { .or_else(|| target.context_window.clone()); target.cost_usd = incoming.cost_usd.or(target.cost_usd); } + +/// Whether a turn prompt is a manual compaction request (`/compact` plus +/// optional instructions), the form `CompactionInput` produces. +pub(crate) fn is_manual_compaction_prompt(prompt: &str) -> bool { + let prompt = prompt.trim_start(); + prompt + .strip_prefix("/compact") + .is_some_and(|rest| rest.is_empty() || rest.starts_with(char::is_whitespace)) +} + +/// Refuse a manual compaction request carrying summary instructions on an +/// adapter whose native compaction cannot take them, instead of dropping them. +#[cfg(any(feature = "codex", feature = "opencode"))] +pub(crate) fn refuse_compaction_instructions( + provider: crate::Provider, + prompt: &str, +) -> crate::Result<()> { + let has_instructions = prompt + .trim_start() + .strip_prefix("/compact") + .is_some_and(|rest| !rest.trim().is_empty()); + if is_manual_compaction_prompt(prompt) && has_instructions { + return Err(crate::RuntimeError::InvalidRequest { + field: "prompt", + message: format!( + "{provider} compacts natively without summary instructions; omit them" + ), + }); + } + Ok(()) +} + +#[cfg(all(test, any(feature = "codex", feature = "opencode")))] +mod compaction_prompt_tests { + use super::*; + + #[test] + fn only_a_bare_compact_is_accepted_without_instruction_support() { + assert!(is_manual_compaction_prompt("/compact")); + assert!(is_manual_compaction_prompt(" /compact keep the API notes")); + assert!(!is_manual_compaction_prompt("/compaction please")); + assert!(refuse_compaction_instructions(crate::Provider::Codex, "/compact").is_ok()); + assert!(refuse_compaction_instructions(crate::Provider::Codex, "/compact ").is_ok()); + assert!(refuse_compaction_instructions(crate::Provider::Codex, "hello").is_ok()); + let error = + refuse_compaction_instructions(crate::Provider::OpenCode, "/compact keep notes") + .unwrap_err() + .to_string(); + assert!(error.contains("without summary instructions"), "{error}"); + } +} diff --git a/src/providers/opencode.rs b/src/providers/opencode.rs index ccb01f0..2879db2 100644 --- a/src/providers/opencode.rs +++ b/src/providers/opencode.rs @@ -46,6 +46,58 @@ pub enum OpenCodeTurnMode { Serve, } +/// Username `opencode serve` checks alongside `OPENCODE_SERVER_PASSWORD`. +const SERVER_USERNAME: &str = "opencode"; + +/// The password `opencode serve` requires on `port`. Derived from a secret this +/// SDK process draws once, so the later turns of a retained server recompute +/// it from the port alone and no credential is stored in turn state. +fn server_password(port: u16) -> Result { + use sha2::Digest; + static SECRET: std::sync::OnceLock> = std::sync::OnceLock::new(); + let secret = SECRET + .get_or_init(|| { + let mut secret = [0_u8; 32]; + getrandom::fill(&mut secret).ok().map(|()| secret) + }) + .ok_or_else(|| RuntimeError::Protocol { + provider: Provider::OpenCode, + message: "could not draw a random password for the OpenCode server".into(), + })?; + let mut hasher = sha2::Sha256::new(); + hasher.update(secret); + hasher.update(b"opencode-serve-password"); + hasher.update(port.to_be_bytes()); + Ok(hasher + .finalize() + .iter() + .map(|byte| format!("{byte:02x}")) + .collect()) +} + +/// `Authorization` header value for HTTP Basic credentials. +fn basic_authorization(password: &str) -> String { + const ALPHABET: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"; + let raw = format!("{SERVER_USERNAME}:{password}"); + let bytes = raw.as_bytes(); + let mut encoded = String::with_capacity(bytes.len().div_ceil(3) * 4); + for chunk in bytes.chunks(3) { + let triple = (u32::from(chunk[0]) << 16) + | (u32::from(*chunk.get(1).unwrap_or(&0)) << 8) + | u32::from(*chunk.get(2).unwrap_or(&0)); + for index in 0..4 { + if index <= chunk.len() { + encoded.push(char::from( + ALPHABET[((triple >> (18 - 6 * index)) & 0x3f) as usize], + )); + } else { + encoded.push('='); + } + } + } + format!("Basic {encoded}") +} + /// OpenCode CLI adapter using `opencode run --format json` or `opencode serve`. #[derive(Debug, Clone, Default)] pub struct OpenCode { @@ -96,6 +148,9 @@ impl OpenCode { "--port".into(), port.to_string().into(), ]); + // The client below talks to this port on the SDK host's loopback; + // a remote transport forwards it there. + spec.loopback_ports.push(port); // The policy travels in the environment rather than argv because it is // this turn's entire enforcement boundary, and `clear_environment` // means nothing reaches the child that was not put here deliberately. @@ -103,6 +158,15 @@ impl OpenCode { "OPENCODE_CONFIG_CONTENT".into(), super::opencode_serve::permission_config(request)?.into(), ); + // The server answers only requests carrying this process's password, + // so another local process (or one reaching a forwarded port on the + // SDK host) cannot drive the session. + spec.environment + .insert("OPENCODE_SERVER_USERNAME".into(), SERVER_USERNAME.into()); + spec.environment.insert( + "OPENCODE_SERVER_PASSWORD".into(), + server_password(port)?.into(), + ); Ok(spec) } @@ -129,6 +193,8 @@ impl AgentAdapter for OpenCode { // summary message, and `session.compacted` over its event stream; // `opencode run --format json` does not. compaction_lifecycle: self.serve_mode(), + // `/session/{id}/summarize` is a server request. + manual_compaction: self.serve_mode(), ..TurnCapabilities::default() } } @@ -292,6 +358,7 @@ impl AgentAdapter for OpenCode { if !self.serve_mode() { return Ok(()); } + super::refuse_compaction_instructions(Provider::OpenCode, &request.prompt)?; // Bind a throwaway listener so the OS picks a free port, then drop it // so the server can bind the same one. `opencode serve --port 0` does // not do this: it falls back to its fixed default port instead. @@ -315,6 +382,7 @@ impl AgentAdapter for OpenCode { if !self.serve_mode() { return self.prepare_turn(request, state); } + super::refuse_compaction_instructions(Provider::OpenCode, &request.prompt)?; let port = process_hint .and_then(|port| u16::try_from(port).ok()) .ok_or_else(|| RuntimeError::Protocol { @@ -347,9 +415,16 @@ impl AgentAdapter for OpenCode { _request: &TurnRequest, state: &AdapterState, ) -> Result> { - Ok(super::opencode_serve::turn_port(state) - .filter(|_| self.serve_mode()) - .map(super::opencode_http::connect)) + let Some(port) = super::opencode_serve::turn_port(state).filter(|_| self.serve_mode()) + else { + return Ok(None); + }; + Ok(Some(super::opencode_http::connect( + super::opencode_http::Endpoint { + port, + authorization: Some(basic_authorization(&server_password(port)?)), + }, + ))) } fn interrupt_request(&self, state: &AdapterState) -> Option> { @@ -604,6 +679,40 @@ impl AgentAdapter for OpenCode { mod tests { use super::*; + #[test] + fn basic_authorization_encodes_the_server_username_and_password() { + // `opencode:s3cret` in Base64, as `curl -u` sends it. + assert_eq!(basic_authorization("s3cret"), "Basic b3BlbmNvZGU6czNjcmV0"); + assert_eq!(basic_authorization(""), "Basic b3BlbmNvZGU6"); + } + + #[test] + fn the_serve_command_requires_a_per_port_server_password() { + let first = server_password(4100).unwrap(); + assert_eq!(first.len(), 64); + assert_eq!( + first, + server_password(4100).unwrap(), + "stable for a retained server" + ); + assert_ne!(first, server_password(4101).unwrap()); + + let adapter = OpenCode::serve(); + let request = TurnRequest::new(Provider::OpenCode, "/tmp/work", "hi"); + let spec = adapter.serve_command(&request, 4100).unwrap(); + assert_eq!( + spec.environment + .get(std::ffi::OsStr::new("OPENCODE_SERVER_PASSWORD")), + Some(&std::ffi::OsString::from(first)) + ); + assert_eq!( + spec.environment + .get(std::ffi::OsStr::new("OPENCODE_SERVER_USERNAME")), + Some(&std::ffi::OsString::from("opencode")) + ); + assert_eq!(spec.loopback_ports, vec![4100]); + } + #[test] fn accumulates_step_usage() { let adapter = OpenCode::default(); diff --git a/src/providers/opencode_http.rs b/src/providers/opencode_http.rs index fd288c3..972cca7 100644 --- a/src/providers/opencode_http.rs +++ b/src/providers/opencode_http.rs @@ -84,11 +84,41 @@ pub(super) const FRAME_ERROR: &str = "@error"; /// missed, and only the bridge can know when the stream is actually open. pub(super) const FRAME_SUBSCRIBED: &str = "@subscribed"; +/// Where the bridge reaches one `opencode serve` process: its loopback port +/// and the `Authorization` header value the server requires, if any. +#[derive(Clone)] +pub(super) struct Endpoint { + pub(super) port: u16, + pub(super) authorization: Option, +} + +impl std::fmt::Debug for Endpoint { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("Endpoint") + .field("port", &self.port) + .field( + "authorization", + &self.authorization.as_ref().map(|_| ""), + ) + .finish() + } +} + +impl Endpoint { + fn authorization_header(&self) -> String { + self.authorization + .as_deref() + .map(|value| format!("Authorization: {value}\r\n")) + .unwrap_or_default() + } +} + /// Start the bridge for a turn and hand the runtime its protocol streams. /// /// The returned streams are live immediately; the task behind them polls the /// server for readiness first and emits [`FRAME_READY`] once it answers. -pub(super) fn connect(port: u16) -> ProtocolStreams { +pub(super) fn connect(endpoint: Endpoint) -> ProtocolStreams { // Two independent pipes rather than the two halves of one. Splitting a // single duplex stream would keep it alive until *both* halves drop, so // the runtime closing its writer at the end of a turn would never reach @@ -97,7 +127,7 @@ pub(super) fn connect(port: u16) -> ProtocolStreams { // direction close on its own. let (runtime_reader, bridge_writer) = tokio::io::duplex(BRIDGE_BUFFER_BYTES); let (bridge_reader, runtime_writer) = tokio::io::duplex(BRIDGE_BUFFER_BYTES); - tokio::spawn(run_bridge(port, bridge_reader, bridge_writer)); + tokio::spawn(run_bridge(endpoint, bridge_reader, bridge_writer)); ProtocolStreams { reader: Box::new(runtime_reader), writer: Box::new(runtime_writer), @@ -110,10 +140,10 @@ fn address(port: u16) -> SocketAddr { /// Drive one turn's HTTP traffic until the state machine stops writing or the /// server stops answering. -async fn run_bridge(port: u16, incoming: DuplexStream, outgoing: DuplexStream) { +async fn run_bridge(endpoint: Endpoint, incoming: DuplexStream, outgoing: DuplexStream) { let outgoing = std::sync::Arc::new(tokio::sync::Mutex::new(outgoing)); - if let Err(error) = wait_until_ready(port).await { + if let Err(error) = wait_until_ready(&endpoint).await { emit_error(&outgoing, &error).await; return; } @@ -145,14 +175,20 @@ async fn run_bridge(port: u16, incoming: DuplexStream, outgoing: DuplexStream) { if command.method == "SUBSCRIBE" { if subscription.is_none() { subscription = Some(tokio::spawn(stream_events( - port, + endpoint.clone(), command.path, std::sync::Arc::clone(&outgoing), ))); } continue; } - let outcome = perform(port, &command.method, &command.path, command.body.as_ref()).await; + let outcome = perform( + &endpoint, + &command.method, + &command.path, + command.body.as_ref(), + ) + .await; let frame = match outcome { Ok((status, body)) => json!({ "type": FRAME_RESPONSE, @@ -177,11 +213,11 @@ async fn run_bridge(port: u16, incoming: DuplexStream, outgoing: DuplexStream) { } /// Poll a cheap, always-safe endpoint until the server answers. -async fn wait_until_ready(port: u16) -> std::result::Result<(), String> { +async fn wait_until_ready(endpoint: &Endpoint) -> std::result::Result<(), String> { for attempt in 0..READINESS_ATTEMPTS { match tokio::time::timeout( Duration::from_secs(1), - perform(port, "GET", "/global/health", None), + perform(endpoint, "GET", "/global/health", None), ) .await { @@ -201,13 +237,15 @@ async fn wait_until_ready(port: u16) -> std::result::Result<(), String> { /// Perform one request on its own connection and read the whole response. async fn perform( - port: u16, + endpoint: &Endpoint, method: &str, path: &str, body: Option<&Value>, ) -> std::io::Result<(u16, Value)> { + let port = endpoint.port; let mut stream = TcpStream::connect(address(port)).await?; stream.set_nodelay(true).ok(); + let authorization = endpoint.authorization_header(); let encoded = body.map(ToString::to_string).unwrap_or_default(); let content_type = if body.is_some() { "Content-Type: application/json\r\n" @@ -215,7 +253,7 @@ async fn perform( "" }; let head = format!( - "{method} {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\n\ + "{method} {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\n{authorization}\ Accept: application/json\r\nConnection: close\r\n{content_type}\ Content-Length: {}\r\n\r\n", encoded.len() @@ -244,10 +282,11 @@ async fn perform( /// and the turn is failed with a real diagnostic instead of hanging until the /// turn deadline. async fn stream_events( - port: u16, + endpoint: Endpoint, path: String, outgoing: std::sync::Arc>, ) { + let port = endpoint.port; let stream = match TcpStream::connect(address(port)).await { Ok(stream) => stream, Err(error) => { @@ -262,7 +301,8 @@ async fn stream_events( stream.set_nodelay(true).ok(); let mut reader = BufReader::new(stream); let head = format!( - "GET {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nAccept: text/event-stream\r\nCache-Control: no-cache\r\n\r\n" + "GET {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\n{}Accept: text/event-stream\r\nCache-Control: no-cache\r\n\r\n", + endpoint.authorization_header() ); { let inner = reader.get_mut(); @@ -642,7 +682,9 @@ mod tests { let mut reader = BufReader::new(&mut socket); let mut request = String::new(); reader.read_line(&mut request).await.unwrap(); - // Drain headers. + // Drain headers; like `opencode serve` with a server + // password, answer nothing without the credential. + let mut authorized = false; loop { let mut header = String::new(); if reader.read_line(&mut header).await.unwrap() == 0 @@ -650,8 +692,11 @@ mod tests { { break; } + authorized |= header.trim_end() == "Authorization: Basic Zml4dHVyZQ=="; } - let response: &[u8] = if request.contains("/event") { + let response: &[u8] = if !authorized { + b"HTTP/1.1 401 Unauthorized\r\nContent-Length: 0\r\n\r\n" + } else if request.contains("/event") { b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n\ 1E\r\ndata: {\"type\":\"session.idle\"}\n\r\n\ 1\r\n\n\r\n" @@ -668,7 +713,10 @@ mod tests { } }); - let streams = connect(port); + let streams = connect(Endpoint { + port, + authorization: Some("Basic Zml4dHVyZQ==".into()), + }); let mut writer = streams.writer; let mut lines = BufReader::new(streams.reader).lines(); @@ -739,7 +787,10 @@ mod tests { } }); - let streams = connect(port); + let streams = connect(Endpoint { + port, + authorization: None, + }); let mut writer = streams.writer; let mut lines = BufReader::new(streams.reader).lines(); assert!(lines.next_line().await.unwrap().is_some(), "ready frame"); diff --git a/src/providers/opencode_serve.rs b/src/providers/opencode_serve.rs index 983eb57..1d87627 100644 --- a/src/providers/opencode_serve.rs +++ b/src/providers/opencode_serve.rs @@ -65,6 +65,11 @@ pub(super) struct TurnState { parts: Vec, /// Native `{providerID, modelID}` selection, when the caller pinned one. model: Option, + /// This invocation summarizes the session (`/session/{id}/summarize`) + /// instead of sending a prompt. + manual_compaction: bool, + /// OpenCode confirmed this invocation's manual compaction. + compacted: bool, /// Session being prompted. session_id: Option, /// Whether every permission must be refused without consulting the @@ -337,6 +342,7 @@ pub(super) fn prepare_turn(request: &TurnRequest, state: &mut AdapterState, port resume: request.session_id.clone(), parts, model: model_selection(request.model.as_deref()), + manual_compaction: super::is_manual_compaction_prompt(&request.prompt), plan_mode: is_plan_mode(request), tools_denied: request .launch_context @@ -457,6 +463,26 @@ fn prompt(turn: &TurnState, output: &mut AdapterOutput) -> Result<()> { let Some(session) = turn.session_id.as_deref() else { return Ok(()); }; + if turn.manual_compaction { + // OpenCode summarizes with an explicit model; its compaction part and + // `session.compacted` then report the lifecycle like an automatic one. + let Some(model) = &turn.model else { + return Err(RuntimeError::InvalidRequest { + field: "model", + message: "Select an OpenCode model before compacting this session.".into(), + }); + }; + output.writes.push(encode(&post( + Some(ID_PROMPT), + format!( + "/session/{}/summarize?directory={}", + query_escape(session), + query_escape(&turn.directory) + ), + model.clone(), + ))?); + return Ok(()); + } let mut body = json!({"parts": turn.parts}); if let Some(model) = &turn.model { body["model"] = model.clone(); @@ -691,6 +717,11 @@ fn event( duration_ms: None, }, }); + if turn.manual_compaction { + // The summary is the whole invocation; it has no reply text. + turn.compacted = true; + output.terminal = true; + } } "session.error" => { turn.error_message = properties @@ -712,6 +743,16 @@ fn event( } if let Some(message) = turn.error_message.clone() { fail(turn, state, output, &message, None); + } else if turn.manual_compaction { + if !turn.compacted { + fail( + turn, + state, + output, + "OpenCode went idle without confirming the compaction.", + None, + ); + } } else if state.result.text.trim().is_empty() && !turn.saw_activity { // A turn that produced no text and no tool call is a real // failure, not an empty success: it is what a silently @@ -1043,6 +1084,120 @@ mod tests { state } + #[test] + fn a_manual_compaction_summarizes_the_session_with_its_model() { + let mut request = request(PermissionMode::Default); + request.prompt = "/compact".into(); + request.model = Some("anthropic/claude-sonnet-4-5".into()); + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state, 4242); + parse_line(&json!({"type": FRAME_READY}).to_string(), &mut state).unwrap(); + parse_line( + &json!({"type": FRAME_RESPONSE, "id": ID_SESSION, "status": 200, + "body": {"id": "session-1"}}) + .to_string(), + &mut state, + ) + .unwrap(); + let sent = parse_line(&json!({"type": FRAME_SUBSCRIBED}).to_string(), &mut state).unwrap(); + let frame = decode(&sent.writes[0]); + assert!(frame["path"] + .as_str() + .unwrap() + .starts_with("/session/session-1/summarize")); + assert_eq!(frame["body"]["providerID"], json!("anthropic")); + assert_eq!(frame["body"]["modelID"], json!("claude-sonnet-4-5")); + + // Without a model there is nothing to summarize with. + let mut request = request.clone(); + request.model = None; + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state, 4242); + parse_line(&json!({"type": FRAME_READY}).to_string(), &mut state).unwrap(); + parse_line( + &json!({"type": FRAME_RESPONSE, "id": ID_SESSION, "status": 200, + "body": {"id": "session-1"}}) + .to_string(), + &mut state, + ) + .unwrap(); + assert!(parse_line(&json!({"type": FRAME_SUBSCRIBED}).to_string(), &mut state).is_err()); + } + + /// A manual compaction invocation with its summarize request accepted. + fn summarizing_turn() -> AdapterState { + let mut request = request(PermissionMode::Default); + request.prompt = "/compact".into(); + request.model = Some("anthropic/claude-sonnet-4-5".into()); + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state, 4242); + for frame in [ + json!({"type": FRAME_READY}), + json!({"type": FRAME_RESPONSE, "id": ID_SESSION, "status": 200, "body": {"id": "session-1"}}), + json!({"type": FRAME_SUBSCRIBED}), + json!({"type": FRAME_RESPONSE, "id": ID_PROMPT, "status": 200, "body": true}), + ] { + parse_line(&frame.to_string(), &mut state).unwrap(); + } + state + } + + #[test] + fn a_confirmed_manual_compaction_ends_the_invocation_without_a_reply() { + let mut state = summarizing_turn(); + let started = parse_line( + &sse( + "message.part.updated", + json!({"sessionID": "session-1", "part": { + "id": "prt-c1", "messageID": "msg-u2", "sessionID": "session-1", + "type": "compaction", "auto": false + }}), + ), + &mut state, + ) + .unwrap(); + assert!(!started.terminal); + let compacted = parse_line( + &sse("session.compacted", json!({"sessionID": "session-1"})), + &mut state, + ) + .unwrap(); + assert!(compacted.terminal); + assert!(compacted.events.iter().any(|event| matches!( + event, + TurnEvent::CompactionCompleted { compaction } if compaction.trigger == CompactionTrigger::Manual + ))); + // An idle that follows is not "finished without producing a reply". + let idle = parse_line( + &sse("session.idle", json!({"sessionID": "session-1"})), + &mut state, + ) + .unwrap(); + assert!(!idle + .events + .iter() + .any(|event| matches!(event, TurnEvent::CompactionFailed { .. }))); + assert!(state.terminal_failure.is_none()); + assert_ne!(state.result.status, RunStatus::Failed); + } + + #[test] + fn an_unconfirmed_manual_compaction_fails_clearly() { + let mut state = summarizing_turn(); + let idle = parse_line( + &sse("session.idle", json!({"sessionID": "session-1"})), + &mut state, + ) + .unwrap(); + assert!(idle.terminal); + assert_eq!(state.result.status, RunStatus::Failed); + let failure = format!("{:?}", state.terminal_failure); + assert!( + failure.contains("without confirming the compaction"), + "{failure}" + ); + } + fn permission_event(kind: &str) -> String { json!({"type": FRAME_EVENT, "event": { "type": "permission.asked", diff --git a/src/retained.rs b/src/retained.rs index d1159fe..b0da970 100644 --- a/src/retained.rs +++ b/src/retained.rs @@ -342,6 +342,10 @@ pub struct RuntimeDriverCapabilities { /// The harness supports provider-native manual compaction of an existing session. #[serde(default)] pub manual_compaction: bool, + /// Manual compaction accepts summary instructions + /// ([`CompactionInput::instructions`]). + #[serde(default)] + pub compaction_instructions: bool, /// The driver emits active context-window occupancy snapshots when available. #[serde(default)] pub context_window_usage: bool, @@ -536,7 +540,12 @@ impl RuntimeTurnExecutor for AgentRuntime { live_interactions: permissions .is_some_and(|support| support.live_approvals || support.live_questions), configurable_auto_compaction: provider == Provider::Claude, - manual_compaction: provider == Provider::Claude, + manual_compaction: self + .turn_capabilities(provider) + .is_ok_and(|capabilities| capabilities.manual_compaction), + compaction_instructions: self + .turn_capabilities(provider) + .is_ok_and(|capabilities| capabilities.compaction_instructions), context_window_usage: provider == Provider::Claude || self .turn_capabilities(provider) @@ -844,6 +853,19 @@ impl RuntimeHandle { ), )); } + if input.instructions.is_some() && !self.driver.compaction_instructions { + return Err(lifecycle_failure( + Some(self.runtime_id.clone()), + Some(input.invocation_id), + RuntimeFailureKind::CapabilityUnavailable, + RetryAdvice::Never, + DeliveryState::NotSent, + format!( + "{} compacts natively without summary instructions; omit them", + self.provider + ), + )); + } Arc::clone(&self.backend) .start_turn(input.into_turn_input(), None) .await @@ -2363,6 +2385,7 @@ mod tests { live_interactions: true, configurable_auto_compaction: true, manual_compaction: true, + compaction_instructions: true, context_window_usage: true, native_image_attachments: true, compaction_lifecycle: true, @@ -2422,6 +2445,7 @@ mod tests { RuntimeDriverCapabilities { session_resume: true, manual_compaction: true, + compaction_instructions: true, compaction_lifecycle: true, ..RuntimeDriverCapabilities::default() } @@ -2989,6 +3013,7 @@ mod tests { retained_process: true, session_resume: true, manual_compaction: true, + compaction_instructions: true, ..RuntimeDriverCapabilities::default() } } diff --git a/src/runtime.rs b/src/runtime.rs index 8f6e6f8..95f3d22 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -1291,12 +1291,22 @@ async fn write_provider_frames( } /// Close a retained turn that reached its terminal frame or quiet completion. +/// Whether this turn was a provider-native manual compaction, which has no +/// reply text by design. Adapters without native compaction (pi, Codex exec, +/// OpenCode run) send `/compact` as an ordinary prompt, so an empty reply +/// there still warns. +fn performed_native_compaction(adapter: &dyn AgentAdapter, request: &TurnRequest) -> bool { + adapter.turn_capabilities().manual_compaction + && crate::providers::is_manual_compaction_prompt(&request.prompt) +} + async fn finish_retained_turn( - provider: Provider, + adapter: &dyn AgentAdapter, request: &TurnRequest, events: &dyn EventSink, mut state: AdapterState, ) -> Result { + let provider = request.provider; if let Some(failure) = state.terminal_failure.take() { return Err(RuntimeError::ProcessFailed { provider, @@ -1307,7 +1317,7 @@ async fn finish_retained_turn( delivery: failure.delivery, }); } - if state.result.text.is_empty() { + if state.result.text.is_empty() && !performed_native_compaction(adapter, request) { events .emit(TurnEvent::Warning { message: format!("{provider} completed without a text response"), @@ -4803,7 +4813,7 @@ impl AgentRuntime { continue; } () = tokio::time::sleep_until(quiet_completion.unwrap_or(deadline)), if quiet_completion.is_some() => { - return finish_retained_turn(provider, request, events, state) + return finish_retained_turn(adapter, request, events, state) .await .map(RetainedTurnEnd::Completed); } @@ -4931,7 +4941,7 @@ impl AgentRuntime { } } if output.terminal { - return finish_retained_turn(provider, request, events, state) + return finish_retained_turn(adapter, request, events, state) .await .map(RetainedTurnEnd::Completed); } @@ -5516,7 +5526,7 @@ impl AgentRuntime { delivery, }); } - if state.result.text.is_empty() { + if state.result.text.is_empty() && !performed_native_compaction(adapter.as_ref(), request) { let _delivery = trace.event_delivery(); events .emit(TurnEvent::Warning { @@ -6152,6 +6162,45 @@ impl EventSink for CompactionGuard<'_> { mod tests { use super::*; + #[test] + fn only_native_compaction_silences_the_empty_reply_warning() { + let compact = TurnRequest::new(Provider::Codex, ".", "/compact"); + let prompt = TurnRequest::new(Provider::Codex, ".", "hello"); + #[cfg(feature = "codex")] + { + assert!(performed_native_compaction( + &crate::providers::Codex::app_server(), + &compact + )); + // `codex exec` sends `/compact` as an ordinary prompt. + assert!(!performed_native_compaction( + &crate::providers::Codex::default(), + &compact + )); + assert!(!performed_native_compaction( + &crate::providers::Codex::app_server(), + &prompt + )); + } + #[cfg(feature = "opencode")] + { + assert!(performed_native_compaction( + &crate::providers::OpenCode::serve(), + &compact + )); + assert!(!performed_native_compaction( + &crate::providers::OpenCode::default(), + &compact + )); + } + #[cfg(feature = "claude")] + assert!(performed_native_compaction( + &crate::providers::Claude::default(), + &compact + )); + let _ = (compact, prompt); + } + fn claude_args(args: &[&str]) -> Vec { args.iter().map(std::ffi::OsString::from).collect() } diff --git a/src/ssh.rs b/src/ssh.rs index e7cee0a..3b726c5 100644 --- a/src/ssh.rs +++ b/src/ssh.rs @@ -337,7 +337,22 @@ impl SshTransport { } fn base_command(&self) -> Command { + self.base_command_forwarding(&[]) + } + + /// `base_command` plus `-L` for each loopback port the remote process + /// serves, so the SDK host reaches it at the same local port. A failed + /// forward fails the connection instead of leaving the client unable to + /// reach the server. + fn base_command_forwarding(&self, loopback_ports: &[u16]) -> Command { let mut command = Command::new(&self.executable); + for port in loopback_ports { + command + .arg("-o") + .arg("ExitOnForwardFailure=yes") + .arg("-L") + .arg(format!("127.0.0.1:{port}:127.0.0.1:{port}")); + } command .arg("-o") .arg(format!( @@ -682,7 +697,7 @@ impl ExecutionTransport for SshTransport { let control_directory = self.stage_remote_launcher(launcher.as_bytes()).await?; let pid_file = format!("{control_directory}/pid"); let remote_command = build_remote_command(&request, &control_directory)?; - let mut command = self.base_command(); + let mut command = self.base_command_forwarding(&request.command.loopback_ports); command .arg(remote_command) .stdin(std::process::Stdio::piped()) @@ -1371,6 +1386,33 @@ mod tests { } } + #[test] + fn loopback_ports_are_forwarded_before_the_destination() { + let transport = SshTransport::builder("example.test") + .user("fleet") + .build() + .unwrap(); + let command = transport.base_command_forwarding(&[4242]); + let args: Vec = command + .as_std() + .get_args() + .map(|arg| arg.to_string_lossy().into_owned()) + .collect(); + let forward = args.iter().position(|arg| arg == "-L").expect("forward"); + assert_eq!(args[forward + 1], "127.0.0.1:4242:127.0.0.1:4242"); + assert!(args.contains(&"ExitOnForwardFailure=yes".to_string())); + let destination = args + .iter() + .position(|arg| arg == "fleet@example.test") + .unwrap(); + assert!(forward < destination); + assert!(!transport + .base_command() + .as_std() + .get_args() + .any(|arg| arg == "-L")); + } + #[test] fn quotes_untrusted_remote_values_without_shell_injection() { let mut command = CommandSpec::new("/usr/local/bin/codex"); diff --git a/src/types.rs b/src/types.rs index e0369a7..768ce76 100644 --- a/src/types.rs +++ b/src/types.rs @@ -56,6 +56,15 @@ pub struct LaunchContextCapabilities { /// not which launch-context fields it can enforce. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] pub struct TurnCapabilities { + /// The adapter performs a provider-native compaction when a turn's prompt + /// is a manual compaction request (see `RuntimeHandle::compact`). + #[serde(default)] + pub manual_compaction: bool, + /// The adapter's native manual compaction accepts summary instructions + /// (`/compact `). Codex and OpenCode compact natively + /// without them, so they refuse instructions rather than drop them. + #[serde(default)] + pub compaction_instructions: bool, /// The adapter delivers image attachments as native provider image inputs. /// /// Applications and the retained runtime stop describing those files in diff --git a/tests/codex_app_server.rs b/tests/codex_app_server.rs index 3c58d86..1c5f170 100644 --- a/tests/codex_app_server.rs +++ b/tests/codex_app_server.rs @@ -291,6 +291,25 @@ async fn serve( ) .await; } + // Played exactly as `codex app-server` 0.159 answers it: an empty + // result, then a turn of its own wrapping one compaction item. + Some("thread/compact/start") => { + send(&mut output, json!({"jsonrpc":"2.0","id":id,"result":{}})).await; + for frame in [ + json!({"jsonrpc":"2.0","method":"turn/started","params":{ + "threadId":"thread-fixture","turn":{"id":"turn-compact","status":"inProgress"}}}), + json!({"jsonrpc":"2.0","method":"item/started","params":{ + "threadId":"thread-fixture","turnId":"turn-compact", + "item":{"type":"contextCompaction","id":"compact-1"}}}), + json!({"jsonrpc":"2.0","method":"item/completed","params":{ + "threadId":"thread-fixture","turnId":"turn-compact", + "item":{"type":"contextCompaction","id":"compact-1"}}}), + json!({"jsonrpc":"2.0","method":"turn/completed","params":{ + "threadId":"thread-fixture","turn":{"id":"turn-compact","status":"completed"}}}), + ] { + send(&mut output, frame).await; + } + } Some("turn/start") => { send( &mut output, @@ -901,6 +920,118 @@ async fn a_retained_process_that_dies_mid_compaction_never_leaves_it_open() { ); } +#[tokio::test] +async fn a_manual_compaction_compacts_the_thread_once() { + let events = Collector::default(); + let mut request = request(); + request.prompt = "/compact".into(); + request.session_id = Some("thread-fixture".into()); + let transport = AppServer::new(Script::AsyncQuestion); + runtime(transport.clone()) + .run(request, &events, None) + .await + .unwrap(); + assert_eq!(transport.method_count("thread/compact/start"), 1); + assert_eq!(transport.method_count("turn/start"), 0); + assert!( + !events + .events() + .iter() + .any(|event| matches!(event, TurnEvent::Warning { .. })), + "a compaction has no reply text to warn about" + ); + let compaction = compaction_events(&events.events()); + assert!( + matches!( + compaction.as_slice(), + [ + TurnEvent::CompactionStarted { .. }, + TurnEvent::CompactionCompleted { .. } + ] + ), + "{compaction:?}" + ); +} + +#[tokio::test] +async fn a_retained_process_compacts_its_thread_after_a_turn() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + assert!(handle.driver_capabilities().manual_compaction); + handle + .start_turn(TurnInput::new( + InvocationId::new("before").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + + let (mut stream, completion) = handle + .compact(temps_agent_runtime::retained::CompactionInput::new( + InvocationId::new("compact").unwrap(), + )) + .await + .unwrap() + .into_parts(); + let mut observed = Vec::new(); + let collect = async { + while let Some(envelope) = stream.next().await { + if let temps_agent_runtime::retained::RuntimeEvent::ProviderEvent { event } = + envelope.event + { + observed.push(event); + } + } + }; + tokio::time::timeout(Duration::from_secs(10), collect) + .await + .expect("the compaction invocation must finish"); + completion.wait().await.unwrap(); + assert_eq!(transport.method_count("thread/compact/start"), 1); + assert_eq!(transport.spawn_count(), 1); + let compaction = compaction_events(&observed); + assert!( + matches!( + compaction.last(), + Some(TurnEvent::CompactionCompleted { .. }) + ), + "{compaction:?}" + ); +} + +#[tokio::test] +async fn compaction_instructions_are_refused_rather_than_dropped() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + assert!(!handle.driver_capabilities().compaction_instructions); + handle + .start_turn(TurnInput::new( + InvocationId::new("before").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + let mut input = + temps_agent_runtime::retained::CompactionInput::new(InvocationId::new("focused").unwrap()); + input.instructions = Some("keep the API notes".into()); + let Err(error) = handle.compact(input).await else { + panic!("Codex cannot take compaction instructions"); + }; + assert!( + error.to_string().contains("without summary instructions"), + "{error}" + ); + assert_eq!(transport.method_count("thread/compact/start"), 0); +} + #[tokio::test] async fn a_crashed_process_is_replaced_for_the_next_turn() { let transport = AppServer::new(Script::Crash);