use anyhow::{anyhow, bail, Context, Result};
use chrono::{SecondsFormat, Utc};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use sha2::{Digest, Sha256};
use std::collections::{HashMap, HashSet};
use std::fmt;
use std::fs::{self, File, OpenOptions};
use std::io::{self, BufRead, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc;
use std::sync::{Arc, Mutex, OnceLock};
use std::thread;
use crate::runtime;
const ACP_PROTOCOL_VERSION: u32 = 1;
const REDACTED: &str = runtime::redaction::REDACTED_MARKER;
static SESSION_COUNTER: AtomicU64 = AtomicU64::new(1);
static SESSION_LOCKS: OnceLock<Mutex<HashMap<PathBuf, Arc<Mutex<()>>>>> = OnceLock::new();
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct AcpSession {
session_id: String,
cwd: PathBuf,
mcp_servers: Vec<Value>,
orchestration_slug: Option<String>,
trace_ids: Vec<String>,
transcript: Vec<AcpTranscriptEntry>,
cancelled: bool,
#[serde(default)]
turn_sequence: u64,
#[serde(default)]
last_event_sequence: u64,
created_at: String,
updated_at: String,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct AcpTranscriptEntry {
role: String,
message_id: String,
text: String,
at: String,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct AcpEvent {
sequence: u64,
kind: String,
at: String,
payload: Value,
}
#[derive(Clone, Copy, Debug)]
enum AcpErrorKind {
InvalidRequest,
MethodNotFound,
InvalidParams,
ResourceNotFound,
Internal,
}
impl AcpErrorKind {
fn code(self) -> i32 {
match self {
Self::InvalidRequest => -32600,
Self::MethodNotFound => -32601,
Self::InvalidParams => -32602,
Self::ResourceNotFound => -32002,
Self::Internal => -32603,
}
}
fn label(self) -> &'static str {
match self {
Self::InvalidRequest => "invalid_request",
Self::MethodNotFound => "method_not_found",
Self::InvalidParams => "invalid_params",
Self::ResourceNotFound => "resource_not_found",
Self::Internal => "internal_error",
}
}
}
#[derive(Debug)]
struct AcpProtocolError {
kind: AcpErrorKind,
message: String,
}
impl fmt::Display for AcpProtocolError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.message)
}
}
impl std::error::Error for AcpProtocolError {}
#[derive(Default)]
struct AcpConnectionState {
initialized: bool,
}
struct AcpPromptExecution {
response_text: String,
tool_title: String,
}
trait AcpPromptExecutor: Send + Sync {
fn execute(
&self,
session: &AcpSession,
prompt: &str,
context: &AcpExecutionContext<'_>,
) -> Result<AcpPromptExecution>;
}
#[allow(dead_code)] #[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum AcpPermissionDecision {
Granted,
Denied,
Cancelled,
}
#[derive(Default)]
struct CancellationToken {
requested: AtomicU64,
consumed: AtomicU64,
}
impl CancellationToken {
fn request(&self) -> u64 {
self.requested.fetch_add(1, Ordering::AcqRel) + 1
}
fn generation(&self) -> u64 {
self.requested.load(Ordering::Acquire)
}
fn pending_after(&self, baseline: u64) -> Option<u64> {
let requested = self.generation();
let consumed = self.consumed.load(Ordering::Acquire);
(requested > baseline && requested > consumed).then_some(requested)
}
fn consume_through(&self, generation: u64) {
self.consumed.fetch_max(generation, Ordering::AcqRel);
}
fn was_consumed(&self, generation: u64) -> bool {
self.consumed.load(Ordering::Acquire) >= generation
}
}
struct PendingPermission {
session_id: String,
sender: mpsc::Sender<Value>,
}
#[derive(Default)]
struct PermissionRegistry {
pending: HashMap<String, PendingPermission>,
early_responses: HashMap<String, Value>,
closed: bool,
}
#[allow(dead_code)] #[derive(Default)]
struct AcpConnectionRuntime {
cancellations: Mutex<HashMap<String, Arc<CancellationToken>>>,
permissions: Mutex<PermissionRegistry>,
request_counter: AtomicU64,
}
#[allow(dead_code)]
impl AcpConnectionRuntime {
fn cancellation_token(&self, session_id: &str) -> Arc<CancellationToken> {
let mut cancellations = self
.cancellations
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
Arc::clone(
cancellations
.entry(session_id.to_string())
.or_insert_with(|| Arc::new(CancellationToken::default())),
)
}
fn next_permission_id(&self) -> String {
let sequence = self.request_counter.fetch_add(1, Ordering::Relaxed) + 1;
format!("permission-{sequence}")
}
fn register_permission(&self, request_id: &str, session_id: &str) -> mpsc::Receiver<Value> {
let (sender, receiver) = mpsc::channel();
let mut permissions = self
.permissions
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if let Some(response) = permissions.early_responses.remove(request_id) {
let _ = sender.send(response);
} else if permissions.closed {
let _ = sender.send(cancelled_permission_response(request_id));
} else {
permissions.pending.insert(
request_id.to_string(),
PendingPermission {
session_id: session_id.to_string(),
sender,
},
);
}
receiver
}
fn resolve_response(&self, response: &Value) {
let Some(request_id) = response_id_key(response) else {
return;
};
let mut permissions = self
.permissions
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if let Some(pending) = permissions.pending.remove(&request_id) {
let _ = pending.sender.send(response.clone());
} else if !permissions.closed {
permissions
.early_responses
.insert(request_id, response.clone());
}
}
fn cancel_permissions(&self, session_id: &str) {
let mut permissions = self
.permissions
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let request_ids = permissions
.pending
.iter()
.filter(|(_, pending)| pending.session_id == session_id)
.map(|(request_id, _)| request_id.clone())
.collect::<Vec<_>>();
for request_id in request_ids {
if let Some(pending) = permissions.pending.remove(&request_id) {
let _ = pending
.sender
.send(cancelled_permission_response(&request_id));
}
}
}
fn remove_permission(&self, request_id: &str) {
self.permissions
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.pending
.remove(request_id);
}
fn close_pending(&self) {
let mut permissions = self
.permissions
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
permissions.closed = true;
let pending = std::mem::take(&mut permissions.pending);
for (request_id, pending) in pending {
let _ = pending
.sender
.send(cancelled_permission_response(&request_id));
}
}
}
#[allow(dead_code)] struct AcpExecutionContext<'a> {
root: &'a Path,
session_id: &'a str,
tool_call_id: &'a str,
sink: Option<&'a dyn AcpResponseSink>,
runtime: &'a AcpConnectionRuntime,
cancellation: Arc<CancellationToken>,
cancel_baseline: u64,
}
#[allow(dead_code)]
impl AcpExecutionContext<'_> {
fn is_cancelled(&self) -> bool {
self.cancellation
.pending_after(self.cancel_baseline)
.is_some()
}
fn request_permission(&self, title: &str) -> Result<AcpPermissionDecision> {
let request_id = self.runtime.next_permission_id();
let receiver = self
.runtime
.register_permission(&request_id, self.session_id);
let request = json!({
"jsonrpc": "2.0",
"id": request_id,
"method": "session/request_permission",
"params": {
"sessionId": self.session_id,
"toolCall": {
"toolCallId": self.tool_call_id,
"title": title,
"kind": "think",
"status": "pending"
},
"options": [
{ "optionId": "allow_once", "name": "Permitir uma vez", "kind": "allow_once" },
{ "optionId": "reject_once", "name": "Negar", "kind": "reject_once" }
]
}
});
append_session_event(
self.root,
self.session_id,
"permission_requested",
json!({ "requestId": request_id, "toolCallId": self.tool_call_id, "title": title }),
)?;
let Some(sink) = self.sink else {
self.runtime.remove_permission(&request_id);
return Ok(AcpPermissionDecision::Denied);
};
sink.emit(request)?;
let response = loop {
if self.is_cancelled() {
self.runtime.remove_permission(&request_id);
return Ok(AcpPermissionDecision::Cancelled);
}
match receiver.recv_timeout(std::time::Duration::from_millis(20)) {
Ok(response) => break response,
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
return Ok(AcpPermissionDecision::Denied);
}
}
};
let decision = permission_decision(&response);
append_session_event(
self.root,
self.session_id,
"permission_resolved",
json!({ "requestId": request_id, "decision": format!("{decision:?}") }),
)?;
Ok(decision)
}
}
fn response_id_key(response: &Value) -> Option<String> {
response.get("id").map(|id| match id {
Value::String(value) => value.clone(),
value => value.to_string(),
})
}
fn cancelled_permission_response(request_id: &str) -> Value {
json!({
"jsonrpc": "2.0",
"id": request_id,
"result": { "outcome": { "outcome": "cancelled" } }
})
}
#[allow(dead_code)]
fn permission_decision(response: &Value) -> AcpPermissionDecision {
if response.get("error").is_some() {
return AcpPermissionDecision::Denied;
}
match response
.pointer("/result/outcome/outcome")
.and_then(Value::as_str)
{
Some("selected") => match response
.pointer("/result/outcome/optionId")
.and_then(Value::as_str)
{
Some("allow_once" | "allow_always") => AcpPermissionDecision::Granted,
_ => AcpPermissionDecision::Denied,
},
Some("cancelled") => AcpPermissionDecision::Cancelled,
_ => AcpPermissionDecision::Denied,
}
}
trait AcpResponseSink: Send + Sync {
fn emit(&self, response: Value) -> Result<()>;
}
#[derive(Clone)]
struct ChannelResponseSink(mpsc::Sender<Value>);
impl AcpResponseSink for ChannelResponseSink {
fn emit(&self, response: Value) -> Result<()> {
self.0
.send(response)
.map_err(|_| anyhow!("ACP response channel closed"))
}
}
struct HandoffPromptExecutor;
impl AcpPromptExecutor for HandoffPromptExecutor {
fn execute(
&self,
session: &AcpSession,
prompt: &str,
_context: &AcpExecutionContext<'_>,
) -> Result<AcpPromptExecution> {
let response_text = session
.orchestration_slug
.as_deref()
.and_then(|orchestration| {
crate::context_service::ContextService::new(&session.cwd)
.build(orchestration, "execution", Some(prompt), false)
.ok()
})
.map(|snapshot| build_agent_response_with_context(session, prompt, &snapshot))
.unwrap_or_else(|| build_agent_response(session, prompt));
Ok(AcpPromptExecution {
response_text,
tool_title: "Preparar handoff SDD supervisionado".to_string(),
})
}
}
pub fn serve_stdio(root: &Path) -> Result<()> {
let stdin = io::stdin();
let mut stdout = io::stdout();
serve_lines(root, stdin.lock(), &mut stdout)
}
pub fn serve_lines<R: BufRead, W: Write + Send>(
root: &Path,
reader: R,
writer: &mut W,
) -> Result<()> {
serve_lines_with_executor(root, reader, writer, &HandoffPromptExecutor)
}
fn serve_lines_with_executor<R: BufRead, W: Write + Send>(
root: &Path,
reader: R,
writer: &mut W,
executor: &dyn AcpPromptExecutor,
) -> Result<()> {
thread::scope(|scope| -> Result<()> {
let (outbound, responses) = mpsc::channel::<Value>();
let runtime = Arc::new(AcpConnectionRuntime::default());
let writer_thread = scope.spawn(move || -> Result<()> {
for response in responses {
writeln!(writer, "{}", serde_json::to_string(&response)?)?;
writer.flush()?;
}
Ok(())
});
let mut state = AcpConnectionState::default();
let mut prompt_threads = Vec::new();
for line in reader.lines() {
let line = line?;
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let request: Value = match serde_json::from_str(trimmed) {
Ok(value) => value,
Err(error) => {
outbound.send(jsonrpc_error(Value::Null, -32700, &error.to_string()))?;
continue;
}
};
if is_jsonrpc_response(&request) {
runtime.resolve_response(&request);
continue;
}
let is_cancel = request.get("jsonrpc").and_then(Value::as_str) == Some("2.0")
&& request.get("method").and_then(Value::as_str) == Some("session/cancel")
&& (state.initialized
|| request
.get("id")
.is_none_or(|request_id| request_id.is_null()));
if is_cancel {
let id = request.get("id").cloned().unwrap_or(Value::Null);
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
let result = cancel_session_with_runtime(root, ¶ms, &runtime);
if !id.is_null() {
let response = match result {
Ok(()) => jsonrpc_result(id, json!({})),
Err(error) => {
let kind = error
.downcast_ref::<AcpProtocolError>()
.map(|error| error.kind)
.unwrap_or(AcpErrorKind::Internal);
structured_error(id, kind, &error.to_string(), "session/cancel")
}
};
outbound.send(response)?;
}
continue;
}
if is_ready_prompt_request(&request, &state) {
let outbound = outbound.clone();
let runtime = Arc::clone(&runtime);
let id = request.get("id").cloned().unwrap_or(Value::Null);
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
prompt_threads.push(scope.spawn(move || {
let sink = ChannelResponseSink(outbound.clone());
let responses = match prompt_session_with_executor_and_sink(
root,
id.clone(),
¶ms,
executor,
Some(&sink),
&runtime,
) {
Ok(responses) => responses,
Err(error) => {
let kind = error
.downcast_ref::<AcpProtocolError>()
.map(|error| error.kind)
.unwrap_or(AcpErrorKind::Internal);
vec![structured_error(
id,
kind,
&error.to_string(),
"session/prompt",
)]
}
};
for response in responses {
if outbound.send(response).is_err() {
break;
}
}
}));
continue;
}
for response in
handle_request_with_state_and_executor(root, request, &mut state, executor)
{
outbound.send(response)?;
}
}
runtime.close_pending();
for prompt_thread in prompt_threads {
prompt_thread
.join()
.map_err(|_| anyhow!("ACP prompt worker panicked"))?;
}
drop(outbound);
writer_thread
.join()
.map_err(|_| anyhow!("ACP writer worker panicked"))??;
Ok(())
})
}
fn is_jsonrpc_response(message: &Value) -> bool {
message.get("jsonrpc").and_then(Value::as_str) == Some("2.0")
&& message.get("method").is_none()
&& (message.get("result").is_some() || message.get("error").is_some())
}
fn is_ready_prompt_request(request: &Value, state: &AcpConnectionState) -> bool {
state.initialized
&& request.get("jsonrpc").and_then(Value::as_str) == Some("2.0")
&& request.get("id").is_some_and(|id| !id.is_null())
&& request.get("method").and_then(Value::as_str) == Some("session/prompt")
}
pub fn doctor(root: &Path) -> Value {
let sessions_dir = sessions_dir(root);
let agent_manifest = runtime::agents::agent_manifest("sdd-orchestrator");
json!({
"status": if agent_manifest.is_ok() { "pass" } else { "fail" },
"protocolVersion": ACP_PROTOCOL_VERSION,
"sdk": acp_sdk_marker(),
"sessionsDir": runtime::platform::display_path(&sessions_dir),
"sessionsDirExists": sessions_dir.exists(),
"agent": "sdd-orchestrator",
"capabilities": initialize_agent_capabilities(),
"error": agent_manifest.err().map(|err| err.to_string()),
})
}
pub fn config_actions(
root: &Path,
targets: &str,
force: bool,
dry_run: bool,
) -> Result<Vec<String>> {
let mut actions = Vec::new();
let selected = parse_targets(targets)?;
if selected
.iter()
.any(|target| target == "zed" || target == "all")
{
let rel = ".zed/settings.json";
let target = root.join(rel);
let content = render_zed_acp_settings(root);
if target.exists() && !force {
actions.push(format!("skip existing {}", target.display()));
} else if dry_run {
actions.push(format!("dry-run write {}", target.display()));
} else {
if let Some(parent) = target.parent() {
fs::create_dir_all(parent)?;
}
fs::write(&target, content)?;
actions.push(format!("write {}", target.display()));
}
}
Ok(actions)
}
#[cfg(test)]
fn handle_request(root: &Path, request: Value) -> Vec<Value> {
let mut state = AcpConnectionState { initialized: true };
handle_request_with_state(root, request, &mut state)
}
#[cfg(test)]
fn handle_request_with_state(
root: &Path,
request: Value,
state: &mut AcpConnectionState,
) -> Vec<Value> {
handle_request_with_state_and_executor(root, request, state, &HandoffPromptExecutor)
}
fn handle_request_with_state_and_executor(
root: &Path,
request: Value,
state: &mut AcpConnectionState,
executor: &dyn AcpPromptExecutor,
) -> Vec<Value> {
let id = request.get("id").cloned().unwrap_or(Value::Null);
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return vec![structured_error(
id,
AcpErrorKind::InvalidRequest,
"request must declare jsonrpc 2.0",
"",
)];
}
let Some(method) = request.get("method").and_then(Value::as_str) else {
return vec![structured_error(
id,
AcpErrorKind::InvalidRequest,
"request method must be a string",
"",
)];
};
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
if id.is_null() {
if method == "session/cancel" {
let _ = cancel_session(root, ¶ms);
}
return Vec::new();
}
if method == "initialize" {
if state.initialized {
return vec![structured_error(
id,
AcpErrorKind::InvalidRequest,
"ACP connection is already initialized",
method,
)];
}
if params
.get("protocolVersion")
.and_then(Value::as_u64)
.filter(|version| *version <= u32::MAX as u64)
.is_none()
{
return vec![structured_error(
id,
AcpErrorKind::InvalidParams,
"initialize requires an unsigned 32-bit protocolVersion",
method,
)];
}
state.initialized = true;
return vec![jsonrpc_result(id, initialize_result(¶ms))];
}
if !state.initialized {
return vec![structured_error(
id,
AcpErrorKind::InvalidRequest,
"initialize must be called before other ACP methods",
method,
)];
}
let result = match method {
"authenticate" => Ok(vec![jsonrpc_result(id.clone(), json!({}))]),
"session/new" => create_session(root, id.clone(), ¶ms),
"session/load" => load_session(root, id.clone(), ¶ms),
"session/list" => list_sessions(root, id.clone(), ¶ms),
"session/delete" => delete_session(root, id.clone(), ¶ms),
"session/prompt" => prompt_session_with_executor(root, id.clone(), ¶ms, executor),
"session/cancel" => {
cancel_session(root, ¶ms).map(|_| vec![jsonrpc_result(id.clone(), json!({}))])
}
_ => Err(protocol_error(
AcpErrorKind::MethodNotFound,
format!("method not found: {method}"),
)),
};
match result {
Ok(responses) => responses,
Err(error) => {
let kind = error
.downcast_ref::<AcpProtocolError>()
.map(|error| error.kind)
.unwrap_or(AcpErrorKind::Internal);
vec![structured_error(id, kind, &error.to_string(), method)]
}
}
}
fn initialize_result(params: &Value) -> Value {
let requested = params.get("protocolVersion").and_then(Value::as_u64);
let negotiated = requested
.filter(|version| *version == ACP_PROTOCOL_VERSION as u64)
.unwrap_or(ACP_PROTOCOL_VERSION as u64);
json!({
"protocolVersion": negotiated,
"agentInfo": {
"name": "sdd-layer",
"version": env!("CARGO_PKG_VERSION")
},
"agentCapabilities": initialize_agent_capabilities(),
"authMethods": [],
"_meta": {
"sdd": {
"agent": "sdd-orchestrator",
"runtimeAgnostic": true,
"artifactStoreCanonical": true
}
}
})
}
fn initialize_agent_capabilities() -> Value {
json!({
"loadSession": true,
"promptCapabilities": {
"image": false,
"audio": false,
"embeddedContext": false
},
"mcpCapabilities": {
"http": false,
"sse": false
},
"sessionCapabilities": {
"list": {},
"delete": {}
},
"auth": {}
})
}
fn create_session(root: &Path, id: Value, params: &Value) -> Result<Vec<Value>> {
let cwd = canonical_cwd(&required_abs_path(params, "cwd")?)?;
let session_id = new_session_id(&cwd);
let now = now();
let session = AcpSession {
session_id: session_id.clone(),
cwd,
mcp_servers: redacted_mcp_servers(params),
orchestration_slug: None,
trace_ids: Vec::new(),
transcript: Vec::new(),
cancelled: false,
turn_sequence: 0,
last_event_sequence: 0,
created_at: now.clone(),
updated_at: now,
};
let mut session = session;
let event_payload = json!({
"cwd": session.cwd,
"mcpServers": session.mcp_servers,
});
commit_session_mutation(root, &mut session, "session_created", event_payload)?;
Ok(vec![jsonrpc_result(id, json!({ "sessionId": session_id }))])
}
fn load_session(root: &Path, id: Value, params: &Value) -> Result<Vec<Value>> {
let session_id = required_string(params, "sessionId")?;
let cwd = canonical_cwd(&required_abs_path(params, "cwd")?)?;
let mut session = read_session(root, &session_id)?;
let stored_cwd = canonical_cwd(&session.cwd)?;
if cwd != stored_cwd {
return Err(protocol_error(
AcpErrorKind::InvalidParams,
format!(
"requested cwd `{}` does not match session cwd `{}`",
cwd.display(),
stored_cwd.display()
),
));
}
session.cwd = stored_cwd;
if let Some(servers) = params.get("mcpServers").and_then(Value::as_array) {
session.mcp_servers = redact_mcp_server_array(servers);
}
session.updated_at = now();
let event_payload = json!({
"cwd": session.cwd,
"mcpServers": session.mcp_servers,
});
commit_session_mutation(root, &mut session, "session_loaded", event_payload)?;
let mut responses = replay_session_updates(root, &session.session_id)?;
if responses.is_empty() {
for entry in &session.transcript {
let update = session_update_message(
&session.session_id,
&entry.message_id,
if entry.role == "user" {
"user_message_chunk"
} else {
"agent_message_chunk"
},
&entry.text,
);
persist_session_update(root, &session.session_id, &update)?;
responses.push(update);
}
}
responses.push(jsonrpc_result(id, json!({})));
Ok(responses)
}
#[cfg(test)]
fn prompt_session(root: &Path, id: Value, params: &Value) -> Result<Vec<Value>> {
prompt_session_with_executor(root, id, params, &HandoffPromptExecutor)
}
fn prompt_session_with_executor(
root: &Path,
id: Value,
params: &Value,
executor: &dyn AcpPromptExecutor,
) -> Result<Vec<Value>> {
let runtime = AcpConnectionRuntime::default();
prompt_session_with_executor_and_sink(root, id, params, executor, None, &runtime)
}
fn prompt_session_with_executor_and_sink(
root: &Path,
id: Value,
params: &Value,
executor: &dyn AcpPromptExecutor,
sink: Option<&dyn AcpResponseSink>,
runtime: &AcpConnectionRuntime,
) -> Result<Vec<Value>> {
let session_id = required_string(params, "sessionId")?;
let prompt = extract_prompt_text(params)?;
let redacted_prompt = runtime::redaction::redact_text(&prompt);
let orchestration_slug = params
.get("orchestration")
.and_then(Value::as_str)
.map(crate::artifact_slug)
.or_else(|| Some(crate::artifact_slug(&redacted_prompt)));
let (session, turn_sequence) = mutate_session(root, &session_id, |session| {
session.orchestration_slug = orchestration_slug;
session.turn_sequence = session.turn_sequence.saturating_add(1);
let turn_sequence = session.turn_sequence;
let user_entry = AcpTranscriptEntry {
role: "user".to_string(),
message_id: message_id("user", &redacted_prompt, turn_sequence),
text: redacted_prompt.clone(),
at: now(),
};
session.transcript.push(user_entry.clone());
session.updated_at = now();
Ok((
"prompt_received",
json!({
"entry": user_entry,
"orchestrationSlug": session.orchestration_slug,
}),
turn_sequence,
))
})?;
let cancellation = runtime.cancellation_token(&session_id);
let cancel_baseline = cancellation.generation();
if session.cancelled {
cancellation.consume_through(cancellation.generation());
let responses = finish_cancelled_turn(root, id, &session_id, &redacted_prompt)?;
return emit_intermediate_responses(responses, sink);
}
let tool_call_id = format!(
"tool_{}",
short_hash(&format!("{session_id}:{turn_sequence}:{redacted_prompt}"))
);
let mut responses = Vec::new();
for update in [
session_update_plan(&session.session_id),
session_update_tool_call(
&session.session_id,
&tool_call_id,
"Preparar handoff SDD supervisionado",
"in_progress",
),
] {
publish_session_update(root, &session.session_id, update, sink, &mut responses)?;
}
let context = AcpExecutionContext {
root,
session_id: &session_id,
tool_call_id: &tool_call_id,
sink,
runtime,
cancellation: Arc::clone(&cancellation),
cancel_baseline,
};
let execution = executor.execute(&session, &redacted_prompt, &context)?;
let response_text = execution.response_text;
let agent_message_id = message_id("agent", &response_text, turn_sequence);
let cancel_generation = cancellation.pending_after(cancel_baseline);
let (session, cancelled) = mutate_session(root, &session_id, |session| {
session.updated_at = now();
if session.cancelled || cancel_generation.is_some() {
if let Some(generation) = cancel_generation {
cancellation.consume_through(generation);
}
session.cancelled = false;
return Ok((
"cancellation_observed",
json!({ "stopReason": "cancelled" }),
true,
));
}
let agent_entry = AcpTranscriptEntry {
role: "agent".to_string(),
message_id: agent_message_id.clone(),
text: response_text.clone(),
at: now(),
};
session.transcript.push(agent_entry.clone());
Ok(("agent_message", json!({ "entry": agent_entry }), false))
})?;
if cancelled {
let cancelled = cancelled_turn_responses(root, id, &session_id, &redacted_prompt)?;
let mut cancelled = emit_intermediate_responses(cancelled, sink)?;
responses.append(&mut cancelled);
return Ok(responses);
}
let used = token_estimate(&redacted_prompt).saturating_add(token_estimate(&response_text));
let completion_text = format!(
"{} concluÃdo; execução real ainda depende do executor compartilhado.",
execution.tool_title
);
let updates = [
session_update_message(
&session.session_id,
&agent_message_id,
"agent_message_chunk",
&response_text,
),
session_update_tool_call_update(
&session.session_id,
&tool_call_id,
"completed",
&completion_text,
),
session_update_usage(&session.session_id, used, 128_000),
];
for update in updates {
publish_session_update(root, &session.session_id, update, sink, &mut responses)?;
}
responses.push(jsonrpc_result(id, json!({ "stopReason": "end_turn" })));
Ok(responses)
}
fn publish_session_update(
root: &Path,
session_id: &str,
update: Value,
sink: Option<&dyn AcpResponseSink>,
collected: &mut Vec<Value>,
) -> Result<()> {
persist_session_update(root, session_id, &update)?;
if let Some(sink) = sink {
sink.emit(update)
} else {
collected.push(update);
Ok(())
}
}
fn emit_intermediate_responses(
mut responses: Vec<Value>,
sink: Option<&dyn AcpResponseSink>,
) -> Result<Vec<Value>> {
let Some(sink) = sink else {
return Ok(responses);
};
if responses.len() <= 1 {
return Ok(responses);
}
let final_response = responses.pop().unwrap();
for response in responses {
sink.emit(response)?;
}
Ok(vec![final_response])
}
fn cancel_session(root: &Path, params: &Value) -> Result<()> {
let session_id = required_string(params, "sessionId")?;
mutate_session(root, &session_id, |session| {
session.cancelled = true;
session.updated_at = now();
Ok(("cancel_requested", json!({ "cooperative": true }), ()))
})
.map(|_| ())
}
fn cancel_session_with_runtime(
root: &Path,
params: &Value,
runtime: &AcpConnectionRuntime,
) -> Result<()> {
let session_id = required_string(params, "sessionId")?;
let cancellation = runtime.cancellation_token(&session_id);
let generation = cancellation.request();
runtime.cancel_permissions(&session_id);
mutate_session(root, &session_id, |session| {
let already_observed = cancellation.was_consumed(generation);
session.cancelled = !already_observed;
session.updated_at = now();
Ok((
"cancel_requested",
json!({
"cooperative": true,
"generation": generation,
"alreadyObserved": already_observed,
}),
(),
))
})
.map(|_| ())
}
fn finish_cancelled_turn(
root: &Path,
id: Value,
session_id: &str,
prompt: &str,
) -> Result<Vec<Value>> {
mutate_session(root, session_id, |session| {
session.cancelled = false;
session.updated_at = now();
Ok((
"cancellation_observed",
json!({ "stopReason": "cancelled" }),
(),
))
})?;
cancelled_turn_responses(root, id, session_id, prompt)
}
fn cancelled_turn_responses(
root: &Path,
id: Value,
session_id: &str,
prompt: &str,
) -> Result<Vec<Value>> {
let usage = session_update_usage(session_id, token_estimate(prompt), 128_000);
persist_session_update(root, session_id, &usage)?;
Ok(vec![
usage,
jsonrpc_result(id, json!({ "stopReason": "cancelled" })),
])
}
fn list_sessions(root: &Path, id: Value, params: &Value) -> Result<Vec<Value>> {
let cwd_filter = params
.get("cwd")
.map(|_| required_abs_path(params, "cwd"))
.transpose()?
.map(|path| canonical_cwd(&path))
.transpose()?;
let cursor = params.get("cursor").and_then(Value::as_str);
let directory = sessions_dir(root);
if !directory.exists() {
return Ok(vec![jsonrpc_result(id, json!({ "sessions": [] }))]);
}
let mut seen = HashSet::new();
let mut sessions = Vec::new();
for entry in fs::read_dir(&directory)
.with_context(|| format!("reading ACP sessions from {}", directory.display()))?
{
let entry = entry?;
let file_type = entry.file_type()?;
let name = entry.file_name().to_string_lossy().to_string();
let session_id = if file_type.is_dir() {
name
} else if file_type.is_file() && name.ends_with(".json") {
name.trim_end_matches(".json").to_string()
} else {
continue;
};
if !seen.insert(session_id.clone()) {
continue;
}
let Ok(session) = read_session(root, &session_id) else {
continue;
};
if let Some(filter) = cwd_filter.as_ref() {
let Ok(session_cwd) = canonical_cwd(&session.cwd) else {
continue;
};
if filter != &session_cwd {
continue;
}
}
sessions.push(session);
}
sessions.sort_by(|left, right| {
right
.updated_at
.cmp(&left.updated_at)
.then_with(|| left.session_id.cmp(&right.session_id))
});
let start = match cursor {
Some(cursor) => sessions
.iter()
.position(|session| session.session_id == cursor)
.map(|index| index + 1)
.ok_or_else(|| {
protocol_error(
AcpErrorKind::InvalidParams,
format!("invalid session/list cursor `{cursor}`"),
)
})?,
None => 0,
};
const PAGE_SIZE: usize = 100;
let page = sessions
.iter()
.skip(start)
.take(PAGE_SIZE)
.collect::<Vec<_>>();
let next_cursor = (start + page.len() < sessions.len())
.then(|| page.last().map(|session| session.session_id.clone()))
.flatten();
let session_infos = page
.into_iter()
.map(|session| {
json!({
"sessionId": session.session_id,
"cwd": session.cwd,
"title": session.orchestration_slug,
"updatedAt": session.updated_at,
"_meta": {
"sdd": {
"derived": true,
"traceIds": session.trace_ids,
}
}
})
})
.collect::<Vec<_>>();
let mut result = json!({ "sessions": session_infos });
if let Some(next_cursor) = next_cursor {
result["nextCursor"] = Value::String(next_cursor);
}
Ok(vec![jsonrpc_result(id, result)])
}
fn delete_session(root: &Path, id: Value, params: &Value) -> Result<Vec<Value>> {
let session_id = required_string(params, "sessionId")?;
let directory = session_dir(root, &session_id)?;
let legacy = legacy_session_path(root, &session_id)?;
let mut deleted = false;
if directory.exists() {
fs::remove_dir_all(&directory)
.with_context(|| format!("deleting ACP session {}", directory.display()))?;
deleted = true;
}
if legacy.exists() {
fs::remove_file(&legacy)
.with_context(|| format!("deleting legacy ACP session {}", legacy.display()))?;
deleted = true;
}
if !deleted {
return Err(protocol_error(
AcpErrorKind::ResourceNotFound,
format!("ACP session `{session_id}` was not found"),
));
}
Ok(vec![jsonrpc_result(id, json!({}))])
}
fn build_agent_response(session: &AcpSession, prompt: &str) -> String {
let slug = session
.orchestration_slug
.as_deref()
.unwrap_or("sdd-orchestration");
format!(
"Recebi a demanda para o agente `sdd-orchestrator`.\n\nPróximo caminho seguro:\n\n```bash\nsdd init \"{slug}\"\nsdd workflow run loop --input {:?} --max-iterations 3 --json\n```\n\nVou tratar ACP como superfÃcie de sessão e manter o artifact store em `docs/{slug}/` como fonte canônica. Nenhuma escrita de código ou aprovação de checkpoint foi executada por este turno ACP.",
prompt
)
}
fn build_agent_response_with_context(
session: &AcpSession,
prompt: &str,
snapshot: &crate::context_service::ContextSnapshot,
) -> String {
let mut response = build_agent_response(session, prompt);
let context = truncate_utf8(&snapshot.content, 4_096);
response.push_str(&format!(
"\n\n## Contexto SDD compartilhado\n\n- Fontes ranqueadas: {}\n- Conflitos estruturados: {}\n- Caracteres considerados: {}\n- Resource derivado: `{}`\n\n{}",
snapshot.sources_included,
snapshot.conflicts,
snapshot.used_chars,
snapshot.path.display(),
context
));
response
}
fn truncate_utf8(value: &str, max_chars: usize) -> String {
if value.chars().count() <= max_chars {
return value.to_string();
}
let mut truncated = value.chars().take(max_chars).collect::<String>();
truncated.push_str("\n...(contexto truncado)\n");
truncated
}
fn session_update_plan(session_id: &str) -> Value {
json!({
"jsonrpc": "2.0",
"method": "session/update",
"params": {
"sessionId": session_id,
"update": {
"sessionUpdate": "plan",
"entries": [
{ "content": "Resolver demanda para um ciclo SDD rastreável", "priority": "high", "status": "completed" },
{ "content": "Usar artifact store local como fonte canônica", "priority": "high", "status": "pending" },
{ "content": "Parar em checkpoints humanos", "priority": "high", "status": "pending" }
]
}
}
})
}
fn session_update_message(
session_id: &str,
message_id: &str,
session_update: &str,
text: &str,
) -> Value {
json!({
"jsonrpc": "2.0",
"method": "session/update",
"params": {
"sessionId": session_id,
"update": {
"sessionUpdate": session_update,
"messageId": message_id,
"content": {
"type": "text",
"text": text
}
}
}
})
}
fn session_update_tool_call(
session_id: &str,
tool_call_id: &str,
title: &str,
status: &str,
) -> Value {
session_update(
session_id,
json!({
"sessionUpdate": "tool_call",
"toolCallId": tool_call_id,
"title": title,
"kind": "think",
"status": status,
"rawInput": {
"executor": "sdd-supervised-handoff"
}
}),
)
}
fn session_update_tool_call_update(
session_id: &str,
tool_call_id: &str,
status: &str,
text: &str,
) -> Value {
session_update(
session_id,
json!({
"sessionUpdate": "tool_call_update",
"toolCallId": tool_call_id,
"status": status,
"content": [{
"type": "content",
"content": {
"type": "text",
"text": text
}
}]
}),
)
}
fn session_update_usage(session_id: &str, used: u64, size: u64) -> Value {
session_update(
session_id,
json!({
"sessionUpdate": "usage_update",
"used": used,
"size": size,
}),
)
}
fn session_update(session_id: &str, update: Value) -> Value {
json!({
"jsonrpc": "2.0",
"method": "session/update",
"params": {
"sessionId": session_id,
"update": update,
}
})
}
fn token_estimate(text: &str) -> u64 {
text.split_whitespace().count().max(1) as u64
}
fn extract_prompt_text(params: &Value) -> Result<String> {
if let Some(text) = params.get("prompt").and_then(Value::as_str) {
return Ok(text.to_string());
}
if let Some(text) = params.get("message").and_then(Value::as_str) {
return Ok(text.to_string());
}
if let Some(prompt) = params.get("prompt").and_then(Value::as_array) {
let mut parts = Vec::new();
for block in prompt {
if let Some(text) = block.get("text").and_then(Value::as_str) {
parts.push(text.to_string());
continue;
}
if let Some(text) = block
.get("resource")
.and_then(|resource| resource.get("text"))
.and_then(Value::as_str)
{
parts.push(text.to_string());
continue;
}
if let Some(link) = resource_link_summary(block) {
parts.push(link);
}
}
if !parts.is_empty() {
return Ok(parts.join("\n\n"));
}
}
Err(protocol_error(
AcpErrorKind::InvalidParams,
"session/prompt requires a text prompt",
))
}
fn resource_link_summary(block: &Value) -> Option<String> {
let resource_link = if block.get("type").and_then(Value::as_str) == Some("resource_link") {
Some(block)
} else {
block
.get("resourceLink")
.or_else(|| block.get("resource_link"))
}?;
let uri = resource_link.get("uri").and_then(Value::as_str)?;
let name = resource_link
.get("name")
.and_then(Value::as_str)
.or_else(|| resource_link.get("title").and_then(Value::as_str))
.unwrap_or("resource");
let mime_type = resource_link
.get("mimeType")
.and_then(Value::as_str)
.or_else(|| resource_link.get("mime_type").and_then(Value::as_str));
Some(match mime_type {
Some(mime_type) => format!("ResourceLink: {name} ({mime_type}) <{uri}>"),
None => format!("ResourceLink: {name} <{uri}>"),
})
}
fn redacted_mcp_servers(params: &Value) -> Vec<Value> {
params
.get("mcpServers")
.and_then(Value::as_array)
.map(|servers| redact_mcp_server_array(servers))
.unwrap_or_default()
}
fn redact_mcp_server_array(servers: &[Value]) -> Vec<Value> {
servers.iter().map(redact_mcp_server_value).collect()
}
fn redact_mcp_server_value(value: &Value) -> Value {
match value {
Value::Array(items) => Value::Array(items.iter().map(redact_mcp_server_value).collect()),
Value::Object(map) => {
let mut redacted = serde_json::Map::new();
for (key, value) in map {
if is_sensitive_mcp_value_key(key) {
redacted.insert(key.clone(), Value::String(REDACTED.to_string()));
} else {
redacted.insert(key.clone(), redact_mcp_server_value(value));
}
}
Value::Object(redacted)
}
Value::String(text) => Value::String(runtime::redaction::redact_text(text)),
other => other.clone(),
}
}
fn is_sensitive_mcp_value_key(key: &str) -> bool {
matches!(
key.to_ascii_lowercase().as_str(),
"value" | "authorization" | "token" | "secret" | "apikey" | "api_key" | "password"
)
}
fn required_abs_path(params: &Value, key: &str) -> Result<PathBuf> {
let value = required_string(params, key)?;
let path = PathBuf::from(&value);
if !path.is_absolute() {
return Err(protocol_error(
AcpErrorKind::InvalidParams,
format!("{key} must be an absolute path"),
));
}
Ok(path)
}
fn required_string(params: &Value, key: &str) -> Result<String> {
params
.get(key)
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string)
.ok_or_else(|| {
protocol_error(
AcpErrorKind::InvalidParams,
format!("missing required field `{key}`"),
)
})
}
fn canonical_cwd(path: &Path) -> Result<PathBuf> {
fs::canonicalize(path).map_err(|error| {
protocol_error(
AcpErrorKind::InvalidParams,
format!("cwd `{}` cannot be canonicalized: {error}", path.display()),
)
})
}
fn read_session(root: &Path, session_id: &str) -> Result<AcpSession> {
let events_path = events_path(root, session_id)?;
let lock = session_lock(&events_path);
let _guard = lock
.lock()
.map_err(|_| anyhow!("ACP session lock poisoned for {}", events_path.display()))?;
read_session_unlocked(root, session_id)
}
fn read_session_unlocked(root: &Path, session_id: &str) -> Result<AcpSession> {
let path = session_path(root, session_id)?;
let events = read_session_events_unlocked(&events_path(root, session_id)?)?;
let snapshot = latest_session_snapshot(&events);
if let Ok(text) = fs::read_to_string(&path) {
if let Ok(mut session) = serde_json::from_str::<AcpSession>(&text) {
if let Some(recovered) = snapshot
.filter(|recovered| recovered.last_event_sequence > session.last_event_sequence)
{
session = recovered;
save_session(root, &session)?;
}
return Ok(session);
}
}
if let Some(session) = snapshot {
save_session(root, &session)?;
return Ok(session);
}
let legacy = legacy_session_path(root, session_id)?;
if !legacy.is_file() {
return Err(protocol_error(
AcpErrorKind::ResourceNotFound,
format!("ACP session `{session_id}` was not found"),
));
}
let text =
fs::read_to_string(&legacy).with_context(|| format!("reading {}", legacy.display()))?;
serde_json::from_str(&text).with_context(|| format!("parsing ACP session {}", legacy.display()))
}
fn latest_session_snapshot(events: &[AcpEvent]) -> Option<AcpSession> {
events.iter().rev().find_map(|event| {
let mut session =
serde_json::from_value::<AcpSession>(event.payload.get("sessionSnapshot")?.clone())
.ok()?;
session.last_event_sequence = event.sequence;
Some(session)
})
}
fn save_session(root: &Path, session: &AcpSession) -> Result<()> {
let path = session_path(root, &session.session_id)?;
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
atomic_write_json(&path, session)
}
fn session_path(root: &Path, session_id: &str) -> Result<PathBuf> {
Ok(session_dir(root, session_id)?.join("meta.json"))
}
fn legacy_session_path(root: &Path, session_id: &str) -> Result<PathBuf> {
validate_session_id(session_id)?;
Ok(sessions_dir(root).join(format!("{session_id}.json")))
}
fn session_dir(root: &Path, session_id: &str) -> Result<PathBuf> {
validate_session_id(session_id)?;
Ok(sessions_dir(root).join(session_id))
}
fn events_path(root: &Path, session_id: &str) -> Result<PathBuf> {
Ok(session_dir(root, session_id)?.join("events.jsonl"))
}
fn validate_session_id(session_id: &str) -> Result<()> {
if session_id.is_empty()
|| session_id.contains('/')
|| session_id.contains('\\')
|| session_id.contains("..")
{
return Err(protocol_error(
AcpErrorKind::InvalidParams,
format!("invalid ACP session id `{session_id}`"),
));
}
Ok(())
}
fn atomic_write_json<T: Serialize>(path: &Path, value: &T) -> Result<()> {
let parent = path
.parent()
.ok_or_else(|| anyhow!("ACP metadata path has no parent: {}", path.display()))?;
fs::create_dir_all(parent)?;
let counter = SESSION_COUNTER.fetch_add(1, Ordering::Relaxed);
let temporary = parent.join(format!(".meta.json.tmp.{}.{}", std::process::id(), counter));
let bytes = format!("{}\n", serde_json::to_string_pretty(value)?);
let result = (|| -> Result<()> {
let mut file = OpenOptions::new()
.create_new(true)
.write(true)
.open(&temporary)?;
file.write_all(bytes.as_bytes())?;
file.sync_all()?;
fs::rename(&temporary, path)?;
File::open(parent)?.sync_all()?;
Ok(())
})();
if result.is_err() {
let _ = fs::remove_file(&temporary);
}
result
}
fn commit_session_mutation(
root: &Path,
session: &mut AcpSession,
kind: &str,
payload: Value,
) -> Result<()> {
let path = events_path(root, &session.session_id)?;
let lock = session_lock(&path);
let _guard = lock
.lock()
.map_err(|_| anyhow!("ACP session lock poisoned for {}", path.display()))?;
commit_session_mutation_unlocked(root, session, kind, payload, &path)
}
fn commit_session_mutation_unlocked(
root: &Path,
session: &mut AcpSession,
kind: &str,
payload: Value,
path: &Path,
) -> Result<()> {
let sequence = next_event_sequence_unlocked(path)?;
session.last_event_sequence = sequence;
let mut payload = match redact_event_value(payload) {
Value::Object(map) => map,
value => serde_json::Map::from_iter([("value".to_string(), value)]),
};
payload.insert(
"sessionSnapshot".to_string(),
redact_event_value(serde_json::to_value(&*session)?),
);
let event = AcpEvent {
sequence,
kind: kind.to_string(),
at: now(),
payload: Value::Object(payload),
};
append_event_record_unlocked(path, &event)?;
save_session(root, session)
}
fn mutate_session<T>(
root: &Path,
session_id: &str,
mutation: impl FnOnce(&mut AcpSession) -> Result<(&'static str, Value, T)>,
) -> Result<(AcpSession, T)> {
let path = events_path(root, session_id)?;
let lock = session_lock(&path);
let _guard = lock
.lock()
.map_err(|_| anyhow!("ACP session lock poisoned for {}", path.display()))?;
let mut session = read_session_unlocked(root, session_id)?;
let (kind, payload, output) = mutation(&mut session)?;
commit_session_mutation_unlocked(root, &mut session, kind, payload, &path)?;
Ok((session, output))
}
fn append_session_event(root: &Path, session_id: &str, kind: &str, payload: Value) -> Result<()> {
let path = events_path(root, session_id)?;
let lock = session_lock(&path);
let _guard = lock
.lock()
.map_err(|_| anyhow!("ACP session lock poisoned for {}", path.display()))?;
let event = AcpEvent {
sequence: next_event_sequence_unlocked(&path)?,
kind: kind.to_string(),
at: now(),
payload: redact_event_value(payload),
};
append_event_record_unlocked(&path, &event)
}
fn session_lock(path: &Path) -> Arc<Mutex<()>> {
let locks = SESSION_LOCKS.get_or_init(|| Mutex::new(HashMap::new()));
let mut locks = locks
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
Arc::clone(
locks
.entry(path.to_path_buf())
.or_insert_with(|| Arc::new(Mutex::new(()))),
)
}
fn append_event_record_unlocked(path: &Path, event: &AcpEvent) -> Result<()> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let mut bytes = serde_json::to_vec(event)?;
bytes.push(b'\n');
let mut file = OpenOptions::new().create(true).append(true).open(path)?;
file.write_all(&bytes)?;
file.sync_data()?;
Ok(())
}
fn next_event_sequence_unlocked(path: &Path) -> Result<u64> {
let Some(last) = read_session_events_unlocked(path)?.last().cloned() else {
return Ok(1);
};
Ok(last.sequence.saturating_add(1))
}
fn read_session_events(path: &Path) -> Result<Vec<AcpEvent>> {
let lock = session_lock(path);
let _guard = lock
.lock()
.map_err(|_| anyhow!("ACP session lock poisoned for {}", path.display()))?;
read_session_events_unlocked(path)
}
fn read_session_events_unlocked(path: &Path) -> Result<Vec<AcpEvent>> {
if !path.is_file() {
return Ok(Vec::new());
}
let bytes = fs::read(path)?;
let complete_len = bytes
.iter()
.rposition(|byte| *byte == b'\n')
.map(|index| index + 1)
.unwrap_or(0);
let mut events: Vec<AcpEvent> = Vec::new();
for (index, line) in bytes[..complete_len]
.split(|byte| *byte == b'\n')
.enumerate()
{
if line.iter().all(u8::is_ascii_whitespace) {
continue;
}
let event = serde_json::from_slice::<AcpEvent>(line)
.with_context(|| format!("parsing ACP event {}:{}", path.display(), index + 1))?;
validate_event_sequence(path, &events, &event)?;
events.push(event);
}
let tail = &bytes[complete_len..];
if !tail.iter().all(u8::is_ascii_whitespace) {
match serde_json::from_slice::<AcpEvent>(tail) {
Ok(event) => {
validate_event_sequence(path, &events, &event)?;
OpenOptions::new()
.append(true)
.open(path)?
.write_all(b"\n")?;
events.push(event);
}
Err(_) => {
let file = OpenOptions::new().write(true).open(path)?;
file.set_len(complete_len as u64)?;
file.sync_data()?;
}
}
}
Ok(events)
}
fn validate_event_sequence(path: &Path, events: &[AcpEvent], event: &AcpEvent) -> Result<()> {
if events
.last()
.is_some_and(|previous| event.sequence <= previous.sequence)
{
return Err(anyhow!(
"ACP events are not strictly ordered in {}",
path.display()
));
}
Ok(())
}
fn persist_session_update(root: &Path, session_id: &str, response: &Value) -> Result<()> {
let update = response
.pointer("/params/update")
.cloned()
.ok_or_else(|| anyhow!("session update response is missing params.update"))?;
append_session_event(
root,
session_id,
"session_update",
json!({ "update": update }),
)
}
fn replay_session_updates(root: &Path, session_id: &str) -> Result<Vec<Value>> {
let events = read_session_events(&events_path(root, session_id)?)?;
Ok(events
.into_iter()
.filter(|event| event.kind == "session_update")
.filter_map(|event| event.payload.get("update").cloned())
.map(|update| session_update(session_id, update))
.collect())
}
fn redact_event_value(value: Value) -> Value {
match value {
Value::String(text) => Value::String(runtime::redaction::redact_text(&text)),
Value::Array(items) => Value::Array(items.into_iter().map(redact_event_value).collect()),
Value::Object(map) => Value::Object(
map.into_iter()
.map(|(key, value)| {
if is_sensitive_mcp_value_key(&key) {
(key, Value::String(REDACTED.to_string()))
} else {
(key, redact_event_value(value))
}
})
.collect(),
),
other => other,
}
}
fn sessions_dir(root: &Path) -> PathBuf {
root.join(".sdd").join("acp").join("sessions")
}
fn new_session_id(cwd: &Path) -> String {
let counter = SESSION_COUNTER.fetch_add(1, Ordering::Relaxed);
let seed = format!(
"{}:{}:{}:{}",
now(),
std::process::id(),
counter,
cwd.display()
);
format!("sess_{}", short_hash(&seed))
}
fn message_id(role: &str, text: &str, turn_sequence: u64) -> String {
format!(
"msg_{role}_{}",
short_hash(&format!("{turn_sequence}:{text}"))
)
}
fn short_hash(input: &str) -> String {
let mut hasher = Sha256::new();
hasher.update(input.as_bytes());
let digest = hasher.finalize();
digest[..8]
.iter()
.map(|byte| format!("{byte:02x}"))
.collect()
}
fn jsonrpc_result(id: Value, result: Value) -> Value {
json!({ "jsonrpc": "2.0", "id": id, "result": result })
}
fn jsonrpc_error(id: Value, code: i32, message: &str) -> Value {
json!({
"jsonrpc": "2.0",
"id": id,
"error": { "code": code, "message": message }
})
}
fn structured_error(id: Value, kind: AcpErrorKind, message: &str, method: &str) -> Value {
json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": kind.code(),
"message": runtime::redaction::redact_text(message),
"data": {
"kind": kind.label(),
"method": method,
"retryable": false,
}
}
})
}
fn protocol_error(kind: AcpErrorKind, message: impl Into<String>) -> anyhow::Error {
AcpProtocolError {
kind,
message: message.into(),
}
.into()
}
fn parse_targets(targets: &str) -> Result<Vec<String>> {
let mut parsed = Vec::new();
for raw in targets.split(',') {
let item = raw.trim().to_ascii_lowercase();
if item.is_empty() {
continue;
}
match item.as_str() {
"all" | "zed" => parsed.push(item),
other => bail!("unsupported ACP config target `{other}` (use zed|all)"),
}
}
if parsed.is_empty() {
parsed.push("all".to_string());
}
Ok(parsed)
}
fn render_zed_acp_settings(root: &Path) -> String {
let root = runtime::platform::display_path(root).replace('\\', "/");
serde_json::to_string_pretty(&json!({
"agent_servers": {
"sdd-layer": {
"command": "sdd",
"args": ["acp", "serve", "--root", root],
"env": {}
}
}
}))
.expect("ACP settings JSON must serialize")
}
fn now() -> String {
Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true)
}
#[cfg(feature = "acp-agent")]
fn acp_sdk_marker() -> &'static str {
std::any::type_name::<agent_client_protocol::Client>()
}
#[cfg(not(feature = "acp-agent"))]
fn acp_sdk_marker() -> &'static str {
"manual-jsonrpc-stdio (compile with acp-agent for agent-client-protocol runtime crate)"
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, AtomicUsize};
use std::sync::{Arc, Barrier};
use std::thread;
use std::time::Duration;
use tempfile::tempdir;
#[derive(Default)]
struct FlushTrackingWriter {
bytes: Vec<u8>,
flush_points: Vec<usize>,
}
impl Write for FlushTrackingWriter {
fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
self.bytes.extend_from_slice(buffer);
Ok(buffer.len())
}
fn flush(&mut self) -> io::Result<()> {
self.flush_points.push(self.bytes.len());
Ok(())
}
}
struct SlowExecutor;
impl AcpPromptExecutor for SlowExecutor {
fn execute(
&self,
session: &AcpSession,
prompt: &str,
_context: &AcpExecutionContext<'_>,
) -> Result<AcpPromptExecution> {
thread::sleep(Duration::from_millis(100));
Ok(AcpPromptExecution {
response_text: build_agent_response(session, prompt),
tool_title: "Slow supervised handoff".to_string(),
})
}
}
struct ObservedSlowExecutor {
finished: Arc<AtomicBool>,
}
impl AcpPromptExecutor for ObservedSlowExecutor {
fn execute(
&self,
session: &AcpSession,
prompt: &str,
_context: &AcpExecutionContext<'_>,
) -> Result<AcpPromptExecution> {
thread::sleep(Duration::from_millis(100));
self.finished.store(true, Ordering::Release);
Ok(AcpPromptExecution {
response_text: build_agent_response(session, prompt),
tool_title: "Observed slow handoff".to_string(),
})
}
}
struct StreamingObservationWriter {
bytes: Vec<u8>,
executor_finished: Arc<AtomicBool>,
flushes_before_executor_finished: Arc<AtomicUsize>,
}
struct PermissionExecutor {
decision: Arc<Mutex<Option<AcpPermissionDecision>>>,
}
impl AcpPromptExecutor for PermissionExecutor {
fn execute(
&self,
session: &AcpSession,
prompt: &str,
context: &AcpExecutionContext<'_>,
) -> Result<AcpPromptExecution> {
let decision = context.request_permission("Aplicar mudança supervisionada")?;
*self.decision.lock().unwrap() = Some(decision);
Ok(AcpPromptExecution {
response_text: build_agent_response(session, prompt),
tool_title: "Permission checkpoint".to_string(),
})
}
}
impl Write for StreamingObservationWriter {
fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
self.bytes.extend_from_slice(buffer);
Ok(buffer.len())
}
fn flush(&mut self) -> io::Result<()> {
if !self.executor_finished.load(Ordering::Acquire) {
self.flushes_before_executor_finished
.fetch_add(1, Ordering::Relaxed);
}
Ok(())
}
}
fn new_session(root: &Path) -> String {
let cwd = root.to_string_lossy();
create_session(root, json!(1), &json!({ "cwd": cwd, "mcpServers": [] })).unwrap()[0]
["result"]["sessionId"]
.as_str()
.unwrap()
.to_string()
}
#[test]
fn serves_initialize_and_session_prompt() {
let root = tempdir().unwrap();
let cwd = root.path().to_string_lossy();
let initialize = json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": { "protocolVersion": 1 }
});
let session_new = json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/new",
"params": { "cwd": cwd, "mcpServers": [] }
});
let input = format!("{initialize}\n{session_new}\n");
let mut output = Vec::new();
serve_lines(root.path(), input.as_bytes(), &mut output).unwrap();
let text = String::from_utf8(output).unwrap();
assert!(text.contains("\"protocolVersion\":1"));
assert!(text.contains("\"sessionId\""));
}
#[test]
fn rejects_relative_cwd() {
let root = tempdir().unwrap();
let input = b"{\"jsonrpc\":\"2.0\",\"id\":0,\"method\":\"initialize\",\"params\":{\"protocolVersion\":1}}\n{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"session/new\",\"params\":{\"cwd\":\"relative\",\"mcpServers\":[]}}\n";
let mut output = Vec::new();
serve_lines(root.path(), &input[..], &mut output).unwrap();
let text = String::from_utf8(output).unwrap();
assert!(text.contains("\"error\""));
assert!(text.contains("absolute path"));
}
#[test]
fn creates_unique_session_ids_for_rapid_new_sessions() {
let root = tempdir().unwrap();
let cwd = root.path().to_string_lossy();
let first = create_session(
root.path(),
json!(1),
&json!({ "cwd": cwd, "mcpServers": [] }),
)
.unwrap();
let second = create_session(
root.path(),
json!(2),
&json!({ "cwd": cwd, "mcpServers": [] }),
)
.unwrap();
let first_id = first[0]["result"]["sessionId"].as_str().unwrap();
let second_id = second[0]["result"]["sessionId"].as_str().unwrap();
assert_ne!(first_id, second_id);
assert!(session_path(root.path(), first_id).unwrap().is_file());
assert!(session_path(root.path(), second_id).unwrap().is_file());
}
#[test]
fn config_supports_zed_dry_run() {
let root = tempdir().unwrap();
let actions = config_actions(root.path(), "zed", false, true).unwrap();
assert_eq!(actions.len(), 1);
assert!(actions[0].contains("dry-run write"));
}
#[test]
fn prompt_persists_redacted_transcript_and_stop_reason() {
let root = tempdir().unwrap();
let cwd = root.path().to_string_lossy();
let responses = create_session(
root.path(),
json!(1),
&json!({ "cwd": cwd, "mcpServers": [{ "name": "sdd" }] }),
)
.unwrap();
let session_id = responses[0]["result"]["sessionId"]
.as_str()
.unwrap()
.to_string();
let prompt =
"implementar feature com OPENAI_API_KEY=sk-abcdefghijklmnopqrstuvwxyz1234567890";
let responses = prompt_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"prompt": [{ "type": "text", "text": prompt }]
}),
)
.unwrap();
assert_eq!(
responses.last().unwrap()["result"]["stopReason"],
"end_turn"
);
let session = read_session(root.path(), &session_id).unwrap();
let transcript = serde_json::to_string(&session.transcript).unwrap();
assert!(!transcript.contains("sk-abcdefghijklmnopqrstuvwxyz1234567890"));
assert!(transcript.contains("[REDACTED]"));
}
#[test]
fn session_persistence_redacts_mcp_server_secrets() {
let root = tempdir().unwrap();
let cwd = root.path().to_string_lossy();
let responses = create_session(
root.path(),
json!(1),
&json!({
"cwd": cwd,
"mcpServers": [{
"type": "stdio",
"name": "private-tools",
"command": "mcp-server",
"args": ["--token", "sk-abcdefghijklmnopqrstuvwxyz1234567890"],
"env": [
{ "name": "OPENAI_API_KEY", "value": "plain-secret-value" }
],
"headers": [
{ "name": "Authorization", "value": "Bearer raw-header-secret" }
]
}]
}),
)
.unwrap();
let session_id = responses[0]["result"]["sessionId"].as_str().unwrap();
let session = read_session(root.path(), session_id).unwrap();
let persisted = serde_json::to_string(&session.mcp_servers).unwrap();
assert!(persisted.contains("OPENAI_API_KEY"));
assert!(persisted.contains("Authorization"));
assert!(persisted.contains("[REDACTED]"));
assert!(!persisted.contains("plain-secret-value"));
assert!(!persisted.contains("raw-header-secret"));
assert!(!persisted.contains("sk-abcdefghijklmnopqrstuvwxyz1234567890"));
load_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"cwd": cwd,
"mcpServers": [{
"type": "stdio",
"name": "replacement",
"command": "mcp-server",
"args": [],
"env": [
{ "name": "ANTHROPIC_API_KEY", "value": "second-secret-value" }
]
}]
}),
)
.unwrap();
let session = read_session(root.path(), session_id).unwrap();
let persisted = serde_json::to_string(&session.mcp_servers).unwrap();
assert!(persisted.contains("ANTHROPIC_API_KEY"));
assert!(persisted.contains("[REDACTED]"));
assert!(!persisted.contains("second-secret-value"));
}
#[test]
fn prompt_accepts_resource_link_content_blocks() {
let root = tempdir().unwrap();
let cwd = root.path().to_string_lossy();
let responses = create_session(
root.path(),
json!(1),
&json!({ "cwd": cwd, "mcpServers": [] }),
)
.unwrap();
let session_id = responses[0]["result"]["sessionId"]
.as_str()
.unwrap()
.to_string();
let responses = prompt_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"prompt": [{
"type": "resource_link",
"uri": "sdd://agents/sdd-orchestrator",
"name": "SDD orchestrator manifest",
"mimeType": "application/json"
}]
}),
)
.unwrap();
assert_eq!(
responses.last().unwrap()["result"]["stopReason"],
"end_turn"
);
let session_id = responses[1]["params"]["sessionId"].as_str().unwrap();
let session = read_session(root.path(), session_id).unwrap();
let transcript = serde_json::to_string(&session.transcript).unwrap();
assert!(transcript.contains("ResourceLink: SDD orchestrator manifest"));
assert!(transcript.contains("sdd://agents/sdd-orchestrator"));
}
#[test]
fn initialize_advertises_supported_session_lifecycle() {
let result = initialize_result(&json!({ "protocolVersion": 1 }));
assert_eq!(result["agentCapabilities"]["loadSession"], true);
assert_eq!(
result["agentCapabilities"]["sessionCapabilities"]["list"],
json!({})
);
assert_eq!(
result["agentCapabilities"]["sessionCapabilities"]["delete"],
json!({})
);
}
#[test]
fn initialize_falls_back_to_v1_when_requested_version_is_not_supported() {
assert_eq!(
initialize_result(&json!({ "protocolVersion": 0 }))["protocolVersion"],
1
);
assert_eq!(
initialize_result(&json!({ "protocolVersion": 99 }))["protocolVersion"],
1
);
}
#[test]
fn load_rejects_a_different_canonical_cwd() {
let root = tempdir().unwrap();
let other = tempdir().unwrap();
let session_id = new_session(root.path());
let error = load_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"cwd": other.path().to_string_lossy(),
"mcpServers": []
}),
)
.unwrap_err();
assert!(error.to_string().contains("does not match session cwd"));
}
#[test]
fn repeated_prompts_receive_unique_turn_and_tool_identifiers() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let params = json!({ "sessionId": session_id, "prompt": "same prompt" });
let first = prompt_session(root.path(), json!(2), ¶ms).unwrap();
let second = prompt_session(root.path(), json!(3), ¶ms).unwrap();
assert_ne!(
first[1]["params"]["update"]["toolCallId"],
second[1]["params"]["update"]["toolCallId"]
);
assert_ne!(
first[2]["params"]["update"]["messageId"],
second[2]["params"]["update"]["messageId"]
);
}
#[test]
fn default_executor_reads_the_shared_ranked_context_service() {
let root = tempdir().unwrap();
let store = root.path().join("docs/flow");
fs::create_dir_all(&store).unwrap();
fs::write(
store.join("traceability-map.yaml"),
"orchestration:\n name: Flow\n slug: flow\nartifacts:\n prd:\n file: 02-prd.md\n state: approved\nrelated_orchestrations: []\n",
)
.unwrap();
fs::write(
store.join("02-prd.md"),
"# PRD\n\n## Resumo\n\nContexto canônico compartilhado pelo ACP.\n",
)
.unwrap();
let session_id = new_session(root.path());
let responses = prompt_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"orchestration": "flow",
"prompt": "prepare o handoff"
}),
)
.unwrap();
let text = responses
.iter()
.find(|response| response["params"]["update"]["sessionUpdate"] == "agent_message_chunk")
.and_then(|response| response["params"]["update"]["content"]["text"].as_str())
.unwrap();
assert!(text.contains("Contexto SDD compartilhado"), "{text}");
assert!(
text.contains("Contexto canônico compartilhado pelo ACP"),
"{text}"
);
}
#[test]
fn serve_lines_flushes_each_prompt_update_before_the_final_result() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let input = format!(
"{}\n{}\n",
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/prompt",
"params": { "sessionId": session_id, "prompt": "stream updates" }
})
);
let mut output = FlushTrackingWriter::default();
serve_lines(root.path(), input.as_bytes(), &mut output).unwrap();
assert_eq!(output.flush_points.len(), 7);
assert!(output
.flush_points
.windows(2)
.all(|window| window[0] < window[1]));
}
#[test]
fn serve_lines_observes_cancel_while_prompt_is_running_on_the_same_connection() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let input = format!(
"{}\n{}\n{}\n",
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/prompt",
"params": { "sessionId": session_id, "prompt": "wait for cancel" }
}),
json!({
"jsonrpc": "2.0",
"method": "session/cancel",
"params": { "sessionId": session_id }
})
);
let mut output = Vec::new();
serve_lines_with_executor(root.path(), input.as_bytes(), &mut output, &SlowExecutor)
.unwrap();
let responses = String::from_utf8(output)
.unwrap()
.lines()
.map(|line| serde_json::from_str::<Value>(line).unwrap())
.collect::<Vec<_>>();
let prompt_result = responses
.iter()
.find(|response| response.get("id") == Some(&json!(2)))
.unwrap();
assert_eq!(prompt_result["result"]["stopReason"], "cancelled");
}
#[test]
fn plan_and_tool_start_are_flushed_before_executor_completion() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let finished = Arc::new(AtomicBool::new(false));
let flushes = Arc::new(AtomicUsize::new(0));
let executor = ObservedSlowExecutor {
finished: Arc::clone(&finished),
};
let input = format!(
"{}\n{}\n",
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/prompt",
"params": { "sessionId": session_id, "prompt": "stream before execution ends" }
})
);
let mut output = StreamingObservationWriter {
bytes: Vec::new(),
executor_finished: Arc::clone(&finished),
flushes_before_executor_finished: Arc::clone(&flushes),
};
serve_lines_with_executor(root.path(), input.as_bytes(), &mut output, &executor).unwrap();
assert!(
flushes.load(Ordering::Relaxed) >= 3,
"initialize, plan and tool start must be flushed before executor completion"
);
}
#[test]
fn jsonrpc_responses_are_not_rejected_as_methodless_requests() {
let root = tempdir().unwrap();
let input = format!(
"{}\n{}\n",
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"id": "permission-1",
"result": { "outcome": { "outcome": "selected", "optionId": "allow_once" } }
})
);
let mut output = Vec::new();
serve_lines(root.path(), input.as_bytes(), &mut output).unwrap();
assert_eq!(String::from_utf8(output).unwrap().lines().count(), 1);
}
#[test]
fn permission_request_waits_for_an_explicit_client_decision_and_rejects_closed() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let decision = Arc::new(Mutex::new(None));
let executor = PermissionExecutor {
decision: Arc::clone(&decision),
};
let input = format!(
"{}\n{}\n{}\n",
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/prompt",
"params": { "sessionId": session_id, "prompt": "needs permission" }
}),
json!({
"jsonrpc": "2.0",
"id": "permission-1",
"result": {
"outcome": { "outcome": "selected", "optionId": "reject_once" }
}
})
);
let mut output = Vec::new();
serve_lines_with_executor(root.path(), input.as_bytes(), &mut output, &executor).unwrap();
let responses = String::from_utf8(output)
.unwrap()
.lines()
.map(|line| serde_json::from_str::<Value>(line).unwrap())
.collect::<Vec<_>>();
assert!(responses.iter().any(|response| {
response["method"] == "session/request_permission"
&& response["params"]["options"]
.as_array()
.unwrap()
.iter()
.any(|option| {
option["optionId"] == "reject_once" && option["kind"] == "reject_once"
})
}));
assert_eq!(
*decision.lock().unwrap(),
Some(AcpPermissionDecision::Denied)
);
#[cfg(feature = "acp-agent")]
{
use agent_client_protocol::schema::v1::RequestPermissionRequest;
let request = responses
.iter()
.find(|response| response["method"] == "session/request_permission")
.unwrap();
serde_json::from_value::<RequestPermissionRequest>(request["params"].clone()).unwrap();
}
}
#[test]
fn concurrent_event_appends_keep_a_strict_monotonic_sequence() {
const WRITERS: usize = 12;
const EVENTS_PER_WRITER: usize = 8;
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let root_path = root.path().to_path_buf();
let barrier = Arc::new(Barrier::new(WRITERS));
let mut handles = Vec::new();
for writer in 0..WRITERS {
let root_path = root_path.clone();
let session_id = session_id.clone();
let barrier = Arc::clone(&barrier);
handles.push(thread::spawn(move || {
for event in 0..EVENTS_PER_WRITER {
barrier.wait();
append_session_event(
&root_path,
&session_id,
"concurrent_test",
json!({ "writer": writer, "event": event }),
)
.unwrap();
}
}));
}
for handle in handles {
handle.join().unwrap();
}
let events = read_session_events(&events_path(root.path(), &session_id).unwrap()).unwrap();
assert_eq!(events.len(), 1 + WRITERS * EVENTS_PER_WRITER);
assert!(events
.windows(2)
.all(|pair| pair[1].sequence == pair[0].sequence + 1));
}
#[test]
fn session_meta_is_rebuilt_from_the_event_log() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
prompt_session(
root.path(),
json!(2),
&json!({ "sessionId": session_id, "prompt": "recover me" }),
)
.unwrap();
let meta = session_path(root.path(), &session_id).unwrap();
fs::remove_file(&meta).unwrap();
let recovered = read_session(root.path(), &session_id).unwrap();
assert_eq!(recovered.transcript.len(), 2);
assert_eq!(recovered.transcript[0].text, "recover me");
assert!(meta.is_file());
}
#[test]
fn truncated_event_tail_is_repaired_before_the_next_append() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let path = events_path(root.path(), &session_id).unwrap();
OpenOptions::new()
.append(true)
.open(&path)
.unwrap()
.write_all(b"{\"sequence\":")
.unwrap();
load_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"cwd": root.path().to_string_lossy(),
"mcpServers": []
}),
)
.unwrap();
let events = read_session_events(&path).unwrap();
assert!(events.len() >= 2);
assert!(events
.windows(2)
.all(|pair| pair[0].sequence < pair[1].sequence));
assert!(fs::read(&path).unwrap().ends_with(b"\n"));
}
#[test]
fn connection_lifecycle_requires_initialize_and_does_not_reply_to_notifications() {
let root = tempdir().unwrap();
let cwd = root.path().to_string_lossy();
let input = format!(
"{}\n{}\n{}\n{}\n{}\n",
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": { "cwd": cwd, "mcpServers": [] }
}),
json!({
"jsonrpc": "2.0",
"id": 2,
"method": "initialize",
"params": { "protocolVersion": "one" }
}),
json!({
"jsonrpc": "2.0",
"id": 3,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"id": 4,
"method": "initialize",
"params": { "protocolVersion": 1 }
}),
json!({
"jsonrpc": "2.0",
"method": "notifications/ignored",
"params": {}
})
);
let mut output = Vec::new();
serve_lines(root.path(), input.as_bytes(), &mut output).unwrap();
let responses = String::from_utf8(output)
.unwrap()
.lines()
.map(|line| serde_json::from_str::<Value>(line).unwrap())
.collect::<Vec<_>>();
assert_eq!(responses.len(), 4);
assert_eq!(responses[0]["error"]["data"]["kind"], "invalid_request");
assert_eq!(responses[1]["error"]["data"]["kind"], "invalid_params");
assert_eq!(responses[2]["result"]["protocolVersion"], 1);
assert_eq!(responses[3]["error"]["data"]["kind"], "invalid_request");
}
#[test]
fn persists_atomic_meta_and_redacted_append_only_events_then_replays_updates() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let session_dir = root.path().join(".sdd/acp/sessions").join(&session_id);
let meta_path = session_dir.join("meta.json");
let events_path = session_dir.join("events.jsonl");
assert!(meta_path.is_file());
assert!(events_path.is_file());
assert!(!root
.path()
.join(".sdd/acp/sessions")
.join(format!("{session_id}.json"))
.exists());
let responses = prompt_session(
root.path(),
json!(2),
&json!({
"sessionId": session_id,
"prompt": "Use OPENAI_API_KEY=sk-abcdefghijklmnopqrstuvwxyz1234567890"
}),
)
.unwrap();
let updates = responses[..responses.len() - 1].to_vec();
let before_cancel = fs::read_to_string(&events_path).unwrap();
assert!(before_cancel.contains("session_update"));
assert!(before_cancel.contains(REDACTED));
assert!(!before_cancel.contains("sk-abcdefghijklmnopqrstuvwxyz1234567890"));
cancel_session(root.path(), &json!({ "sessionId": session_id })).unwrap();
let after_cancel = fs::read_to_string(&events_path).unwrap();
assert!(after_cancel.starts_with(&before_cancel));
assert!(after_cancel.len() > before_cancel.len());
assert!(after_cancel.contains("cancel_requested"));
assert!(fs::read_dir(&session_dir).unwrap().all(|entry| !entry
.unwrap()
.file_name()
.to_string_lossy()
.contains(".tmp")));
let replay = load_session(
root.path(),
json!(3),
&json!({
"sessionId": session_id,
"cwd": root.path().to_string_lossy(),
"mcpServers": []
}),
)
.unwrap();
let replay_updates = replay[..replay.len() - 1].to_vec();
assert_eq!(replay_updates, updates);
}
#[test]
fn lists_new_and_legacy_sessions_and_deletes_both_formats() {
let root = tempdir().unwrap();
let new_id = new_session(root.path());
let legacy_id = "sess_legacy";
let now = now();
let legacy = AcpSession {
session_id: legacy_id.to_string(),
cwd: root.path().to_path_buf(),
mcp_servers: Vec::new(),
orchestration_slug: Some("legacy-title".to_string()),
trace_ids: Vec::new(),
transcript: Vec::new(),
cancelled: false,
turn_sequence: 0,
last_event_sequence: 0,
created_at: now.clone(),
updated_at: now,
};
let legacy_path = root
.path()
.join(".sdd/acp/sessions")
.join(format!("{legacy_id}.json"));
fs::write(
&legacy_path,
format!("{}\n", serde_json::to_string_pretty(&legacy).unwrap()),
)
.unwrap();
let listed = handle_request(
root.path(),
json!({
"jsonrpc": "2.0",
"id": 10,
"method": "session/list",
"params": { "cwd": root.path().to_string_lossy() }
}),
);
let ids = listed[0]["result"]["sessions"]
.as_array()
.unwrap()
.iter()
.map(|session| session["sessionId"].as_str().unwrap())
.collect::<Vec<_>>();
assert!(ids.contains(&new_id.as_str()));
assert!(ids.contains(&legacy_id));
for session_id in [new_id.as_str(), legacy_id] {
let deleted = handle_request(
root.path(),
json!({
"jsonrpc": "2.0",
"id": 11,
"method": "session/delete",
"params": { "sessionId": session_id }
}),
);
assert_eq!(deleted[0]["result"], json!({}));
}
assert!(!legacy_path.exists());
assert!(!root.path().join(".sdd/acp/sessions").join(new_id).exists());
}
#[test]
fn cancellation_is_observed_as_a_cancelled_prompt_turn() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
cancel_session(root.path(), &json!({ "sessionId": session_id })).unwrap();
let responses = prompt_session(
root.path(),
json!(2),
&json!({ "sessionId": session_id, "prompt": "continue" }),
)
.unwrap();
assert_eq!(
responses.last().unwrap()["result"]["stopReason"],
"cancelled"
);
let session = read_session(root.path(), &session_id).unwrap();
assert!(
!session.cancelled,
"o cancelamento deve ser consumido pelo turno"
);
let events = fs::read_to_string(
root.path()
.join(".sdd/acp/sessions")
.join(session_id)
.join("events.jsonl"),
)
.unwrap();
assert!(events.contains("cancellation_observed"));
}
#[test]
fn prompt_emits_protocol_compatible_plan_message_tool_and_usage_updates() {
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let responses = prompt_session(
root.path(),
json!(2),
&json!({ "sessionId": session_id, "prompt": "prepare handoff" }),
)
.unwrap();
let kinds = responses[..responses.len() - 1]
.iter()
.map(|response| {
response["params"]["update"]["sessionUpdate"]
.as_str()
.unwrap()
})
.collect::<Vec<_>>();
assert_eq!(
kinds,
vec![
"plan",
"tool_call",
"agent_message_chunk",
"tool_call_update",
"usage_update"
]
);
assert_eq!(responses[1]["params"]["update"]["status"], "in_progress");
assert_eq!(responses[3]["params"]["update"]["status"], "completed");
assert!(responses[4]["params"]["update"]["size"].as_u64().unwrap() > 0);
}
#[cfg(feature = "acp-agent")]
#[test]
fn emitted_updates_deserialize_with_agent_client_protocol_v1() {
use agent_client_protocol::schema::v1::SessionNotification;
let root = tempdir().unwrap();
let session_id = new_session(root.path());
let responses = prompt_session(
root.path(),
json!(2),
&json!({ "sessionId": session_id, "prompt": "prepare typed handoff" }),
)
.unwrap();
for response in &responses[..responses.len() - 1] {
serde_json::from_value::<SessionNotification>(response["params"].clone()).unwrap();
}
}
#[test]
fn returns_structured_protocol_errors() {
let root = tempdir().unwrap();
let unknown = handle_request(
root.path(),
json!({ "jsonrpc": "2.0", "id": 1, "method": "unknown", "params": {} }),
);
assert_eq!(unknown[0]["error"]["code"], -32601);
assert_eq!(unknown[0]["error"]["data"]["kind"], "method_not_found");
let invalid = handle_request(
root.path(),
json!({ "jsonrpc": "2.0", "id": 2, "method": "session/new", "params": {} }),
);
assert_eq!(invalid[0]["error"]["code"], -32602);
assert_eq!(invalid[0]["error"]["data"]["kind"], "invalid_params");
let missing = handle_request(
root.path(),
json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/load",
"params": {
"sessionId": "sess_missing",
"cwd": root.path().to_string_lossy()
}
}),
);
assert_eq!(missing[0]["error"]["code"], -32002);
assert_eq!(missing[0]["error"]["data"]["kind"], "resource_not_found");
}
}