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
25 changes: 25 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
Expand Down
1 change: 1 addition & 0 deletions examples/fullstack-chat/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

7 changes: 7 additions & 0 deletions src/adapter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,10 @@ pub struct CommandSpec {
pub initial_stdin: Option<Vec<u8>>,
/// 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<u16>,
}

impl fmt::Debug for CommandSpec {
Expand All @@ -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()
}
}
Expand Down Expand Up @@ -91,6 +96,7 @@ impl CommandSpec {
clear_environment: true,
initial_stdin: None,
interactive_stdin: false,
loopback_ports: Vec::new(),
}
}

Expand All @@ -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,
}
}
}
Expand Down
10 changes: 4 additions & 6 deletions src/providers/claude.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
}
}
Expand Down Expand Up @@ -1991,12 +1994,7 @@ fn claude_account_usage(value: &Value) -> Option<AccountUsageSnapshot> {
}

/// 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;
Expand Down
4 changes: 4 additions & 0 deletions src/providers/codex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}
Expand Down
200 changes: 196 additions & 4 deletions src/providers/codex_app_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,9 @@ struct TurnState {
/// Latest active-context occupancy, the pre-compaction size when a
/// compaction starts.
last_context_tokens: Option<u64>,
/// This invocation compacts the thread (`thread/compact/start`) instead of
/// starting a turn.
manual_compaction: bool,
}

pub(super) fn mark_retained(state: &mut AdapterState) {
Expand Down Expand Up @@ -154,6 +157,7 @@ pub(super) fn handshake() -> Vec<u8> {
/// 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"));
Expand Down Expand Up @@ -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()
},
);
Expand Down Expand Up @@ -296,6 +301,31 @@ pub(super) fn parse_line(line: &str, state: &mut AdapterState) -> Result<Adapter
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
// `thread/compact/start` answers with an empty result; the compaction's
// turn id first appears on `turn/started`, so a manual compaction adopts
// it there instead of discarding its own turn as uncorrelated.
if turn.manual_compaction
&& turn.turn_id.is_none()
&& method == "turn/started"
&& turn.thread_id.is_some()
&& value.pointer("/params/threadId").and_then(Value::as_str) == turn.thread_id.as_deref()
{
if let Some(started) = reported_turn_id(&value) {
turn.turn_id = Some(started.to_string());
store(state, &turn);
}
}
// Older app servers report a finished compaction only through
// `thread/compacted`, which names the thread but no turn.
let own_thread_compacted = turn.manual_compaction
&& method == "thread/compacted"
&& turn.thread_id.is_some()
&& value.pointer("/params/threadId").and_then(Value::as_str) == turn.thread_id.as_deref();
if own_thread_compacted {
notification(&value, &method, &mut turn, state, &mut output);
store(state, &turn);
return Ok(output);
}
if turn.retained
&& !method.is_empty()
&& value.get("id").is_none()
Expand Down Expand Up @@ -442,6 +472,18 @@ fn parse_response(
});
}
turn.thread_id = Some(thread_id.clone());
if turn.manual_compaction {
turn.open_compaction = Some("manual".to_string());
output.events.push(TurnEvent::CompactionStarted {
trigger: CompactionTrigger::Manual,
});
output.writes.push(encode(&request(
ID_TURN,
"thread/compact/start",
json!({"threadId": thread_id}),
))?);
return Ok(());
Comment thread
greptile-apps[bot] marked this conversation as resolved.
}
let mut params = turn.turn_params.clone();
params["threadId"] = json!(thread_id);
output
Expand Down Expand Up @@ -638,6 +680,18 @@ fn notification(
output.events.push(TurnEvent::Usage(usage));
}
}
"thread/compacted" if turn.manual_compaction => {
// 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
Expand Down Expand Up @@ -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()
})),
Expand Down Expand Up @@ -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),
Comment thread
greptile-apps[bot] marked this conversation as resolved.
});
}

fn completed_compaction(pre_tokens: Option<u64>) -> ContextCompaction {
completed_compaction_with(CompactionTrigger::Automatic, pre_tokens)
}

fn completed_compaction_with(
trigger: CompactionTrigger,
pre_tokens: Option<u64>,
) -> ContextCompaction {
ContextCompaction {
trigger: CompactionTrigger::Automatic,
trigger,
pre_tokens,
post_tokens: None,
dropped_tokens: None,
Expand Down Expand Up @@ -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");
Expand Down
Loading
Loading