use async_trait::async_trait;
use serde_json::{json, Value};
use std::io;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll};
use supercode_harness::acp_server::{self, AcpCompatibilityProfile, AcpServer, HttpAcpRuntime};
use supercode_harness::server::RpcEngine;
use supercode_harness::{
Agent, ApprovalPolicy, ChatMessage, ChatRequest, Config, FrontendActions,
FrontendAttachSnapshot, FrontendAttachment, FrontendConnectionState,
FrontendDisplayCapabilities, FrontendEvent, FrontendFacadeMethod, FrontendResponse,
FrontendRuntime, FrontendRuntimeDescriptor, FrontendRuntimeError, FrontendRuntimeMetadata,
FrontendTurnState, FunctionCall, HttpFrontendRuntime, Provider, Role, Session, SessionFormat,
ToolCall, Usage, FRONTEND_RUNTIME_SCHEMA_VERSION,
};
use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite, AsyncWriteExt, BufReader};
use tokio::sync::broadcast;
struct StreamingProvider;
struct BlockingProvider {
entered: tokio::sync::Notify,
release: tokio::sync::Notify,
}
struct SharedBlockingProvider(Arc<BlockingProvider>);
#[async_trait]
impl Provider for SharedBlockingProvider {
async fn complete(
&self,
_request: &ChatRequest,
on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
self.0.entered.notify_one();
self.0.release.notified().await;
on_delta("same reply");
Ok((ChatMessage::assistant("same reply"), Usage::default()))
}
}
#[async_trait]
impl Provider for StreamingProvider {
async fn complete(
&self,
_request: &ChatRequest,
on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
on_delta("hello over ACP");
Ok((ChatMessage::assistant("hello over ACP"), Usage::default()))
}
}
async fn write_request<W: AsyncWrite + Unpin>(writer: &mut W, value: Value) {
writer
.write_all(format!("{value}\n").as_bytes())
.await
.unwrap();
writer.flush().await.unwrap();
}
async fn read_message<R: AsyncBufRead + Unpin>(reader: &mut R) -> Value {
let mut line = String::new();
reader.read_line(&mut line).await.unwrap();
assert!(
!line.is_empty(),
"ACP server closed before sending a message"
);
serde_json::from_str(&line).unwrap()
}
#[tokio::test]
async fn acp_preserves_the_complete_sdk_script_contract() {
let fixture: Value = serde_json::from_str(include_str!(
"../../../sdk/frontend/test/fixtures/conformance.json"
))
.unwrap();
let session_id = fixture["session_id"].as_str().unwrap();
let prompt = fixture["prompt"].as_str().unwrap();
let competing_prompt = fixture["competing_prompt"].as_str().unwrap();
let reply = fixture["reply"].as_str().unwrap();
assert_eq!(reply, "same reply");
let provider = Arc::new(BlockingProvider {
entered: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let persisted = Arc::new(Mutex::new(String::new()));
let hook = persisted.clone();
let runtime = RpcEngine::new_named(
Agent::with_provider(
Config::builder().system_prompt("sdk conformance").build(),
Box::new(SharedBlockingProvider(provider.clone())),
),
session_id,
Some(Box::new(move |agent: &supercode_harness::SdkAgent| {
*hook.lock().unwrap() = serde_json::to_string(agent.history()).unwrap();
})),
);
let mut canonical = runtime.attach(50).await.unwrap();
let acp = AcpServer::new(runtime.clone()).await.unwrap();
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
acp,
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({"jsonrpc":"2.0", "id":1, "method":"initialize", "params":{"protocolVersion":1}}),
)
.await;
assert_eq!(
read_message(&mut client_read).await["result"]["protocolVersion"],
1
);
write_request(
&mut client_write,
json!({"jsonrpc":"2.0", "id":2, "method":"session/new", "params":{}}),
)
.await;
assert_eq!(
read_message(&mut client_read).await["result"]["sessionId"],
session_id
);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":3, "method":"session/prompt",
"params":{"sessionId":session_id, "prompt":[{"type":"text", "text":prompt}]}
,"_meta":{"surface":"ACP_ENVELOPE_SENTINEL"}
}),
)
.await;
provider.entered.notified().await;
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":4, "method":"session/prompt",
"params":{"sessionId":session_id, "prompt":[{"type":"text", "text":competing_prompt}]}
}),
)
.await;
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":5, "method":"session/respond",
"params":{"sessionId":session_id, "response":{
"kind":"other", "request_id":99, "action":"cancel", "content":null
}}
}),
)
.await;
let mut projected = Vec::new();
let mut saw_busy = false;
let mut saw_unsupported = false;
while !saw_busy || !saw_unsupported {
let message = read_message(&mut client_read).await;
if let (Some(sequence), Some(kind)) = (
message
.pointer("/params/update/_meta/supercode/sdkSequence")
.and_then(Value::as_u64),
message
.pointer("/params/update/_meta/supercode/sdkKind")
.and_then(Value::as_str),
) {
projected.push((sequence, kind.to_string()));
}
saw_busy |= message["id"] == 4
&& message["error"]["name"] == "busy"
&& message["error"]["operation"] == "input";
saw_unsupported |= message["id"] == 5
&& message["error"]["name"] == "unsupported_action"
&& message["error"]["operation"] == "respond";
}
provider.release.notify_waiters();
loop {
let message = read_message(&mut client_read).await;
if let (Some(sequence), Some(kind)) = (
message
.pointer("/params/update/_meta/supercode/sdkSequence")
.and_then(Value::as_u64),
message
.pointer("/params/update/_meta/supercode/sdkKind")
.and_then(Value::as_str),
) {
projected.push((sequence, kind.to_string()));
}
if message["id"] == 3 {
assert_eq!(message["result"]["stopReason"], "end_turn");
break;
}
}
let mut expected = Vec::new();
loop {
let event = canonical.next_event().await.unwrap();
let terminal = event.kind == "turn_succeeded";
expected.push((event.sequence, event.kind));
if terminal {
break;
}
}
assert_eq!(projected, expected);
assert_eq!(
expected
.iter()
.map(|(_, kind)| kind.as_str())
.collect::<Vec<_>>(),
fixture["event_kinds"]
.as_array()
.unwrap()
.iter()
.map(|value| value.as_str().unwrap())
.collect::<Vec<_>>()
);
let expected_persisted = serde_json::to_string(&vec![
ChatMessage::system("sdk conformance"),
ChatMessage::user(prompt),
ChatMessage::assistant(reply),
])
.unwrap();
assert_eq!(*persisted.lock().unwrap(), expected_persisted);
let content = Session::from_native_messages(runtime.attach(50).await.unwrap().history)
.to_jsonl(SessionFormat::Pi)
.unwrap();
let export_path =
std::env::temp_dir().join(format!("supercode-acp-export-{}.jsonl", std::process::id()));
std::fs::write(&export_path, &content).unwrap();
let reloaded = Session::load(&export_path).unwrap();
assert!(!content.contains("ACP_ENVELOPE_SENTINEL"));
assert!(!reloaded
.to_jsonl(SessionFormat::Pi)
.unwrap()
.contains("ACP_ENVELOPE_SENTINEL"));
std::fs::remove_file(export_path).ok();
client_write.shutdown().await.unwrap();
drop(client_write);
server_task.await.unwrap().unwrap();
runtime.shutdown().await;
}
#[tokio::test]
async fn acp_negotiates_opens_prompts_and_streams_from_the_canonical_runtime() {
let cwd = std::env::temp_dir().join(format!("supercode-acp-server-{}", std::process::id()));
std::fs::create_dir_all(&cwd).unwrap();
let agent = Agent::with_provider(
Config::builder().cwd(cwd.clone()).build(),
Box::new(StreamingProvider),
);
let runtime = RpcEngine::new_named(agent, "shared-runtime", None);
let acp = AcpServer::new(runtime.clone()).await.unwrap();
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
acp,
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {"protocolVersion": 1, "clientCapabilities": {}}
}),
)
.await;
let initialized = read_message(&mut client_read).await;
assert_eq!(initialized["result"]["protocolVersion"], 1);
assert_eq!(
initialized["result"]
.pointer("/agentCapabilities/_meta/supercode/frontend/eventMethod")
.and_then(Value::as_str),
Some("frontend.v2.event")
);
assert_eq!(
initialized["result"]
.pointer("/agentCapabilities/_meta/supercode/frontend/contract")
.and_then(Value::as_str),
Some("supercode.frontend.contract.v2")
);
let advertised = initialized["result"]
.pointer("/agentCapabilities/_meta/supercode/frontend/methods")
.and_then(Value::as_array)
.unwrap();
for method in FrontendFacadeMethod::ALL {
assert!(advertised.iter().any(|value| value == method.wire_name()));
}
assert_eq!(
initialized["result"]["agentCapabilities"]["loadSession"],
true
);
assert_eq!(
initialized["result"]["agentCapabilities"]["sessionCapabilities"]["resume"],
json!({})
);
write_request(
&mut client_write,
json!({
"jsonrpc": "2.0", "id": 2, "method": "session/new",
"params": {"cwd": "/tmp", "mcpServers": []}
}),
)
.await;
let opened = read_message(&mut client_read).await;
assert_eq!(opened["result"]["sessionId"], "shared-runtime");
assert_eq!(runtime.status().session_id, "shared-runtime");
write_request(
&mut client_write,
json!({
"jsonrpc": "2.0", "id": 3, "method": "session/prompt",
"params": {
"sessionId": "shared-runtime",
"prompt": [{"type": "text", "text": "say hello"}],
"_meta": {"surface": "ACP_ENVELOPE_SENTINEL"}
}
}),
)
.await;
let mut messages = Vec::new();
loop {
let message = read_message(&mut client_read).await;
let terminal = message.get("id").and_then(Value::as_u64) == Some(3);
messages.push(message);
if terminal {
break;
}
}
assert!(messages.iter().any(|message| {
message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
== Some("agent_message_chunk")
&& message
.pointer("/params/update/content/text")
.and_then(Value::as_str)
== Some("hello over ACP")
}));
assert!(messages.iter().any(|message| {
message.get("id").and_then(Value::as_u64) == Some(3)
&& message
.pointer("/result/stopReason")
.and_then(Value::as_str)
== Some("end_turn")
}));
assert!(!serde_json::to_string(&runtime.history(50).await)
.unwrap()
.contains("ACP_ENVELOPE_SENTINEL"));
client_write.shutdown().await.unwrap();
drop(client_write);
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.expect("ACP server should stop when its client stream closes")
.unwrap()
.unwrap();
std::fs::remove_dir_all(cwd).ok();
}
struct ScriptedRuntime {
events: broadcast::Sender<FrontendEvent>,
next_sequence: std::sync::atomic::AtomicU64,
submits: std::sync::atomic::AtomicUsize,
reply: String,
emitted: Vec<Value>,
history: Vec<ChatMessage>,
replay: std::collections::VecDeque<FrontendEvent>,
}
impl ScriptedRuntime {
fn new(reply: &str, emitted: Vec<Value>) -> Self {
let (events, _) = broadcast::channel(32);
Self {
events,
next_sequence: std::sync::atomic::AtomicU64::new(1),
submits: std::sync::atomic::AtomicUsize::new(0),
reply: reply.into(),
emitted,
history: Vec::new(),
replay: std::collections::VecDeque::new(),
}
}
fn with_history(mut self, history: Vec<ChatMessage>) -> Self {
self.history = history;
self.next_sequence
.store(78, std::sync::atomic::Ordering::SeqCst);
self
}
fn with_replay(mut self, replay: Vec<FrontendEvent>) -> Self {
if let Some(maximum) = replay.iter().map(|event| event.sequence).max() {
self.next_sequence
.store(maximum + 1, std::sync::atomic::Ordering::SeqCst);
}
self.replay = replay.into();
self
}
}
#[async_trait]
impl FrontendRuntime for ScriptedRuntime {
async fn describe(&self) -> Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
Ok(scripted_descriptor())
}
async fn attach(
&self,
_history_limit: usize,
) -> Result<FrontendAttachment, FrontendRuntimeError> {
Ok(FrontendAttachment::from_snapshot(
FrontendAttachSnapshot {
descriptor: scripted_descriptor(),
history: self.history.clone(),
history_cursor: u64::from(!self.history.is_empty()) * 77,
replay: self.replay.clone(),
},
self.events.subscribe(),
))
}
async fn send_input(self: Arc<Self>, prompt: String) -> Result<(), FrontendRuntimeError> {
self.submit(prompt).await.map(|_| ())
}
async fn submit(&self, _prompt: String) -> Result<String, FrontendRuntimeError> {
self.submits
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
for payload in &self.emitted {
let sequence = self
.next_sequence
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let kind = payload
.get("type")
.and_then(Value::as_str)
.unwrap_or("unknown")
.to_string();
let _ = self.events.send(FrontendEvent {
sequence,
kind,
payload: payload.clone(),
});
}
let sequence = self
.next_sequence
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let _ = self.events.send(FrontendEvent {
sequence,
kind: "turn_succeeded".into(),
payload: json!({"type":"turn_succeeded", "reply":self.reply}),
});
Ok(self.reply.clone())
}
async fn interrupt(&self) -> Result<bool, FrontendRuntimeError> {
Ok(false)
}
async fn steer(&self, _prompt: String) -> Result<(), FrontendRuntimeError> {
Err(FrontendRuntimeError::UnsupportedAction("steer"))
}
async fn respond(&self, _response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
Err(FrontendRuntimeError::UnsupportedAction("respond"))
}
}
fn scripted_descriptor() -> FrontendRuntimeDescriptor {
FrontendRuntimeDescriptor {
schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
session_id: "scripted-runtime".into(),
source_harness: None,
emulation_profile: None,
active_modules: Vec::new(),
commands: Vec::new(),
operations: Vec::new(),
actions: FrontendActions {
submit: true,
interrupt: true,
steer: false,
respond: false,
detach: true,
close: false,
},
display: FrontendDisplayCapabilities {
event_kinds: vec!["text_delta".into(), "runtime_disconnected".into()],
opaque_fallback: true,
},
model: "scripted".into(),
turn_state: FrontendTurnState::Idle,
connection_state: FrontendConnectionState::Connected,
extensions: Default::default(),
}
}
struct GooseApprovalProvider(std::sync::atomic::AtomicUsize);
#[async_trait]
impl Provider for GooseApprovalProvider {
async fn complete(
&self,
request: &ChatRequest,
on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
if self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst) == 0 {
return Ok((
ChatMessage {
role: Role::Assistant,
content: None,
content_parts: None,
tool_calls: Some(vec![ToolCall {
id: "goose-call-1".into(),
kind: "function".into(),
function: FunctionCall {
name: "bash".into(),
arguments: json!({"command":"echo goose-approved"}).to_string(),
},
}]),
tool_call_id: None,
name: None,
metadata: Default::default(),
},
Usage::default(),
));
}
let tool_result = request
.messages
.last()
.expect("tool result reaches provider");
assert_eq!(tool_result.role, Role::Tool);
assert!(tool_result
.content
.as_deref()
.unwrap_or_default()
.contains("goose-approved"));
on_delta("permission handled");
Ok((
ChatMessage::assistant("permission handled"),
Usage::default(),
))
}
}
struct WriterFailureProvider(std::sync::atomic::AtomicUsize);
#[async_trait]
impl Provider for WriterFailureProvider {
async fn complete(
&self,
request: &ChatRequest,
on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
if self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst) == 0 {
return Ok((
ChatMessage {
role: Role::Assistant,
content: None,
content_parts: None,
tool_calls: Some(vec![ToolCall {
id: "writer-failure-call".into(),
kind: "function".into(),
function: FunctionCall {
name: "bash".into(),
arguments: json!({"command":"printf MUST_NOT_RUN"}).to_string(),
},
}]),
tool_call_id: None,
name: None,
metadata: Default::default(),
},
Usage::default(),
));
}
let tool_result = request.messages.last().expect("denied tool result");
assert_eq!(tool_result.role, Role::Tool);
assert!(tool_result
.content
.as_deref()
.unwrap_or_default()
.contains("was not approved for execution"));
on_delta("writer failure denied");
Ok((
ChatMessage::assistant("writer failure denied"),
Usage::default(),
))
}
}
struct FailOnPermissionWriter {
seen: Vec<u8>,
failed: Arc<std::sync::atomic::AtomicBool>,
}
impl AsyncWrite for FailOnPermissionWriter {
fn poll_write(
mut self: Pin<&mut Self>,
_cx: &mut Context<'_>,
bytes: &[u8],
) -> Poll<io::Result<usize>> {
self.seen.extend_from_slice(bytes);
if self
.seen
.windows(b"session/request_permission".len())
.any(|window| window == b"session/request_permission")
{
self.failed.store(true, std::sync::atomic::Ordering::SeqCst);
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"injected ACP permission writer failure",
)));
}
Poll::Ready(Ok(bytes.len()))
}
fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<io::Result<()>> {
Poll::Ready(Ok(()))
}
fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<io::Result<()>> {
Poll::Ready(Ok(()))
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn pinned_goose_permission_response_unblocks_the_canonical_tool_turn() {
let root = std::env::temp_dir().join(format!(
"supercode-goose-acp-permission-{}",
std::process::id()
));
std::fs::remove_dir_all(&root).ok();
std::fs::create_dir_all(&root).unwrap();
let mut config = Config::builder()
.cwd(root.clone())
.approval(ApprovalPolicy::OnRequest)
.build();
config.permissions_enabled = true;
let runtime = RpcEngine::new_named_with_frontend_requests(
Agent::with_provider(config, Box::new(GooseApprovalProvider(0.into()))),
"goose-permission-runtime",
FrontendRuntimeMetadata::default(),
None,
);
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
AcpServer::new_with_profile(runtime.clone(), AcpCompatibilityProfile::GooseAcd3c135)
.await
.unwrap(),
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({"jsonrpc":"2.0", "id":1, "method":"initialize", "params":{"protocolVersion":1}}),
)
.await;
assert_eq!(read_message(&mut client_read).await["id"], 1);
write_request(
&mut client_write,
json!({"jsonrpc":"2.0", "id":2, "method":"session/new", "params":{"mcpServers":[]}}),
)
.await;
assert_eq!(read_message(&mut client_read).await["id"], 2);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":3, "method":"session/prompt",
"params":{
"sessionId":"goose-permission-runtime",
"prompt":[{"type":"text", "text":"run the reviewed command"}]
}
}),
)
.await;
let mut saw_tool_start = false;
let permission = loop {
let message = tokio::time::timeout(
std::time::Duration::from_secs(5),
read_message(&mut client_read),
)
.await
.expect("Goose permission exchange timed out");
saw_tool_start |= message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
== Some("tool_call");
if message.get("method").and_then(Value::as_str) == Some("session/request_permission") {
assert_eq!(message["params"]["sessionId"], "goose-permission-runtime");
assert_eq!(
message["params"]["toolCall"]["_meta"]["supercode"]["tool"],
"bash"
);
assert_eq!(message["params"]["options"][0]["optionId"], "allow_once");
break message["id"].clone();
}
};
assert!(saw_tool_start);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0",
"id":permission,
"result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}}
}),
)
.await;
let mut saw_tool_completion = false;
loop {
let message = tokio::time::timeout(
std::time::Duration::from_secs(5),
read_message(&mut client_read),
)
.await
.expect("approved Goose turn timed out");
saw_tool_completion |= message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
== Some("tool_call_update");
if message.get("id") == Some(&json!(3)) {
assert_eq!(message["result"]["stopReason"], "end_turn");
break;
}
}
assert!(saw_tool_completion);
assert_eq!(
runtime.history(20).await.last().unwrap().content.as_deref(),
Some("permission handled")
);
client_write.shutdown().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.unwrap()
.unwrap()
.unwrap();
runtime.shutdown().await;
std::fs::remove_dir_all(root).ok();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn permission_writer_failure_unblocks_held_input_and_denies_the_pending_tool() {
let root = std::env::temp_dir().join(format!(
"supercode-acp-writer-failure-{}",
std::process::id()
));
std::fs::remove_dir_all(&root).ok();
std::fs::create_dir_all(&root).unwrap();
let mut config = Config::builder()
.cwd(root.clone())
.approval(ApprovalPolicy::OnRequest)
.build();
config.permissions_enabled = true;
let runtime = RpcEngine::new_named_with_frontend_requests(
Agent::with_provider(config, Box::new(WriterFailureProvider(0.into()))),
"writer-failure-runtime",
FrontendRuntimeMetadata::default(),
None,
);
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, _server_write) = tokio::io::split(server_io);
let (_client_read, mut client_write) = tokio::io::split(client_io);
let failed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let server_task = tokio::spawn(acp_server::run_stdio(
AcpServer::new_with_profile(runtime.clone(), AcpCompatibilityProfile::GooseAcd3c135)
.await
.unwrap(),
BufReader::new(server_read),
FailOnPermissionWriter {
seen: Vec::new(),
failed: failed.clone(),
},
));
for request in [
json!({"jsonrpc":"2.0", "id":1, "method":"initialize", "params":{"protocolVersion":1}}),
json!({"jsonrpc":"2.0", "id":2, "method":"session/new", "params":{"mcpServers":[]}}),
json!({
"jsonrpc":"2.0", "id":3, "method":"session/prompt",
"params":{
"sessionId":"writer-failure-runtime",
"prompt":[{"type":"text", "text":"exercise writer failure"}]
}
}),
] {
write_request(&mut client_write, request).await;
}
let error = tokio::time::timeout(std::time::Duration::from_secs(5), server_task)
.await
.expect("run_stdio hung with its input peer held open")
.unwrap()
.expect_err("writer failure must reach the ACP caller");
assert_eq!(error.kind(), io::ErrorKind::BrokenPipe);
assert!(failed.load(std::sync::atomic::Ordering::SeqCst));
tokio::time::timeout(std::time::Duration::from_secs(2), async {
while runtime.status().busy {
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await
.expect("pending approval did not fail closed after writer loss");
assert_eq!(
runtime.history(20).await.last().unwrap().content.as_deref(),
Some("writer failure denied")
);
drop(client_write);
runtime.shutdown().await;
std::fs::remove_dir_all(root).ok();
}
#[tokio::test]
async fn pinned_goose_profile_maps_new_session_to_one_runtime_and_replays_history() {
let mut assistant_with_tool = ChatMessage::assistant("I will inspect it.");
assistant_with_tool.tool_calls = Some(vec![ToolCall {
id: "call-read".into(),
kind: "function".into(),
function: FunctionCall {
name: "read_file".into(),
arguments: r#"{"path":"README.md"}"#.into(),
},
}]);
let history = vec![
ChatMessage::system("private runtime authority"),
ChatMessage::user_with_images(
"Prior user question",
&["data:image/png;base64,aGVsbG8=".into()],
),
assistant_with_tool,
ChatMessage::tool_result("call-read", "read_file", "file contents"),
ChatMessage::assistant("Prior final answer"),
];
let runtime = Arc::new(
ScriptedRuntime::new("fresh reply", Vec::new())
.with_history(history)
.with_replay(vec![FrontendEvent {
sequence: 78,
kind: "text_delta".into(),
payload: json!({
"type":"text_delta",
"text":"reply completed between snapshot and session/new"
}),
}]),
);
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
AcpServer::new_with_profile(runtime.clone(), AcpCompatibilityProfile::GooseAcd3c135)
.await
.unwrap(),
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":0, "method":"initialize",
"params":{
"protocolVersion":1,
"clientInfo":{"name":"goose-text","version":"0.1.0"},
"clientCapabilities":{}
}
}),
)
.await;
let initialized = read_message(&mut client_read).await;
assert_eq!(initialized["result"]["authMethods"], json!([]));
assert_eq!(initialized["result"]["agentInfo"]["name"], "supercode");
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":1,
"method":"_goose/unstable/defaults/read", "params":{}
}),
)
.await;
let defaults = read_message(&mut client_read).await;
assert_eq!(defaults["result"]["providerId"], "supercode");
assert_eq!(defaults["result"]["modelId"], "scripted");
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":2, "method":"session/new",
"params":{"cwd":"/remote/client/path", "mcpServers":[]}
}),
)
.await;
let mut updates = Vec::new();
loop {
let message = read_message(&mut client_read).await;
if message.get("id") == Some(&json!(2)) {
assert_eq!(message["result"]["sessionId"], "scripted-runtime");
break;
}
updates.push(message);
}
assert_eq!(
updates
.iter()
.filter_map(|message| {
message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
})
.collect::<Vec<_>>(),
[
"user_message_chunk",
"user_message_chunk",
"agent_message_chunk",
"tool_call",
"tool_call_update",
"agent_message_chunk",
"agent_message_chunk",
]
);
assert_eq!(
updates[0]
.pointer("/params/update/_meta/supercode/historyCursor")
.and_then(Value::as_u64),
Some(77)
);
assert_eq!(
updates[1]
.pointer("/params/update/content/mimeType")
.and_then(Value::as_str),
Some("image/png")
);
assert_eq!(
updates[3]
.pointer("/params/update/rawInput/path")
.and_then(Value::as_str),
Some("README.md")
);
assert_eq!(
updates[4]
.pointer("/params/update/toolCallId")
.and_then(Value::as_str),
Some("call-read")
);
assert!(
updates.iter().any(|message| {
message
.pointer("/params/update/content/text")
.and_then(Value::as_str)
== Some("reply completed between snapshot and session/new")
}),
"attachment replay newer than the history cursor must cross the open boundary"
);
assert_eq!(
runtime.submits.load(std::sync::atomic::Ordering::SeqCst),
0,
"session/new must attach a display, not create a continuation"
);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":3,
"method":"_goose/unstable/session/extensions/add",
"params":{"sessionId":"scripted-runtime","name":"forbidden"}
}),
)
.await;
let unsupported = read_message(&mut client_read).await;
assert_eq!(unsupported["error"]["code"], -32601);
assert_eq!(runtime.submits.load(std::sync::atomic::Ordering::SeqCst), 0);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":4, "method":"session/new",
"params":{"cwd":"/remote/client/path", "mcpServers":[]}
}),
)
.await;
let reopened = read_message(&mut client_read).await;
assert_eq!(reopened["id"], 4);
assert_eq!(reopened["result"]["sessionId"], "scripted-runtime");
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":5, "method":"session/prompt",
"params":{
"sessionId":"scripted-runtime",
"prompt":[{"type":"text", "text":"continue"}]
}
}),
)
.await;
let mut live = Vec::new();
loop {
let message = read_message(&mut client_read).await;
let terminal = message.get("id") == Some(&json!(5));
live.push(message);
if terminal {
break;
}
}
assert_eq!(
live.iter()
.filter(|message| {
message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
== Some("agent_message_chunk")
})
.filter_map(|message| {
message
.pointer("/params/update/content/text")
.and_then(Value::as_str)
})
.collect::<Vec<_>>(),
["fresh reply"],
"reopening must not duplicate history and the live reply must appear once; wire={live:?}"
);
assert_eq!(live.last().unwrap()["result"]["stopReason"], "end_turn");
assert_eq!(runtime.submits.load(std::sync::atomic::Ordering::SeqCst), 1);
client_write.shutdown().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.unwrap()
.unwrap()
.unwrap();
}
#[tokio::test]
async fn goose_profile_rejects_runtime_configuration_and_standard_acp_stays_neutral() {
let runtime = Arc::new(
ScriptedRuntime::new("unused", Vec::new())
.with_history(vec![ChatMessage::assistant("history must stay hidden")]),
);
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
AcpServer::new(runtime.clone()).await.unwrap(),
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":1,
"method":"_goose/unstable/defaults/read", "params":{}
}),
)
.await;
assert_eq!(
read_message(&mut client_read).await["error"]["code"],
-32601,
"standard ACP must treat donor extensions exactly like any unknown method"
);
write_request(
&mut client_write,
json!({"jsonrpc":"2.0", "id":2, "method":"session/new", "params":{}}),
)
.await;
let opened = read_message(&mut client_read).await;
assert_eq!(
opened["id"], 2,
"standard ACP must not receive Goose history replay"
);
assert_eq!(runtime.submits.load(std::sync::atomic::Ordering::SeqCst), 0);
client_write.shutdown().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.unwrap()
.unwrap()
.unwrap();
let runtime = Arc::new(ScriptedRuntime::new("unused", Vec::new()));
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
AcpServer::new_with_profile(runtime.clone(), AcpCompatibilityProfile::GooseAcd3c135)
.await
.unwrap(),
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0", "id":3, "method":"session/new",
"params":{"mcpServers":[{"name":"forbidden"}]}
}),
)
.await;
let rejected = read_message(&mut client_read).await;
assert_eq!(rejected["error"]["code"], -32000);
assert!(rejected["error"]["message"]
.as_str()
.unwrap()
.contains("cannot mutate"));
assert_eq!(runtime.submits.load(std::sync::atomic::Ordering::SeqCst), 0);
client_write.shutdown().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.unwrap()
.unwrap()
.unwrap();
}
async fn scripted_prompt(reply: &str, emitted: Vec<Value>) -> Vec<Value> {
let runtime = std::sync::Arc::new(ScriptedRuntime::new(reply, emitted));
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
AcpServer::new(runtime).await.unwrap(),
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({"jsonrpc":"2.0","id":1,"method":"session/new","params":{}}),
)
.await;
let _ = read_message(&mut client_read).await;
write_request(
&mut client_write,
json!({
"jsonrpc":"2.0","id":3,"method":"session/prompt",
"params": {
"sessionId":"scripted-runtime",
"prompt":[{"type":"text","text":"go"}]
}
}),
)
.await;
let mut messages = Vec::new();
loop {
let message = tokio::time::timeout(
std::time::Duration::from_secs(2),
read_message(&mut client_read),
)
.await
.expect("scripted ACP response timed out");
let terminal = message.get("id").and_then(Value::as_u64) == Some(3);
messages.push(message);
if terminal {
break;
}
}
client_write.shutdown().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.unwrap()
.unwrap()
.unwrap();
messages
}
#[tokio::test]
async fn acp_projects_rpc_reply_when_runtime_emits_no_text_deltas() {
let messages = scripted_prompt("reply-only", vec![]).await;
assert_eq!(
messages
.iter()
.filter(|message| {
message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
== Some("agent_message_chunk")
})
.filter_map(|message| message.pointer("/params/update/content/text"))
.collect::<Vec<_>>(),
vec![&json!("reply-only")]
);
assert_eq!(messages.last().unwrap()["result"]["stopReason"], "end_turn");
}
#[tokio::test]
async fn acp_does_not_duplicate_streamed_reply_and_rejects_mismatches() {
let messages = scripted_prompt(
"streamed",
vec![json!({"type":"text_delta","text":"streamed"})],
)
.await;
let chunks = messages
.iter()
.filter(|message| {
message
.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str)
== Some("agent_message_chunk")
})
.filter_map(|message| {
message
.pointer("/params/update/content/text")
.and_then(Value::as_str)
})
.collect::<Vec<_>>();
assert_eq!(chunks, vec!["streamed"]);
assert_eq!(
messages
.iter()
.filter_map(|message| {
message
.pointer("/params/update/_meta/supercode/sdkSequence")
.and_then(Value::as_u64)
})
.collect::<Vec<_>>(),
vec![1, 2]
);
assert_eq!(messages.last().unwrap()["result"]["stopReason"], "end_turn");
let mismatch = scripted_prompt(
"different",
vec![json!({"type":"text_delta","text":"streamed"})],
)
.await;
assert!(mismatch.last().unwrap().get("error").is_some());
assert!(mismatch.last().unwrap().get("result").is_none());
}
#[tokio::test]
async fn acp_frontend_events_preserve_unknown_payloads_opaquely() {
let unknown = json!({
"type":"future_semantic_event",
"extension":{"nested":[1, {"untouched":true}]},
"_meta":{"newPeer":{"version":9}}
});
let messages = scripted_prompt("reply", vec![unknown.clone()]).await;
let projected = messages
.iter()
.find(|message| {
message.get("method").and_then(Value::as_str) == Some("frontend.v2.event")
&& message
.pointer("/params/event/kind")
.and_then(Value::as_str)
== Some("future_semantic_event")
})
.expect("canonical ACP frontend event");
assert_eq!(projected.pointer("/params/event/payload"), Some(&unknown));
}
#[tokio::test]
async fn acp_disconnect_never_reports_end_turn() {
let messages = scripted_prompt(
"reply",
vec![json!({
"type":"runtime_disconnected",
"message":"event transport vanished"
})],
)
.await;
let terminal = messages.last().unwrap();
assert!(terminal.get("error").is_some());
assert!(messages.iter().all(|message| {
message
.pointer("/result/stopReason")
.and_then(Value::as_str)
!= Some("end_turn")
}));
}
#[tokio::test]
async fn acp_bridge_joins_an_existing_http_runtime_instead_of_resuming_a_copy() {
let cwd = std::env::temp_dir().join(format!("supercode-acp-http-{}", std::process::id()));
std::fs::create_dir_all(&cwd).unwrap();
let agent = Agent::with_provider(
Config::builder().cwd(cwd.clone()).build(),
Box::new(StreamingProvider),
);
let persisted = Arc::new(Mutex::new(String::new()));
let persisted_hook = persisted.clone();
let runtime = RpcEngine::new_named(
agent,
"live-http-runtime",
Some(Box::new(move |agent: &supercode_harness::SdkAgent| {
*persisted_hook.lock().unwrap() = serde_json::to_string(agent.history()).unwrap();
})),
);
let token = std::sync::Arc::<str>::from("test-token");
let address =
supercode_harness::server::run_http(runtime.clone(), "127.0.0.1:0", token.clone())
.await
.unwrap();
let remote = HttpAcpRuntime::connect(format!("http://{address}"), token.to_string())
.await
.unwrap();
assert_eq!(
remote.describe().await.unwrap().session_id,
"live-http-runtime"
);
let acp = AcpServer::new(remote).await.unwrap();
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let (server_read, server_write) = tokio::io::split(server_io);
let (client_read, mut client_write) = tokio::io::split(client_io);
let server_task = tokio::spawn(acp_server::run_stdio(
acp,
BufReader::new(server_read),
server_write,
));
let mut client_read = BufReader::new(client_read);
write_request(
&mut client_write,
json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {"protocolVersion": 1, "clientCapabilities": {}}
}),
)
.await;
assert_eq!(
read_message(&mut client_read).await["result"]["protocolVersion"],
1
);
write_request(
&mut client_write,
json!({
"jsonrpc": "2.0", "id": 2, "method": "session/load",
"params": {"sessionId": "live-http-runtime", "cwd": "/tmp", "mcpServers": []}
}),
)
.await;
assert_eq!(
read_message(&mut client_read).await["result"]["sessionId"],
"live-http-runtime"
);
write_request(
&mut client_write,
json!({
"jsonrpc": "2.0", "id": 3, "method": "session/prompt",
"params": {
"sessionId": "live-http-runtime",
"prompt": [{"type": "text", "text": "through remote ACP"}]
}
}),
)
.await;
let mut saw_update = false;
let mut saw_completion = false;
loop {
let message = read_message(&mut client_read).await;
saw_update |= message
.pointer("/params/update/content/text")
.and_then(Value::as_str)
== Some("hello over ACP");
saw_completion |= message.get("id").and_then(Value::as_u64) == Some(3)
&& message
.pointer("/result/stopReason")
.and_then(Value::as_str)
== Some("end_turn");
if message.get("id").and_then(Value::as_u64) == Some(3) {
break;
}
}
assert!(saw_update && saw_completion);
assert!(!runtime.status().shutting_down);
client_write.shutdown().await.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), server_task)
.await
.unwrap()
.unwrap()
.unwrap();
assert!(!runtime.status().shutting_down);
let local_history = runtime.attach(50).await.unwrap().history;
let http = HttpFrontendRuntime::connect(format!("http://{address}"), token.to_string())
.await
.unwrap();
let http_history = http.attach(50).await.unwrap().history;
http.detach().await.unwrap();
let second_remote = HttpAcpRuntime::connect(format!("http://{address}"), token.to_string())
.await
.unwrap();
let second_acp = AcpServer::new(second_remote).await.unwrap();
let (second_server_io, second_client_io) = tokio::io::duplex(64 * 1024);
let (second_server_read, second_server_write) = tokio::io::split(second_server_io);
let (second_client_read, mut second_client_write) = tokio::io::split(second_client_io);
let second_server_task = tokio::spawn(acp_server::run_stdio(
second_acp,
BufReader::new(second_server_read),
second_server_write,
));
let mut second_client_read = BufReader::new(second_client_read);
write_request(
&mut second_client_write,
json!({"jsonrpc":"2.0","id":10,"method":"initialize","params":{"protocolVersion":1}}),
)
.await;
let _ = read_message(&mut second_client_read).await;
write_request(
&mut second_client_write,
json!({"jsonrpc":"2.0","id":11,"method":"session/load","params":{"sessionId":"live-http-runtime"}}),
)
.await;
let _ = read_message(&mut second_client_read).await;
write_request(
&mut second_client_write,
json!({
"jsonrpc":"2.0","id":12,"method":"frontend.v2.attach",
"params":{"sessionId":"live-http-runtime","limit":50,"after_sequence":0}
}),
)
.await;
let attach_response = loop {
let message = read_message(&mut second_client_read).await;
if message["id"] == 12 {
break message;
}
assert!(
matches!(
message["method"].as_str(),
Some("frontend.v2.event" | "session/update")
),
"only sequenced replay may precede the correlated attach response: {message}"
);
};
let acp_history: Vec<ChatMessage> =
serde_json::from_value(attach_response["result"]["history"].clone()).unwrap();
second_client_write.shutdown().await.unwrap();
second_server_task.await.unwrap().unwrap();
let persisted_local_history = serde_json::to_string(&local_history).unwrap();
assert_eq!(
persisted_local_history,
serde_json::to_string(&http_history).unwrap()
);
assert_eq!(
persisted_local_history,
serde_json::to_string(&acp_history).unwrap()
);
assert_eq!(
persisted_local_history,
*persisted.lock().unwrap(),
"local, HTTP, and ACP reattachments must converge on the persisted transcript"
);
runtime.shutdown().await;
std::fs::remove_dir_all(cwd).ok();
}