use std::collections::{HashSet, VecDeque};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use serde_json::{json, Value};
use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite, AsyncWriteExt};
use tokio::sync::{broadcast, mpsc};
use crate::server::{RuntimeSubmitError, SERVER_MAX_LINE_BYTES};
use crate::{
ChatMessage, FrontendApprovalDecision, FrontendEvent, FrontendRequest, FrontendRequestKind,
FrontendResponse, FrontendRuntimeError, Role, SdkRuntime, FRONTEND_REPLAY_CAPACITY,
};
const ACP_APPROVAL_REQUEST_ID_PREFIX: &str = "supercode-approval-";
pub const ACP_PROTOCOL_VERSION: u64 = 1;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum AcpCompatibilityProfile {
#[default]
Standard,
GooseAcd3c135,
}
#[cfg(feature = "adapter-api")]
pub type HttpAcpRuntime = crate::HttpFrontendRuntime;
struct AcpRuntimeBridge {
runtime: Arc<dyn SdkRuntime>,
session_id: String,
model: String,
history: Vec<ChatMessage>,
history_cursor: u64,
initial_replay_cursor: u64,
events: broadcast::Sender<FrontendEvent>,
routing: Mutex<EventRoutingState>,
}
#[derive(Default)]
struct EventRoutingState {
active: bool,
pending: VecDeque<FrontendEvent>,
failure: Option<String>,
}
impl EventRoutingState {
fn buffer(&mut self, event: FrontendEvent) -> Result<(), ()> {
if self.pending.len() >= FRONTEND_REPLAY_CAPACITY {
self.pending.clear();
self.failure = Some(format!(
"ACP attachment received more than {FRONTEND_REPLAY_CAPACITY} events before session activation; reconnect required to preserve a gap-free stream"
));
return Err(());
}
self.pending.push_back(event);
Ok(())
}
fn activate(&mut self, acknowledged_cursor: u64) -> Result<Vec<FrontendEvent>, String> {
if let Some(error) = &self.failure {
return Err(error.clone());
}
self.active = true;
Ok(self
.pending
.drain(..)
.filter(|event| event.sequence > acknowledged_cursor)
.collect())
}
}
impl AcpRuntimeBridge {
async fn connect(runtime: Arc<dyn SdkRuntime>) -> Result<Arc<Self>, FrontendRuntimeError> {
let mut attachment = runtime.attach(200).await?;
let session_id = attachment.descriptor.session_id.clone();
let model = attachment.descriptor.model.clone();
let history = std::mem::take(&mut attachment.history);
let history_cursor = attachment.history_cursor;
let initial_replay_cursor = attachment
.replay
.iter()
.map(|event| event.sequence)
.max()
.unwrap_or(history_cursor);
let (events, _) = broadcast::channel(1024);
let bridge = Arc::new(Self {
runtime,
session_id,
model,
history,
history_cursor,
initial_replay_cursor,
events,
routing: Mutex::new(EventRoutingState::default()),
});
let weak = Arc::downgrade(&bridge);
tokio::spawn(async move {
loop {
let event = match attachment.next_event().await {
Ok(event) => event,
Err(error) => FrontendEvent {
sequence: u64::MAX,
kind: "runtime_disconnected".into(),
payload: json!({
"type": "runtime_disconnected",
"message": error.to_string(),
}),
},
};
let terminal = disconnect_message(&event).is_some();
let Some(bridge) = weak.upgrade() else {
return;
};
let mut routing = bridge
.routing
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if routing.active {
drop(routing);
let _ = bridge.events.send(event);
} else if routing.buffer(event).is_err() {
return;
}
if terminal {
return;
}
}
});
Ok(bridge)
}
fn subscribe(&self) -> broadcast::Receiver<FrontendEvent> {
self.events.subscribe()
}
fn session_id(&self) -> &str {
&self.session_id
}
fn activate_and_route(
&self,
history: Vec<Value>,
acknowledged_cursor: u64,
tx: &mpsc::UnboundedSender<Value>,
) -> Result<(), String> {
let mut routing = self
.routing
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if routing.active {
return Ok(());
}
let pending = routing.activate(acknowledged_cursor)?;
for update in history {
let _ = tx.send(update);
}
let mut terminal = None;
for event in pending {
let _ = route_event_projection(tx, self.session_id(), &event);
if disconnect_message(&event).is_some() {
terminal = Some(event);
break;
}
}
drop(routing);
if let Some(event) = terminal {
let _ = self.events.send(event);
}
Ok(())
}
async fn submit(
&self,
prompt: String,
image_urls: Vec<String>,
) -> Result<String, FrontendRuntimeError> {
let mut events = self.subscribe();
let reply = self.runtime.submit_with_images(prompt, image_urls).await?;
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
match tokio::time::timeout_at(deadline, events.recv()).await {
Ok(Ok(event)) if is_turn_boundary(&event) => return Ok(reply),
Ok(Ok(event)) => {
if let Some(message) = disconnect_message(&event) {
return Err(FrontendRuntimeError::Transport(message.to_string()));
}
}
Ok(Err(broadcast::error::RecvError::Lagged(skipped))) => {
return Err(FrontendRuntimeError::Transport(format!(
"SDK event stream lagged by {skipped} event(s)"
)));
}
Ok(Err(broadcast::error::RecvError::Closed)) => {
return Err(FrontendRuntimeError::Closed);
}
Err(_) => {
return Err(FrontendRuntimeError::Transport(
"SDK event stream did not confirm turn completion".into(),
));
}
}
}
}
async fn interrupt(&self) -> bool {
self.runtime.interrupt().await.unwrap_or(false)
}
}
fn is_turn_boundary(event: &FrontendEvent) -> bool {
matches!(
event.payload.get("type").and_then(Value::as_str),
Some("turn_succeeded" | "turn_failed" | "turn_interrupted")
)
}
pub struct AcpServer {
runtime: Arc<AcpRuntimeBridge>,
compatibility: AcpCompatibilityProfile,
session_open: AtomicBool,
history_replayed: AtomicBool,
prompt_active: AtomicBool,
}
enum RouterCommand {
Prompt {
id: Value,
prompt: String,
image_urls: Vec<String>,
frontend_reply: bool,
},
}
async fn project_prompt(
server: &AcpServer,
events: &mut tokio::sync::broadcast::Receiver<FrontendEvent>,
prompt: String,
image_urls: Vec<String>,
frontend_reply: bool,
tx: &mpsc::UnboundedSender<Value>,
) -> (Result<Value, FrontendRuntimeError>, bool) {
let runtime = server.runtime.clone();
let mut submit = tokio::spawn(async move { runtime.submit(prompt, image_urls).await });
let mut streamed = String::new();
let mut terminal = false;
let submit_result = loop {
tokio::select! {
result = &mut submit => break match result {
Ok(result) => result,
Err(error) => Err(FrontendRuntimeError::Transport(format!(
"ACP runtime submit task failed: {error}"
))),
},
event = events.recv() => match event {
Ok(event) => {
if let Some(message) = disconnect_message(&event) {
submit.abort();
terminal = true;
break Err(FrontendRuntimeError::Transport(message.to_string()));
}
if let Some(text) = assistant_event_text(&event) {
streamed.push_str(text);
}
if !route_event_projection(tx, server.runtime.session_id(), &event) {
submit.abort();
terminal = true;
break Err(FrontendRuntimeError::Closed);
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
submit.abort();
terminal = true;
break Err(FrontendRuntimeError::Transport(format!(
"ACP event stream lagged by {skipped} events"
)));
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
submit.abort();
terminal = true;
break Err(FrontendRuntimeError::Closed);
}
}
}
};
while !terminal {
match events.try_recv() {
Ok(event) => {
if let Some(message) = disconnect_message(&event) {
terminal = true;
return (
Err(FrontendRuntimeError::Transport(message.to_string())),
terminal,
);
}
if let Some(text) = assistant_event_text(&event) {
streamed.push_str(text);
}
if !route_event_projection(tx, server.runtime.session_id(), &event) {
terminal = true;
return (Err(FrontendRuntimeError::Closed), terminal);
}
}
Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break,
Err(tokio::sync::broadcast::error::TryRecvError::Lagged(skipped)) => {
terminal = true;
return (
Err(FrontendRuntimeError::Transport(format!(
"ACP event stream lagged by {skipped} events"
))),
terminal,
);
}
Err(tokio::sync::broadcast::error::TryRecvError::Closed) => {
terminal = true;
return (Err(FrontendRuntimeError::Closed), terminal);
}
}
}
let result = match submit_result {
Ok(reply) if reply.starts_with(&streamed) => {
let missing = &reply[streamed.len()..];
if !missing.is_empty() {
let _ = tx.send(session_update(
server.runtime.session_id(),
json!({
"sessionUpdate": "agent_message_chunk",
"content": {"type": "text", "text": missing}
}),
));
streamed.push_str(missing);
}
if streamed.is_empty() {
Err(FrontendRuntimeError::Execution {
operation: crate::SdkOperation::Input,
message: "runtime completed without assistant output".into(),
})
} else {
if frontend_reply {
Ok(json!({"reply": reply}))
} else {
Ok(json!({"stopReason": "end_turn"}))
}
}
}
Ok(reply) => Err(FrontendRuntimeError::Execution {
operation: crate::SdkOperation::Input,
message: format!(
"runtime reply did not match streamed assistant output (reply {} bytes, stream {} bytes)",
reply.len(),
streamed.len()
),
}),
Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Interrupted)) => {
Ok(json!({"stopReason": "cancelled"}))
}
Err(error) => Err(error),
};
(result, terminal)
}
impl AcpServer {
pub async fn new(runtime: Arc<dyn SdkRuntime>) -> Result<Arc<Self>, FrontendRuntimeError> {
Self::new_with_profile(runtime, AcpCompatibilityProfile::Standard).await
}
pub async fn new_with_profile(
runtime: Arc<dyn SdkRuntime>,
compatibility: AcpCompatibilityProfile,
) -> Result<Arc<Self>, FrontendRuntimeError> {
Ok(Arc::new(Self {
runtime: AcpRuntimeBridge::connect(runtime).await?,
compatibility,
session_open: AtomicBool::new(false),
history_replayed: AtomicBool::new(false),
prompt_active: AtomicBool::new(false),
}))
}
pub fn runtime(&self) -> &Arc<dyn SdkRuntime> {
&self.runtime.runtime
}
fn initialize(&self) -> Value {
let mut frontend_methods = crate::FrontendFacadeMethod::ALL
.into_iter()
.map(crate::FrontendFacadeMethod::wire_name)
.collect::<Vec<_>>();
frontend_methods.extend(["session/cancel", "session/steer", "session/respond"]);
json!({
"protocolVersion": ACP_PROTOCOL_VERSION,
"agentCapabilities": {
"loadSession": true,
"sessionCapabilities": {"resume": {}},
"promptCapabilities": {"image": false, "embeddedContext": false},
"_meta": {
"supercode": {
"frontend": {
"schemaVersion": crate::frontend::FRONTEND_RUNTIME_SCHEMA_VERSION,
"contract": "supercode.frontend.contract.v2",
"eventMethod": "frontend.v2.event",
"runtimeOwnedByClient": false,
"methods": frontend_methods
}
}
}
},
"authMethods": [],
"agentInfo": {
"name": "supercode",
"title": "Supercode",
"version": env!("CARGO_PKG_VERSION")
}
})
}
fn open_session(&self, requested: Option<&str>) -> Result<Value, String> {
if let Some(requested) = requested {
if requested != self.runtime.session_id() {
return Err(format!(
"runtime `{}` is not session `{requested}`",
self.runtime.session_id()
));
}
}
self.session_open.store(true, Ordering::SeqCst);
Ok(json!({"sessionId": self.runtime.session_id()}))
}
fn goose_defaults(&self) -> Result<Value, String> {
if self.compatibility != AcpCompatibilityProfile::GooseAcd3c135 {
return Err("Goose defaults are unavailable on the standard ACP profile".into());
}
Ok(json!({
"providerId": "supercode",
"modelId": self.runtime.model,
}))
}
fn validate_new_session(&self, params: &Value) -> Result<(), String> {
if self.compatibility != AcpCompatibilityProfile::GooseAcd3c135 {
return Ok(());
}
if params
.get("mcpServers")
.and_then(Value::as_array)
.is_some_and(|servers| !servers.is_empty())
{
return Err("Goose frontend cannot mutate the selected runtime's MCP servers".into());
}
Ok(())
}
fn history_updates_once(&self) -> Vec<Value> {
if self.compatibility != AcpCompatibilityProfile::GooseAcd3c135
|| self.history_replayed.swap(true, Ordering::SeqCst)
{
return Vec::new();
}
project_history(
self.runtime.session_id(),
&self.runtime.history,
self.runtime.history_cursor,
)
}
fn prompt_text(params: &Value) -> Result<String, String> {
let prompt = params
.get("prompt")
.and_then(Value::as_array)
.ok_or_else(|| "session/prompt requires a `prompt` content array".to_string())?;
let text = prompt
.iter()
.filter(|part| part.get("type").and_then(Value::as_str) == Some("text"))
.filter_map(|part| part.get("text").and_then(Value::as_str))
.collect::<Vec<_>>()
.join("\n");
if text.is_empty() {
Err("session/prompt contains no text content".into())
} else {
Ok(text)
}
}
fn validate_session(&self, params: &Value) -> Result<(), String> {
if !self.session_open.load(Ordering::SeqCst) {
return Err("open a session before prompting".into());
}
let requested = params
.get("sessionId")
.and_then(Value::as_str)
.ok_or_else(|| "request omitted `sessionId`".to_string())?;
if requested != self.runtime.session_id() {
return Err(format!(
"runtime `{}` is not session `{requested}`",
self.runtime.session_id()
));
}
Ok(())
}
}
fn response(id: Value, result: Result<Value, String>) -> Value {
match result {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(message) => json!({
"jsonrpc": "2.0",
"id": id,
"error": {"code": -32000, "message": message}
}),
}
}
fn sdk_response(id: Value, result: Result<Value, FrontendRuntimeError>) -> Value {
match result {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(error) => {
let code = match error.code() {
crate::SdkErrorCode::Busy => -32000,
crate::SdkErrorCode::Unauthenticated => -32030,
crate::SdkErrorCode::Unauthorized => -32031,
crate::SdkErrorCode::ControllerRequired => -32032,
crate::SdkErrorCode::LeaseExpired => -32033,
_ => -32002,
};
let mut envelope = json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": code,
"name": error.code(),
"operation": error.operation(),
"message": error.to_string(),
}
});
if let Some(detail) = envelope.get_mut("error").and_then(Value::as_object_mut) {
match error {
FrontendRuntimeError::Unauthorized { permission } => {
detail.insert("permission".into(), Value::String(permission));
}
FrontendRuntimeError::ControllerRequired {
holder,
expires_at_ms,
} => {
if let Some(holder) = holder {
detail.insert("holder".into(), Value::String(holder));
}
if let Some(expires_at_ms) = expires_at_ms {
detail.insert("expiresAtMs".into(), json!(expires_at_ms));
}
}
_ => {}
}
}
envelope
}
}
}
fn session_update(session_id: &str, update: Value) -> Value {
json!({
"jsonrpc": "2.0",
"method": "session/update",
"params": {"sessionId": session_id, "update": update}
})
}
fn projected_session_update(session_id: &str, event: &FrontendEvent, mut update: Value) -> Value {
if let Some(update) = update.as_object_mut() {
update.insert(
"_meta".into(),
json!({
"supercode": {
"sdkSequence": event.sequence,
"sdkKind": event.kind,
}
}),
);
}
session_update(session_id, update)
}
fn frontend_event_notification(session_id: &str, event: &FrontendEvent) -> Value {
json!({
"jsonrpc": "2.0",
"method": "frontend.v2.event",
"params": {"sessionId": session_id, "event": event}
})
}
fn approval_request_id(request_id: u64) -> String {
format!("{ACP_APPROVAL_REQUEST_ID_PREFIX}{request_id}")
}
fn approval_request_id_from_wire(message: &Value) -> Option<u64> {
let wire_id = message.get("id")?.as_str()?;
let request_id = wire_id
.strip_prefix(ACP_APPROVAL_REQUEST_ID_PREFIX)?
.parse::<u64>()
.ok()?;
(approval_request_id(request_id) == wire_id).then_some(request_id)
}
fn project_approval_request(session_id: &str, event: &FrontendEvent) -> Option<Value> {
if event.payload.get("type").and_then(Value::as_str) != Some("request") {
return None;
}
let request: FrontendRequest =
serde_json::from_value(event.payload.get("request")?.clone()).ok()?;
if request.kind != FrontendRequestKind::Approval {
return None;
}
let tool = request
.payload
.get("tool")
.and_then(Value::as_str)
.unwrap_or("tool");
let title = request
.payload
.get("subject")
.and_then(Value::as_str)
.filter(|subject| !subject.is_empty())
.unwrap_or(tool);
let raw_input = request
.payload
.get("raw_args")
.cloned()
.unwrap_or(Value::Null);
Some(json!({
"jsonrpc":"2.0",
"id":approval_request_id(request.id),
"method":"session/request_permission",
"params":{
"sessionId":session_id,
"toolCall":{
"toolCallId":format!("supercode-request-{}", request.id),
"title":title,
"kind":"other",
"status":"pending",
"rawInput":raw_input,
"_meta":{"supercode":{"requestId":request.id, "tool":tool}},
},
"options":[
{"optionId":"allow_once", "name":"Allow once", "kind":"allow_once"},
{"optionId":"allow_for_session", "name":"Allow for session", "kind":"allow_always"},
{"optionId":"deny", "name":"Deny", "kind":"reject_once"},
],
},
}))
}
fn decode_approval_response(message: &Value) -> Option<FrontendResponse> {
let request_id = approval_request_id_from_wire(message)?;
let result = message.get("result");
let error = message.get("error");
let exact_top_level = message.as_object().is_some_and(|object| {
object.len() == 3
&& object.contains_key("jsonrpc")
&& object.contains_key("id")
&& (object.contains_key("result") ^ object.contains_key("error"))
});
let valid_envelope = message.get("jsonrpc").and_then(Value::as_str) == Some("2.0")
&& message.get("method").is_none()
&& exact_top_level
&& (result.is_some() ^ error.is_some());
let decision =
if valid_envelope && error.is_none() {
match message.pointer("/result/outcome") {
Some(outcome)
if result.and_then(Value::as_object).is_some_and(|object| {
object.len() == 1 && object.contains_key("outcome")
}) && outcome.as_object().is_some_and(|object| {
object.len() == 2
&& object.contains_key("outcome")
&& object.contains_key("optionId")
}) && outcome.get("outcome").and_then(Value::as_str) == Some("selected") =>
{
match outcome.get("optionId").and_then(Value::as_str) {
Some("allow_once") => FrontendApprovalDecision::Allow,
Some("allow_for_session") => FrontendApprovalDecision::AllowForSession,
_ => FrontendApprovalDecision::Deny,
}
}
_ => FrontendApprovalDecision::Deny,
}
} else {
FrontendApprovalDecision::Deny
};
Some(FrontendResponse::Approval {
request_id,
decision,
})
}
fn route_event_projection(
tx: &mpsc::UnboundedSender<Value>,
session_id: &str,
event: &FrontendEvent,
) -> bool {
if tx
.send(frontend_event_notification(session_id, event))
.is_err()
{
return false;
}
if let Some(request) = project_approval_request(session_id, event) {
if tx.send(request).is_err() {
return false;
}
}
if let Some(update) = project_event(session_id, event) {
if tx.send(update).is_err() {
return false;
}
}
true
}
fn project_event(session_id: &str, event: &FrontendEvent) -> Option<Value> {
let payload = &event.payload;
match payload.get("type").and_then(Value::as_str)? {
"text_delta" => Some(projected_session_update(
session_id,
event,
json!({
"sessionUpdate": "agent_message_chunk",
"content": {"type": "text", "text": payload.get("text")?}
}),
)),
"tool_call_started" => Some(projected_session_update(
session_id,
event,
json!({
"sessionUpdate": "tool_call",
"toolCallId": payload.get("id")?,
"title": payload.get("name")?,
"kind": "other",
"status": "in_progress",
"rawInput": payload.get("arguments").cloned().unwrap_or(Value::Null)
}),
)),
"tool_call_completed" => Some(projected_session_update(
session_id,
event,
json!({
"sessionUpdate": "tool_call_update",
"toolCallId": payload.get("id")?,
"status": if payload.get("is_error").and_then(Value::as_bool).unwrap_or(false) {
"failed"
} else {
"completed"
},
"content": [{
"type": "content",
"content": {
"type": "text",
"text": payload.get("output").cloned().unwrap_or(Value::String(String::new()))
}
}]
}),
)),
"background_output" => Some(projected_session_update(
session_id,
event,
json!({
"sessionUpdate": "agent_message_chunk",
"content": {"type": "text", "text": payload.get("chunk")?},
"nativeEvent": payload
}),
)),
_ => Some(projected_session_update(
session_id,
event,
json!({
"sessionUpdate": "agent_thought_chunk",
"content": {"type": "text", "text": ""},
"nativeEvent": payload,
}),
)),
}
}
fn canonical_content_blocks(message: &ChatMessage) -> Vec<Value> {
if let Some(content) = &message.content {
return if content.is_empty() {
Vec::new()
} else {
vec![json!({"type":"text", "text":content})]
};
}
message
.content_parts
.as_deref()
.unwrap_or_default()
.iter()
.filter_map(|part| match part.get("type").and_then(Value::as_str) {
Some("text") => part
.get("text")
.and_then(Value::as_str)
.map(|text| json!({"type":"text", "text":text})),
Some("image_url") => {
let url = part.pointer("/image_url/url").and_then(Value::as_str)?;
if let Some(data) = url.strip_prefix("data:") {
let (mime_type, data) = data.split_once(";base64,")?;
Some(json!({
"type":"image",
"data":data,
"mimeType":mime_type,
}))
} else {
Some(json!({
"type":"resource_link",
"name":"image attachment",
"uri":url,
}))
}
}
_ => None,
})
.collect()
}
fn history_update(
session_id: &str,
history_cursor: u64,
history_index: usize,
role: Role,
update: Value,
) -> Value {
let mut notification = session_update(session_id, update);
notification["params"]["update"]["_meta"] = json!({
"supercode": {
"historyCursor": history_cursor,
"historyIndex": history_index,
"canonicalRole": match role {
Role::System => "system",
Role::User => "user",
Role::Assistant => "assistant",
Role::Tool => "tool",
},
}
});
notification
}
fn project_history(session_id: &str, history: &[ChatMessage], history_cursor: u64) -> Vec<Value> {
let mut projected = Vec::new();
for (index, message) in history.iter().enumerate() {
match message.role {
Role::System => {}
Role::User | Role::Assistant => {
let session_update_kind = if message.role == Role::User {
"user_message_chunk"
} else {
"agent_message_chunk"
};
for content in canonical_content_blocks(message) {
projected.push(history_update(
session_id,
history_cursor,
index,
message.role,
json!({
"sessionUpdate":session_update_kind,
"content":content,
}),
));
}
if message.role == Role::Assistant {
for call in message.tool_calls() {
let raw_input = serde_json::from_str::<Value>(&call.function.arguments)
.unwrap_or_else(|_| Value::String(call.function.arguments.clone()));
projected.push(history_update(
session_id,
history_cursor,
index,
message.role,
json!({
"sessionUpdate":"tool_call",
"toolCallId":call.id,
"title":call.function.name,
"kind":"other",
"status":"in_progress",
"rawInput":raw_input,
}),
));
}
}
}
Role::Tool => {
let Some(tool_call_id) = message.tool_call_id.as_deref() else {
continue;
};
let content = canonical_content_blocks(message)
.into_iter()
.map(|content| json!({"type":"content", "content":content}))
.collect::<Vec<_>>();
projected.push(history_update(
session_id,
history_cursor,
index,
message.role,
json!({
"sessionUpdate":"tool_call_update",
"toolCallId":tool_call_id,
"status":"completed",
"content":content,
}),
));
}
}
}
projected
}
fn assistant_event_text(event: &FrontendEvent) -> Option<&str> {
match event.payload.get("type").and_then(Value::as_str)? {
"text_delta" => event.payload.get("text").and_then(Value::as_str),
_ => None,
}
}
fn disconnect_message(event: &FrontendEvent) -> Option<&str> {
(event.payload.get("type").and_then(Value::as_str) == Some("runtime_disconnected"))
.then(|| event.payload.get("message").and_then(Value::as_str))
.flatten()
}
pub async fn run_stdio<R, W>(
server: Arc<AcpServer>,
mut reader: R,
writer: W,
) -> std::io::Result<()>
where
R: AsyncBufRead + Unpin + Send + 'static,
W: AsyncWrite + Unpin + Send + 'static,
{
let (out_tx, mut out_rx) = mpsc::unbounded_channel::<Value>();
let pending_approvals = Arc::new(Mutex::new(HashSet::<u64>::new()));
let writer_pending_approvals = pending_approvals.clone();
let mut writer_task = tokio::spawn(async move {
let mut writer = writer;
while let Some(value) = out_rx.recv().await {
if value.get("method").and_then(Value::as_str) == Some("session/request_permission") {
if let Some(request_id) = approval_request_id_from_wire(&value) {
writer_pending_approvals
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(request_id);
}
}
writer.write_all(format!("{value}\n").as_bytes()).await?;
writer.flush().await?;
}
Ok::<(), std::io::Error>(())
});
let mut events = server.runtime.subscribe();
let event_server = server.clone();
let event_tx = out_tx.clone();
let (router_tx, mut router_rx) = mpsc::unbounded_channel::<RouterCommand>();
let (disconnect_tx, mut disconnect_rx) = mpsc::unbounded_channel::<()>();
let event_task = tokio::spawn(async move {
loop {
tokio::select! {
command = router_rx.recv() => match command {
Some(RouterCommand::Prompt { id, prompt, image_urls, frontend_reply }) => {
let (result, terminal) = project_prompt(
&event_server,
&mut events,
prompt,
image_urls,
frontend_reply,
&event_tx,
).await;
event_server.prompt_active.store(false, Ordering::SeqCst);
let _ = event_tx.send(sdk_response(id, result));
if terminal {
let _ = disconnect_tx.send(());
break;
}
}
None => break,
},
event = events.recv() => match event {
Ok(event) if disconnect_message(&event).is_some() => {
let _ = disconnect_tx.send(());
break;
}
Ok(event) if event_server.session_open.load(Ordering::SeqCst) => {
if !route_event_projection(
&event_tx,
event_server.runtime.session_id(),
&event,
) {
break;
}
}
Ok(_) => {}
Err(tokio::sync::broadcast::error::RecvError::Lagged(_))
| Err(tokio::sync::broadcast::error::RecvError::Closed) => {
let _ = disconnect_tx.send(());
break;
}
},
}
}
});
let mut writer_finished = false;
let mut terminal_error = None;
loop {
let mut line = String::new();
let bytes = tokio::select! {
bytes = reader.read_line(&mut line) => bytes?,
_ = disconnect_rx.recv() => break,
result = &mut writer_task => {
writer_finished = true;
terminal_error = match result {
Ok(Ok(())) => None,
Ok(Err(error)) => Some(error),
Err(error) => Some(std::io::Error::other(format!(
"ACP writer task failed: {error}"
))),
};
break;
}
};
if bytes == 0 {
break;
}
if line.len() > SERVER_MAX_LINE_BYTES {
let _ = out_tx.send(response(
Value::Null,
Err("ACP request exceeded size limit".into()),
));
continue;
}
let request: Value = match serde_json::from_str(line.trim()) {
Ok(request) => request,
Err(error) => {
let _ = out_tx.send(json!({
"jsonrpc": "2.0", "id": null,
"error": {"code": -32700, "message": error.to_string()}
}));
continue;
}
};
if request.get("method").is_none() {
if let Some(response) = decode_approval_response(&request) {
pending_approvals
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&response.request_id());
let _ = server.runtime.runtime.respond(response).await;
}
continue;
}
let method = request.get("method").and_then(Value::as_str).unwrap_or("");
let id = request.get("id").cloned();
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
if method == "session/cancel" {
let server = server.clone();
tokio::spawn(async move {
let _ = server.runtime.interrupt().await;
});
continue;
}
let Some(id) = id else {
continue;
};
match method {
"initialize" => {
let version = params.get("protocolVersion").and_then(Value::as_u64);
let result = if version == Some(ACP_PROTOCOL_VERSION) {
Ok(server.initialize())
} else {
Err(format!("unsupported ACP protocol version {version:?}"))
};
let _ = out_tx.send(response(id, result));
}
"_goose/unstable/defaults/read"
if server.compatibility == AcpCompatibilityProfile::GooseAcd3c135 =>
{
let _ = out_tx.send(response(id, server.goose_defaults()));
}
"session/new" => {
let mut opened = server
.validate_new_session(¶ms)
.and_then(|()| server.open_session(None));
if opened.is_ok() {
let goose = server.compatibility == AcpCompatibilityProfile::GooseAcd3c135;
let acknowledged_cursor = if goose {
server.runtime.history_cursor
} else {
server.runtime.initial_replay_cursor
};
if let Err(error) = server.runtime.activate_and_route(
server.history_updates_once(),
acknowledged_cursor,
&out_tx,
) {
server.session_open.store(false, Ordering::SeqCst);
opened = Err(error);
}
}
let _ = out_tx.send(response(id, opened));
}
"session/load" | "session/resume" => {
let requested = params.get("sessionId").and_then(Value::as_str);
let mut opened = server.open_session(requested);
if opened.is_ok() {
if let Err(error) = server.runtime.activate_and_route(
Vec::new(),
server.runtime.initial_replay_cursor,
&out_tx,
) {
server.session_open.store(false, Ordering::SeqCst);
opened = Err(error);
}
}
let _ = out_tx.send(response(id, opened));
}
"session/prompt" => {
let validation = server
.validate_session(¶ms)
.and_then(|_| AcpServer::prompt_text(¶ms));
let prompt = match validation {
Ok(prompt) => prompt,
Err(error) => {
let _ = out_tx.send(response(id, Err(error)));
continue;
}
};
if server
.prompt_active
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.is_err()
{
let _ = out_tx.send(sdk_response(
id,
Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)),
));
continue;
}
if let Err(error) = router_tx.send(RouterCommand::Prompt {
id,
prompt,
image_urls: Vec::new(),
frontend_reply: false,
}) {
server.prompt_active.store(false, Ordering::SeqCst);
let RouterCommand::Prompt { id, .. } = error.0;
let _ = out_tx.send(response(
id,
Err("ACP runtime event router is unavailable".into()),
));
}
}
"session/steer" | "frontend.v2.steer" => {
let result = match server.validate_session(¶ms) {
Ok(()) => match params.get("text").and_then(Value::as_str) {
Some(text) => server.runtime.runtime.steer(text.to_string()).await,
None => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Steer,
message: "session/steer requires string `text`".into(),
}),
},
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Steer,
message,
}),
};
let _ = out_tx.send(sdk_response(id, result.map(|()| json!({}))));
}
"session/respond" | "frontend.v2.respond" => {
let result = match server.validate_session(¶ms) {
Ok(()) => serde_json::from_value::<FrontendResponse>(
params.get("response").cloned().unwrap_or(Value::Null),
)
.map_err(|error| FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Respond,
message: error.to_string(),
}),
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Respond,
message,
}),
};
let result = match result {
Ok(response) => server.runtime.runtime.respond(response).await,
Err(error) => Err(error),
};
let _ = out_tx.send(sdk_response(id, result.map(|()| json!({}))));
}
"frontend.v2.describe" | "supercode/frontend/describe" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server
.runtime
.runtime
.describe()
.await
.and_then(|descriptor| {
serde_json::to_value(descriptor)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
}),
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Resume,
message,
}),
};
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.attach" | "supercode/frontend/attach" => {
let limit = params
.get("limit")
.and_then(Value::as_u64)
.unwrap_or(50)
.clamp(1, crate::server::SERVER_HISTORY_CAPACITY as u64)
as usize;
let after = params
.get("after_sequence")
.or_else(|| params.get("afterSequence"))
.and_then(Value::as_u64)
.unwrap_or_default();
let result = match server.validate_session(¶ms) {
Ok(()) => match server.runtime.runtime.attach(limit).await {
Ok(attachment) => {
let mut snapshot = crate::FrontendAttachSnapshot {
descriptor: attachment.descriptor,
history: attachment.history,
history_cursor: attachment.history_cursor,
replay: attachment.replay,
};
snapshot.replay.retain(|event| event.sequence > after);
serde_json::to_value(snapshot)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
}
Err(error) => Err(error),
},
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Events,
message,
}),
};
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.send_input" | "supercode/frontend/send_input" => {
let result = match server.validate_session(¶ms) {
Ok(()) => match params.get("prompt").and_then(Value::as_str) {
Some(prompt) => {
server
.runtime
.runtime
.clone()
.send_input(prompt.to_string())
.await
}
None => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message: "supercode/frontend/send_input requires string `prompt`"
.into(),
}),
},
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message,
}),
};
let _ = out_tx.send(sdk_response(id, result.map(|()| json!({"accepted":true}))));
}
"frontend.v2.submit" | "supercode/frontend/submit" => {
let prompt = match server.validate_session(¶ms).and_then(|_| {
params
.get("prompt")
.and_then(Value::as_str)
.map(str::to_owned)
.ok_or_else(|| {
"supercode/frontend/submit requires string `prompt`".to_string()
})
}) {
Ok(prompt) => prompt,
Err(error) => {
let _ = out_tx.send(sdk_response(
id,
Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message: error,
}),
));
continue;
}
};
let image_urls = match params.get("image_urls") {
None => Vec::new(),
Some(Value::Array(values)) => {
match values.iter().map(Value::as_str).collect::<Option<Vec<_>>>() {
Some(values) => values.into_iter().map(str::to_owned).collect(),
None => {
let _ = out_tx.send(sdk_response(
id,
Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message: "supercode/frontend/submit requires string entries in `image_urls`".into(),
}),
));
continue;
}
}
}
Some(_) => {
let _ = out_tx.send(sdk_response(
id,
Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message: "supercode/frontend/submit requires array `image_urls`"
.into(),
}),
));
continue;
}
};
if server
.prompt_active
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.is_err()
{
let _ = out_tx.send(sdk_response(
id,
Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)),
));
continue;
}
if let Err(error) = router_tx.send(RouterCommand::Prompt {
id,
prompt,
image_urls,
frontend_reply: true,
}) {
server.prompt_active.store(false, Ordering::SeqCst);
let RouterCommand::Prompt { id, .. } = error.0;
let _ = out_tx.send(sdk_response(
id,
Err(FrontendRuntimeError::Transport(
"ACP runtime event router is unavailable".into(),
)),
));
}
}
"frontend.v2.interrupt" | "supercode/frontend/interrupt" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server
.runtime
.runtime
.interrupt()
.await
.map(|interrupted| json!({"interrupted":interrupted})),
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Interrupt,
message,
}),
};
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.invoke" | "supercode/frontend/invoke" => {
let operation = params
.get("operation")
.cloned()
.ok_or_else(|| FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message: "supercode/frontend/invoke requires `operation`".into(),
})
.and_then(|value| {
serde_json::from_value(value).map_err(|error| {
FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message: error.to_string(),
}
})
});
let result = match (server.validate_session(¶ms), operation) {
(Ok(()), Ok(operation)) => server.runtime.runtime.invoke(operation).await,
(Err(message), _) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message,
}),
(_, Err(error)) => Err(error),
};
let result = result.and_then(|result| {
serde_json::to_value(result)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
});
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.lease" | "supercode/frontend/lease" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server.runtime.runtime.lease_snapshot().await,
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Events,
message,
}),
}
.and_then(|snapshot| {
serde_json::to_value(snapshot)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
});
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.take_control" | "supercode/frontend/take_control" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server.runtime.runtime.take_control().await,
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Input,
message,
}),
}
.and_then(|snapshot| {
serde_json::to_value(snapshot)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
});
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.heartbeat" | "supercode/frontend/heartbeat" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server.runtime.runtime.heartbeat().await,
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Events,
message,
}),
}
.and_then(|snapshot| {
serde_json::to_value(snapshot)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
});
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.detach" | "supercode/frontend/detach" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server.runtime.runtime.detach().await,
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Events,
message,
}),
}
.and_then(|snapshot| {
serde_json::to_value(snapshot)
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
});
if result.is_ok() {
server.session_open.store(false, Ordering::SeqCst);
}
let _ = out_tx.send(sdk_response(id, result));
}
"frontend.v2.close" | "supercode/frontend/close" => {
let result = match server.validate_session(¶ms) {
Ok(()) => server
.runtime
.runtime
.close()
.await
.map(|()| json!({"closed":true})),
Err(message) => Err(FrontendRuntimeError::InvalidArgument {
operation: crate::SdkOperation::Close,
message,
}),
};
let _ = out_tx.send(sdk_response(id, result));
}
other => {
let _ = out_tx.send(json!({
"jsonrpc": "2.0", "id": id,
"error": {"code": -32601, "message": format!("unknown ACP method `{other}`")}
}));
}
}
}
drop(out_tx);
event_task.abort();
let _ = event_task.await;
if !writer_finished {
terminal_error = match writer_task.await {
Ok(Ok(())) => terminal_error,
Ok(Err(error)) => Some(error),
Err(error) => Some(std::io::Error::other(format!(
"ACP writer task failed: {error}"
))),
};
}
let unresolved = pending_approvals
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.drain()
.collect::<Vec<_>>();
for request_id in unresolved {
let _ = server
.runtime
.runtime
.respond(FrontendResponse::Approval {
request_id,
decision: FrontendApprovalDecision::Deny,
})
.await;
}
match terminal_error {
Some(error) => Err(error),
None => Ok(()),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn projects_text_and_tool_events_without_transport_state() {
let sdk_event = FrontendEvent {
sequence: 41,
kind: "text_delta".into(),
payload: json!({"type": "text_delta", "text": "hello"}),
};
let text = project_event("s1", &sdk_event).unwrap();
assert_eq!(
text.pointer("/params/update/sessionUpdate")
.and_then(Value::as_str),
Some("agent_message_chunk")
);
assert_eq!(
text.pointer("/params/update/content/text")
.and_then(Value::as_str),
Some("hello")
);
assert_eq!(
text.pointer("/params/update/_meta/supercode/sdkSequence")
.and_then(Value::as_u64),
Some(41)
);
assert_eq!(
text.pointer("/params/update/_meta/supercode/sdkKind")
.and_then(Value::as_str),
Some("text_delta")
);
let tool = project_event(
"s1",
&FrontendEvent {
sequence: 42,
kind: "tool_call_started".into(),
payload: json!({
"type": "tool_call_started", "id": "call-1",
"name": "read_file", "arguments": "{}"
}),
},
)
.unwrap();
assert_eq!(
tool.pointer("/params/update/toolCallId")
.and_then(Value::as_str),
Some("call-1")
);
}
#[test]
fn preserves_named_sdk_errors_in_acp_responses() {
let response = sdk_response(
json!(7),
Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)),
);
assert_eq!(response["error"]["name"], "busy");
assert_eq!(response["error"]["operation"], "input");
}
#[test]
fn approval_requests_and_responses_preserve_runtime_correlation() {
let event = FrontendEvent {
sequence: 43,
kind: "request".into(),
payload: json!({
"type":"request",
"request":{
"id":17,
"kind":"approval",
"payload":{
"tool":"bash",
"subject":"echo reviewed",
"raw_args":{"command":"echo reviewed"},
},
},
}),
};
let request = project_approval_request("s1", &event).unwrap();
assert_eq!(request["id"], "supercode-approval-17");
assert_eq!(request["method"], "session/request_permission");
assert_eq!(request["params"]["sessionId"], "s1");
assert_eq!(request["params"]["toolCall"]["title"], "echo reviewed");
assert_eq!(request["params"]["options"][0]["optionId"], "allow_once");
assert_eq!(
decode_approval_response(&json!({
"jsonrpc":"2.0",
"id":"supercode-approval-17",
"result":{"outcome":{"outcome":"selected", "optionId":"allow_for_session"}},
})),
Some(FrontendResponse::Approval {
request_id: 17,
decision: FrontendApprovalDecision::AllowForSession,
})
);
assert_eq!(
decode_approval_response(&json!({
"jsonrpc":"2.0",
"id":"supercode-approval-17",
"result":{"outcome":{"outcome":"cancelled"}},
})),
Some(FrontendResponse::Approval {
request_id: 17,
decision: FrontendApprovalDecision::Deny,
})
);
for malformed in [
json!({
"id":"supercode-approval-17",
"result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
}),
json!({
"jsonrpc":"2.0",
"id":"supercode-approval-17",
"result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
"error":{"code":-32603, "message":"invalid"},
}),
json!({
"jsonrpc":"2.0",
"id":"supercode-approval-17",
"result":{"outcome":{"outcome":"selected", "optionId":"unknown"}},
}),
json!({
"jsonrpc":"2.0",
"id":"supercode-approval-17",
"result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
"params":{},
}),
json!({
"jsonrpc":"2.0",
"id":"supercode-approval-17",
"result":{
"outcome":{"outcome":"selected", "optionId":"allow_once"},
"extra":true,
},
}),
] {
assert_eq!(
decode_approval_response(&malformed),
Some(FrontendResponse::Approval {
request_id: 17,
decision: FrontendApprovalDecision::Deny,
})
);
}
assert!(decode_approval_response(&json!({
"jsonrpc":"2.0",
"id":"supercode-approval-017",
"result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
}))
.is_none());
assert!(decode_approval_response(&json!({"jsonrpc":"2.0", "id":99})).is_none());
}
#[test]
fn inactive_event_routing_is_bounded_and_fails_instead_of_truncating() {
let mut routing = EventRoutingState::default();
for sequence in 1..=FRONTEND_REPLAY_CAPACITY as u64 {
routing
.buffer(FrontendEvent {
sequence,
kind: "loop_tick".into(),
payload: json!({"type":"loop_tick", "sequence":sequence}),
})
.unwrap();
}
assert_eq!(routing.pending.len(), FRONTEND_REPLAY_CAPACITY);
assert!(routing
.buffer(FrontendEvent {
sequence: FRONTEND_REPLAY_CAPACITY as u64 + 1,
kind: "loop_tick".into(),
payload: json!({"type":"loop_tick"}),
})
.is_err());
assert!(routing.pending.is_empty(), "failed replay is never partial");
let error = routing.activate(0).unwrap_err();
assert!(error.contains(&FRONTEND_REPLAY_CAPACITY.to_string()));
assert!(error.contains("reconnect required"));
assert!(!routing.active, "a gapped attachment must never activate");
}
}