#![recursion_limit = "256"]
use std::sync::Arc;
use agent_client_protocol as acp;
use tempfile::TempDir;
use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
use zeph_acp::{AcpServerConfig, AgentSpawner, serve_connection};
use zeph_core::channel::Channel as _;
fn noop_spawner() -> AgentSpawner {
Arc::new(|channel, _ctx, _session| {
Box::pin(async move {
drop(channel);
})
})
}
fn echo_spawner() -> AgentSpawner {
Arc::new(|mut channel, _ctx, _session| {
Box::pin(async move {
let _ = channel.recv().await;
let _ = channel.flush_chunks().await;
})
})
}
fn multi_turn_echo_spawner() -> AgentSpawner {
Arc::new(|mut channel, _ctx, _session| {
Box::pin(async move {
while let Ok(Some(_)) = channel.recv().await {
if channel.flush_chunks().await.is_err() {
break;
}
}
})
})
}
fn hanging_spawner(
started: Arc<tokio::sync::Notify>,
alive: Arc<std::sync::atomic::AtomicBool>,
) -> AgentSpawner {
Arc::new(move |_channel, _ctx, _session| {
let started = Arc::clone(&started);
let alive = Arc::clone(&alive);
Box::pin(async move {
struct ClearOnDrop(Arc<std::sync::atomic::AtomicBool>);
impl Drop for ClearOnDrop {
fn drop(&mut self) {
self.0.store(false, std::sync::atomic::Ordering::SeqCst);
}
}
alive.store(true, std::sync::atomic::Ordering::SeqCst);
let _clear_on_drop = ClearOnDrop(alive);
started.notify_one();
std::future::pending::<()>().await;
})
})
}
fn text_chunks_spawner(chunks: Vec<&'static str>) -> AgentSpawner {
Arc::new(move |mut channel, _ctx, _session| {
let chunks = chunks.clone();
Box::pin(async move {
let _ = channel.recv().await;
for chunk in chunks {
let _ = channel.send_chunk(chunk).await;
}
let _ = channel.flush_chunks().await;
})
})
}
fn permission_gated_spawner(decision_tx: tokio::sync::mpsc::UnboundedSender<bool>) -> AgentSpawner {
Arc::new(move |mut channel, ctx, session| {
let decision_tx = decision_tx.clone();
Box::pin(async move {
let _ = channel.recv().await;
let gate = ctx
.expect("AcpContext must be present")
.permission_gate
.expect("permission gate must be present");
let tool_call = acp::schema::v1::ToolCallUpdate::new(
"tc-perm-1".to_owned(),
acp::schema::v1::ToolCallUpdateFields::new().title("shell_execute".to_owned()),
);
let allowed = gate
.check_permission(session.session_id, tool_call)
.await
.unwrap_or(false);
let _ = decision_tx.send(allowed);
let _ = channel.flush_chunks().await;
})
})
}
fn test_config(name: &str) -> AcpServerConfig {
AcpServerConfig {
agent_name: name.to_owned(),
agent_version: "0.0.1".to_owned(),
max_sessions: 8,
..AcpServerConfig::default()
}
}
fn test_config_with_models(name: &str, models: Vec<&str>) -> AcpServerConfig {
let factory: zeph_acp::ProviderFactory = Arc::new(|_key: &str| {
Some(zeph_llm::any::AnyProvider::Mock(
zeph_llm::mock::MockProvider::default(),
))
});
AcpServerConfig {
agent_name: name.to_owned(),
agent_version: "0.0.1".to_owned(),
max_sessions: 8,
provider_factory: Some(factory),
available_models: Arc::new(parking_lot::RwLock::new(
models.into_iter().map(str::to_owned).collect(),
)),
..AcpServerConfig::default()
}
}
#[cfg(feature = "unstable-llm-providers")]
fn test_config_with_provider_names(
name: &str,
providers: Vec<(&str, zeph_acp::LlmProtocol)>,
) -> AcpServerConfig {
AcpServerConfig {
agent_name: name.to_owned(),
agent_version: "0.0.1".to_owned(),
max_sessions: 8,
provider_names: providers
.into_iter()
.map(|(n, p)| (n.to_owned(), p))
.collect(),
..AcpServerConfig::default()
}
}
fn duplex_pair() -> (
impl futures::AsyncWrite + Unpin + Send + 'static,
impl futures::AsyncRead + Unpin + Send + 'static,
impl futures::AsyncWrite + Unpin + Send + 'static,
impl futures::AsyncRead + Unpin + Send + 'static,
) {
let (s_tok, c_tok) = tokio::io::duplex(64 * 1024);
let (s_read, s_write) = tokio::io::split(s_tok);
let (c_read, c_write) = tokio::io::split(c_tok);
(
s_write.compat_write(),
s_read.compat(),
c_write.compat_write(),
c_read.compat(),
)
}
fn temp_workdir() -> TempDir {
tempfile::tempdir().expect("failed to create temp dir")
}
fn select_current_value(option: &acp::schema::v1::SessionConfigOption) -> &str {
match &option.kind {
acp::schema::v1::SessionConfigKind::Select(select) => select.current_value.0.as_ref(),
#[allow(unreachable_patterns)]
other => panic!("expected a Select config option, got {other:?}"),
}
}
#[tokio::test(flavor = "current_thread")]
async fn initialize_handshake() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
let resp = cx
.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
assert!(resp.agent_info.is_some(), "agent_info missing");
let info = resp.agent_info.unwrap();
assert_eq!(info.name, "test-agent");
assert_eq!(info.version, "0.0.1");
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "initialize failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn new_session_returns_session_id() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let resp = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?;
assert!(
!resp.session_id.0.is_empty(),
"session_id must not be empty"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "new_session failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn cancel_notification_does_not_panic() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_resp = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?;
cx.send_notification(acp::schema::v1::CancelNotification::new(
session_resp.session_id,
))?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "cancel notification failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn unknown_ext_method_returns_null() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let raw_params =
Arc::from(serde_json::value::RawValue::from_string("{}".to_owned()).unwrap());
let resp = cx
.send_request(acp::schema::v1::ClientRequest::ExtMethodRequest(
acp::schema::v1::ExtRequest::new("_unknown_method", raw_params),
))
.block_task()
.await?;
assert_eq!(
resp.to_string(),
"null",
"unknown ext method must return null"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "ext_method failed: {result:?}");
}
}
})
.await;
}
#[cfg(feature = "unstable-llm-providers")]
#[tokio::test(flavor = "current_thread")]
async fn providers_list_reflects_configured_provider_names() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let (sw, sr, cw, cr) = duplex_pair();
let config = test_config_with_provider_names(
"test-agent",
vec![("openai", zeph_acp::LlmProtocol::OpenAi)],
);
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let raw_params =
Arc::from(serde_json::value::RawValue::from_string("{}".to_owned()).unwrap());
let resp = cx
.send_request(acp::schema::v1::ClientRequest::ExtMethodRequest(
acp::schema::v1::ExtRequest::new("providers/list", raw_params),
))
.block_task()
.await?;
let body = resp.to_string();
assert!(
body.contains("openai"),
"providers/list must include the configured provider name, got: {body}"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "providers/list ext_method failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn load_session_unknown_id_returns_error() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let err = cx
.send_request(acp::schema::v1::LoadSessionRequest::new(
"non-existent-session-id",
workdir.path(),
))
.block_task()
.await;
assert!(
err.is_err(),
"load_session of unknown id must return an error"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "client connection failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn session_list_contains_created_sessions() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let id_a = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let id_b = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let resp = cx
.send_request(acp::schema::v1::ListSessionsRequest::new())
.block_task()
.await?;
let ids: Vec<&acp::schema::v1::SessionId> =
resp.sessions.iter().map(|s| &s.session_id).collect();
assert!(ids.contains(&&id_a), "session A not in list: {ids:?}");
assert!(ids.contains(&&id_b), "session B not in list: {ids:?}");
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "list_sessions failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn prompt_round_trip_returns_end_turn() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
echo_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("hello"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(
resp.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"expected EndTurn, got {:?}",
resp.stop_reason,
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "prompt round-trip failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn drain_until_stop_collects_text_chunks() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
text_chunks_spawner(vec!["hello", " ", "world"]),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("go"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(resp.stop_reason, acp::schema::v1::StopReason::EndTurn);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "drain_until_stop text test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn review_command_output_is_delivered_to_client() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
text_chunks_spawner(vec!["review", " ", "output"]),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let collected = Arc::new(std::sync::Mutex::new(String::new()));
let collected_for_handler = Arc::clone(&collected);
let client_fut = acp::Client
.builder()
.on_receive_notification(
move |notification: acp::schema::v1::SessionNotification, _cx| {
let collected = Arc::clone(&collected_for_handler);
async move {
if let acp::schema::v1::SessionUpdate::AgentMessageChunk(chunk) =
notification.update
&& let acp::schema::v1::ContentBlock::Text(text) = chunk.content
{
collected.lock().unwrap().push_str(&text.text);
}
Ok(())
}
},
acp::on_receive_notification!(),
)
.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("/review"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(resp.stop_reason, acp::schema::v1::StopReason::EndTurn);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "review round-trip test failed: {result:?}");
}
}
assert_eq!(
*collected.lock().unwrap(),
"review output",
"the spawner's chunks must have reached the client as AgentMessageChunk \
notifications — before #6673 this would be empty since /review never called \
acquire_prompt_channels"
);
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn cancel_before_prompt_is_a_no_op() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
echo_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_notification(acp::schema::v1::CancelNotification::new(session_id.clone()))?;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("go"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(
resp.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"a cancel notification sent before any prompt is in flight must not \
retroactively cancel the next prompt, got {:?}",
resp.stop_reason,
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "cancel_before_prompt test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn late_cancel_after_prompt_completion_does_not_affect_next_prompt() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
multi_turn_echo_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let first_content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("first"),
)];
let first = cx
.send_request(acp::schema::v1::PromptRequest::new(
session_id.clone(),
first_content,
))
.block_task()
.await?;
assert_eq!(
first.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"first prompt must complete normally before the late cancel arrives"
);
cx.send_notification(acp::schema::v1::CancelNotification::new(session_id.clone()))?;
let second_content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("second"),
)];
let second = cx
.send_request(acp::schema::v1::PromptRequest::new(
session_id,
second_content,
))
.block_task()
.await?;
assert_eq!(
second.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"a cancel notification that arrives between two prompts must not cancel \
the next, unrelated prompt, got {:?}",
second.stop_reason,
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "late cancel regression test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn authenticate_returns_default_response() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::AuthenticateRequest::new("agent"))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "authenticate failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn logout_returns_default_response() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::LogoutRequest::new())
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "logout failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn close_session_removes_session_from_memory() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_request(acp::schema::v1::CloseSessionRequest::new(
session_id.clone(),
))
.block_task()
.await?;
let load_err = cx
.send_request(acp::schema::v1::LoadSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await;
assert!(
load_err.is_err(),
"loading a closed session must fail when no store is configured"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "close_session test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn delete_session_removes_session_from_list() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-delete-test.db")
.to_string_lossy()
.into_owned();
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_request(acp::schema::v1::DeleteSessionRequest::new(
session_id.clone(),
))
.block_task()
.await?;
let resp = cx
.send_request(acp::schema::v1::ListSessionsRequest::new())
.block_task()
.await?;
let ids: Vec<&acp::schema::v1::SessionId> =
resp.sessions.iter().map(|s| &s.session_id).collect();
assert!(
!ids.contains(&&session_id),
"deleted session must not appear in session/list: {ids:?}"
);
Ok(session_id)
});
let session_id = tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
result.expect("delete_session test failed")
}
};
let store = zeph_memory::store::SqliteStore::new(&sqlite_path)
.await
.expect("SqliteStore::new");
assert!(
!store
.acp_session_exists(&session_id.to_string())
.await
.expect("acp_session_exists query failed"),
"deleted session must not survive in the persistence store"
);
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn close_session_aborts_agent_loop_task() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let started = Arc::new(tokio::sync::Notify::new());
let alive = Arc::new(std::sync::atomic::AtomicBool::new(false));
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
hanging_spawner(Arc::clone(&started), Arc::clone(&alive)),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let started_for_client = Arc::clone(&started);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
started_for_client.notified().await;
cx.send_request(acp::schema::v1::CloseSessionRequest::new(session_id))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "close_session test failed: {result:?}");
}
}
assert!(
!alive.load(std::sync::atomic::Ordering::SeqCst),
"agent-loop task must be aborted and joined before session/close returns"
);
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn delete_session_aborts_agent_loop_task() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let started = Arc::new(tokio::sync::Notify::new());
let alive = Arc::new(std::sync::atomic::AtomicBool::new(false));
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
hanging_spawner(Arc::clone(&started), Arc::clone(&alive)),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let started_for_client = Arc::clone(&started);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
started_for_client.notified().await;
cx.send_request(acp::schema::v1::DeleteSessionRequest::new(session_id))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "delete_session test failed: {result:?}");
}
}
assert!(
!alive.load(std::sync::atomic::Ordering::SeqCst),
"agent-loop task must be aborted and joined before session/delete returns"
);
})
.await;
}
#[cfg(feature = "unstable-session-fork")]
#[tokio::test(flavor = "current_thread")]
async fn fork_session_creates_distinct_session_id() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let source_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let forked = cx
.send_request(acp::schema::v1::ForkSessionRequest::new(
source_id.clone(),
workdir.path(),
))
.block_task()
.await?;
assert_ne!(
forked.session_id, source_id,
"forked session must have a distinct id from the source"
);
assert!(!forked.session_id.0.is_empty());
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "fork_session test failed: {result:?}");
}
}
})
.await;
}
#[cfg(feature = "unstable-session-fork")]
#[tokio::test(flavor = "current_thread")]
async fn fork_session_copies_event_log_when_persistence_enabled() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-fork-test.db")
.to_string_lossy()
.into_owned();
let session_data_dir = db_dir.path().join("sessions");
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path),
session_data_dir: Some(session_data_dir.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let source_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let source_dir =
zeph_session::session_dir(&session_data_dir, &source_id.to_string());
let log = zeph_session::SessionEventLog::open(&source_dir)
.await
.expect("open source event log");
log.append(
None,
None,
zeph_session::SessionEvent::SessionStarted {
session_id: source_id.to_string(),
cwd: workdir.path().to_string_lossy().into_owned(),
provider_name: "claude".to_owned(),
model: "opus".to_owned(),
forked_from: None,
},
)
.await
.expect("append SessionStarted");
log.append(
None,
None,
zeph_session::SessionEvent::UserMessage {
text: "hello".to_owned(),
image_refs: vec![],
},
)
.await
.expect("append UserMessage");
let forked = cx
.send_request(acp::schema::v1::ForkSessionRequest::new(
source_id.clone(),
workdir.path(),
))
.block_task()
.await?;
let child_dir =
zeph_session::session_dir(&session_data_dir, &forked.session_id.to_string());
let child_log = zeph_session::SessionEventLog::open(&child_dir)
.await
.expect("open child event log");
let events = child_log.read_all().await.expect("read child event log");
assert_eq!(events.len(), 3, "child log must contain the copied events");
assert!(matches!(
events[0].kind,
zeph_session::SessionEvent::SessionStarted {
forked_from: Some(_),
..
}
));
let parent_log = zeph_session::SessionEventLog::open(&source_dir)
.await
.expect("reopen source event log");
let parent_events = parent_log.read_all().await.expect("read source event log");
assert!(
matches!(
parent_events.last().expect("parent has events").kind,
zeph_session::SessionEvent::ForkPoint { .. }
),
"parent log must record a ForkPoint provenance event"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "fork_session persistence test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn resume_session_reconnects_to_existing_session() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_request(acp::schema::v1::ResumeSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "resume_session test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn resume_session_reconstructs_from_store_after_restart() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-resume-test.db")
.to_string_lossy()
.into_owned();
let session_id = {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut =
acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
};
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(echo_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::ResumeSessionRequest::new(
session_id.clone(),
workdir.path(),
))
.block_task()
.await?;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("hello again"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(
resp.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"expected EndTurn from the store-reconstructed session, got {:?}",
resp.stop_reason,
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(
result.is_ok(),
"store-backed resume_session test failed: {result:?}"
);
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn load_session_succeeds_from_event_log_with_no_legacy_rows() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-load-test.db")
.to_string_lossy()
.into_owned();
let session_data_dir = db_dir.path().join("sessions");
let session_id = {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
session_data_dir: Some(session_data_dir.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut =
acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
};
let session_dir = zeph_session::session_dir(&session_data_dir, &session_id.to_string());
let log = zeph_session::SessionEventLog::open(&session_dir)
.await
.expect("open event log");
log.append(
None,
None,
zeph_session::SessionEvent::UserMessage {
text: "hello".to_owned(),
image_refs: vec![],
},
)
.await
.expect("append UserMessage");
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path),
session_data_dir: Some(session_data_dir),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::LoadSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "session/load from event log failed: {result:?}");
}
}
})
.await;
}
#[cfg(feature = "unstable-session-fork")]
#[tokio::test(flavor = "current_thread")]
async fn fork_session_inherits_config_from_in_memory_source() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config_with_models("test-agent", vec!["claude:sonnet"]),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let source_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
source_id.clone(),
"temperature",
"creative",
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
source_id.clone(),
"thinking",
"on",
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
source_id.clone(),
"auto_approve",
"auto-edit",
))
.block_task()
.await?;
let forked = cx
.send_request(acp::schema::v1::ForkSessionRequest::new(
source_id,
workdir.path(),
))
.block_task()
.await?;
let options = forked.config_options.unwrap_or_default();
let get = |id: &str| {
select_current_value(options.iter().find(|o| o.id.0.as_ref() == id).unwrap())
.to_owned()
};
assert_eq!(
get("temperature"),
"creative",
"forked session must inherit the source's temperature preset, not reset to \
the configured default"
);
assert_eq!(
get("thinking"),
"on",
"forked session must inherit the source's thinking toggle"
);
assert_eq!(
get("auto_approve"),
"auto-edit",
"forked session must inherit the source's auto-approve level"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "fork inheritance test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn resume_session_inherits_temperature_preset_after_close() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let captured: Arc<std::sync::Mutex<Option<zeph_llm::provider::GenerationOverrides>>> =
Arc::new(std::sync::Mutex::new(None));
let captured_for_factory = Arc::clone(&captured);
let factory: zeph_acp::ProviderFactory = Arc::new(move |_key: &str| {
Some(zeph_llm::any::AnyProvider::Mock(
zeph_llm::mock::MockProvider::default()
.with_overrides_capture(Arc::clone(&captured_for_factory)),
))
});
let mut config = test_config_with_models("test-agent", vec!["claude:sonnet"]);
config.provider_factory = Some(factory);
config.sqlite_path = Some(":memory:".to_owned());
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
session_id.clone(),
"temperature",
"creative",
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::CloseSessionRequest::new(
session_id.clone(),
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::ResumeSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "resume inheritance test failed: {result:?}");
}
}
let applied_temperature = captured
.lock()
.expect("capture mutex poisoned")
.as_ref()
.and_then(|o| o.temperature);
assert_eq!(
applied_temperature,
Some(zeph_config::AcpTemperaturePreset::Creative.temperature()),
"resumed session must inherit the closed source's persisted temperature preset, \
not reset to the configured default"
);
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn set_session_mode_switches_mode() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
cx.send_request(acp::schema::v1::SetSessionModeRequest::new(
session_id,
"architect",
))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "set_session_mode test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn set_session_config_option_model_switches_active_model() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config_with_models("test-agent", vec!["claude:sonnet", "ollama:llama3"]),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let resp = cx
.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
session_id,
"model",
"ollama:llama3",
))
.block_task()
.await?;
let model_option = resp
.config_options
.iter()
.find(|o| o.id.0.as_ref() == "model")
.expect("model option must be present");
assert_eq!(
select_current_value(model_option),
"ollama:llama3",
"model option must reflect the switched model"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "set_config_option model test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn set_session_config_option_temperature_preset_changes() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config_with_models("test-agent", vec!["claude:sonnet"]),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_resp = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?;
let session_id = session_resp.session_id;
let initial_temperature = session_resp
.config_options
.unwrap_or_default()
.into_iter()
.find(|o| o.id.0.as_ref() == "temperature")
.expect("temperature option must be advertised in new_session response");
assert_eq!(select_current_value(&initial_temperature), "balanced");
let resp = cx
.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
session_id,
"temperature",
"creative",
))
.block_task()
.await?;
let temperature_option = resp
.config_options
.iter()
.find(|o| o.id.0.as_ref() == "temperature")
.expect("temperature option must be present");
assert_eq!(
select_current_value(temperature_option),
"creative",
"temperature option must reflect the switched preset"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "set_config_option temperature test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn default_temperature_preset_is_primed_at_session_creation() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let captured: Arc<std::sync::Mutex<Option<zeph_llm::provider::GenerationOverrides>>> =
Arc::new(std::sync::Mutex::new(None));
let captured_for_factory = Arc::clone(&captured);
let factory: zeph_acp::ProviderFactory = Arc::new(move |_key: &str| {
Some(zeph_llm::any::AnyProvider::Mock(
zeph_llm::mock::MockProvider::default()
.with_overrides_capture(Arc::clone(&captured_for_factory)),
))
});
let mut config = test_config_with_models("test-agent", vec!["claude:sonnet"]);
config.provider_factory = Some(factory);
config.model_config = zeph_config::AcpModelConfigConfig {
default_temperature_preset: zeph_config::AcpTemperaturePreset::Creative,
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "acp-local".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "client connection failed: {result:?}");
}
}
let applied_temperature = captured
.lock()
.expect("capture mutex poisoned")
.as_ref()
.and_then(|o| o.temperature);
assert_eq!(
applied_temperature,
Some(zeph_config::AcpTemperaturePreset::Creative.temperature()),
"default_temperature_preset must be primed into the effective provider at \
session creation, with no session/set_config_option call made"
);
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn set_session_config_option_unknown_config_id_errors() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
noop_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let err = cx
.send_request(acp::schema::v1::SetSessionConfigOptionRequest::new(
session_id,
"nonexistent_option",
"whatever",
))
.block_task()
.await;
assert!(err.is_err(), "unknown config_id must return an error");
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "client connection failed: {result:?}");
}
}
})
.await;
}
#[cfg(feature = "unstable-cancel-request")]
#[tokio::test(flavor = "current_thread")]
async fn cancel_request_during_prompt_cancels() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let server_fut = serve_connection(
echo_spawner(),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("go"),
)];
let request =
cx.send_request(acp::schema::v1::PromptRequest::new(session_id, content));
request.cancel()?;
let resp = request.block_task().await;
match resp {
Ok(r) => {
assert!(matches!(
r.stop_reason,
acp::schema::v1::StopReason::EndTurn
| acp::schema::v1::StopReason::Cancelled
));
}
Err(e) => {
assert_eq!(
i32::from(e.code),
-32800,
"expected request_cancelled error"
);
}
}
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "cancel_request test failed: {result:?}");
}
}
})
.await;
}
async fn create_owned_session(
sqlite_path: &str,
owner: &str,
workdir: &std::path::Path,
) -> acp::schema::v1::SessionId {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.to_owned()),
..test_config("test-agent")
};
let server_fut = serve_connection(noop_spawner(), config, sw, sr, owner.to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
}
#[tokio::test(flavor = "current_thread")]
async fn list_sessions_isolated_by_owner_across_stdio_connections() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-owner-list-test.db")
.to_string_lossy()
.into_owned();
let alice_session = create_owned_session(&sqlite_path, "alice", workdir.path()).await;
let bob_session = create_owned_session(&sqlite_path, "bob", workdir.path()).await;
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut = serve_connection(noop_spawner(), config, sw, sr, "alice".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let resp = cx
.send_request(acp::schema::v1::ListSessionsRequest::new())
.block_task()
.await?;
let ids: Vec<&acp::schema::v1::SessionId> =
resp.sessions.iter().map(|s| &s.session_id).collect();
assert!(
ids.contains(&&alice_session),
"alice's own session missing from her list: {ids:?}"
);
assert!(
!ids.contains(&&bob_session),
"bob's session leaked into alice's list: {ids:?}"
);
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(result.is_ok(), "cross-owner list_sessions test failed: {result:?}");
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn resume_session_cross_owner_fails() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-owner-resume-test.db")
.to_string_lossy()
.into_owned();
let session_id = {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "alice".to_owned());
let client_fut =
acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
};
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut = serve_connection(echo_spawner(), config, sw, sr, "bob".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::ResumeSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(
result.is_err(),
"bob must not be able to resume alice's session, got: {result:?}"
);
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn load_session_cross_owner_fails() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-owner-load-test.db")
.to_string_lossy()
.into_owned();
let session_id = {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "alice".to_owned());
let client_fut =
acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
};
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut = serve_connection(echo_spawner(), config, sw, sr, "bob".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::LoadSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(
result.is_err(),
"bob must not be able to load alice's session, got: {result:?}"
);
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn delete_session_cross_owner_fails() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-owner-delete-test.db")
.to_string_lossy()
.into_owned();
let session_id = {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "alice".to_owned());
let client_fut =
acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
};
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut = serve_connection(noop_spawner(), config, sw, sr, "bob".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let _ = cx
.send_request(acp::schema::v1::DeleteSessionRequest::new(
session_id.clone(),
))
.block_task()
.await;
Ok(())
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
result.expect("bob's delete_session request errored unexpectedly");
}
}
let store = zeph_memory::store::SqliteStore::new(&sqlite_path)
.await
.expect("SqliteStore::new");
assert!(
store
.acp_session_exists(&session_id.to_string())
.await
.expect("acp_session_exists query failed"),
"bob must not be able to delete alice's session from the store"
);
})
.await;
}
#[cfg(feature = "unstable-session-fork")]
#[tokio::test(flavor = "current_thread")]
async fn fork_session_cross_owner_fails() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let db_dir = tempfile::tempdir().expect("failed to create temp db dir");
let sqlite_path = db_dir
.path()
.join("acp-owner-fork-test.db")
.to_string_lossy()
.into_owned();
let session_id = {
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut =
serve_connection(noop_spawner(), config, sw, sr, "alice".to_owned());
let client_fut =
acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
Ok(session_id)
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result.expect("session/new failed"),
}
};
let (sw, sr, cw, cr) = duplex_pair();
let config = AcpServerConfig {
sqlite_path: Some(sqlite_path.clone()),
..test_config("test-agent")
};
let server_fut = serve_connection(noop_spawner(), config, sw, sr, "bob".to_owned());
let client_fut = acp::Client.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
cx.send_request(acp::schema::v1::ForkSessionRequest::new(
session_id,
workdir.path(),
))
.block_task()
.await
});
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => {
assert!(
result.is_err(),
"bob must not be able to fork alice's session, got: {result:?}"
);
}
}
})
.await;
}
#[tokio::test(flavor = "current_thread")]
#[allow(clippy::large_futures)]
async fn permission_gated_prompt_round_trip_does_not_deadlock() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let (decision_tx, mut decision_rx) = tokio::sync::mpsc::unbounded_channel::<bool>();
let server_fut = serve_connection(
permission_gated_spawner(decision_tx),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client
.builder()
.on_receive_request(
async |_req: acp::schema::v1::RequestPermissionRequest,
responder: acp::Responder<
acp::schema::v1::RequestPermissionResponse,
>,
_cx| {
responder.respond(acp::schema::v1::RequestPermissionResponse::new(
acp::schema::v1::RequestPermissionOutcome::Selected(
acp::schema::v1::SelectedPermissionOutcome::new("allow_once"),
),
))
},
acp::on_receive_request!(),
)
.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("run a gated tool"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(
resp.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"expected EndTurn, got {:?}",
resp.stop_reason,
);
Ok(())
});
let outcome = tokio::time::timeout(std::time::Duration::from_secs(10), async {
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result,
}
})
.await;
let result = outcome.expect(
"permission-gated prompt round-trip timed out — likely a #6656 deadlock regression",
);
assert!(result.is_ok(), "prompt round-trip failed: {result:?}");
let decision = decision_rx
.recv()
.await
.expect("spawner must report a permission decision");
assert!(
decision,
"IDE selected allow_once, expected the gate to allow"
);
})
.await;
}
#[tokio::test(flavor = "current_thread")]
#[allow(clippy::large_futures)]
async fn permission_gated_prompt_denial_does_not_deadlock() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let workdir = temp_workdir();
let (sw, sr, cw, cr) = duplex_pair();
let (decision_tx, mut decision_rx) = tokio::sync::mpsc::unbounded_channel::<bool>();
let server_fut = serve_connection(
permission_gated_spawner(decision_tx),
test_config("test-agent"),
sw,
sr,
"acp-local".to_owned(),
);
let client_fut = acp::Client
.builder()
.on_receive_request(
async |_req: acp::schema::v1::RequestPermissionRequest,
responder: acp::Responder<
acp::schema::v1::RequestPermissionResponse,
>,
_cx| {
responder.respond(acp::schema::v1::RequestPermissionResponse::new(
acp::schema::v1::RequestPermissionOutcome::Selected(
acp::schema::v1::SelectedPermissionOutcome::new("reject_once"),
),
))
},
acp::on_receive_request!(),
)
.connect_with(acp::ByteStreams::new(cw, cr), async |cx| {
cx.send_request(acp::schema::v1::InitializeRequest::new(
acp::schema::ProtocolVersion::LATEST,
))
.block_task()
.await?;
let session_id = cx
.send_request(acp::schema::v1::NewSessionRequest::new(workdir.path()))
.block_task()
.await?
.session_id;
let content = vec![acp::schema::v1::ContentBlock::Text(
acp::schema::v1::TextContent::new("run a gated tool"),
)];
let resp = cx
.send_request(acp::schema::v1::PromptRequest::new(session_id, content))
.block_task()
.await?;
assert_eq!(
resp.stop_reason,
acp::schema::v1::StopReason::EndTurn,
"expected EndTurn, got {:?}",
resp.stop_reason,
);
Ok(())
});
let outcome = tokio::time::timeout(std::time::Duration::from_secs(10), async {
tokio::select! {
res = server_fut => panic!("server exited before client: {res:?}"),
result = client_fut => result,
}
})
.await;
let result = outcome.expect(
"permission-gated prompt round-trip timed out — likely a #6656 deadlock regression",
);
assert!(result.is_ok(), "prompt round-trip failed: {result:?}");
let decision = decision_rx
.recv()
.await
.expect("spawner must report a permission decision");
assert!(
!decision,
"IDE selected reject_once, expected the gate to deny"
);
})
.await;
}