use super::*;
mod plan_document;
fn install_test_agent_event_log_sink(
log: &std::sync::Arc<harn_vm::event_log::AnyEventLog>,
session_id: &str,
) {
harn_vm::agent_events::clear_session_sinks(session_id);
harn_vm::agent_events::register_sink(
session_id.to_string(),
harn_vm::agent_events::EventLogSink::new(log.clone(), session_id),
);
}
async fn run_mock_llm_prompt(
request_tx: &mpsc::UnboundedSender<serde_json::Value>,
response_rx: &mut mpsc::UnboundedReceiver<String>,
session_id: &str,
id: u64,
) -> serde_json::Value {
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"method": "session/prompt",
"params": {
"sessionId": session_id,
"prompt": [{
"type": "text",
"text": "harness.llm.mock_clear()\nharness.llm.mock_enqueue({text: \"ok\", input_tokens: 1, output_tokens: 1, model: \"mock\", provider: \"mock\"})\nlet r = harness.llm.call(\"hello\", nil, {provider: \"mock\", model: \"mock\"})\nharness.stdio.println(r.text)",
}],
},
}))
.expect("send session/prompt");
for _ in 0..32 {
let message = recv_json(response_rx).await;
if message["method"] == "host/capabilities" {
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": message["id"].clone(),
"result": {},
}))
.expect("send host capabilities response");
continue;
}
if message["id"] == id {
return message;
}
}
panic!("timed out waiting for session/prompt response {id}");
}
async fn recv_response_with_id(
response_rx: &mut mpsc::UnboundedReceiver<String>,
id: u64,
) -> serde_json::Value {
for _ in 0..32 {
let message = recv_json(response_rx).await;
if message["id"] == id {
return message;
}
}
panic!("timed out waiting for response {id}");
}
#[tokio::test(flavor = "current_thread")]
async fn acp_authenticate_uses_shared_auth_policy() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let config = AcpServerConfig::new(None).with_auth_policy(AuthPolicy {
methods: vec![AuthMethodConfig::ApiKey(ApiKeyAuthConfig::single("secret"))],
mcp_allowlist: None,
});
let (request_tx, request_rx) = mpsc::unbounded_channel();
let (response_tx, mut response_rx) = mpsc::unbounded_channel();
let server = tokio::task::spawn_local(super::run_acp_channel_server(
config,
request_rx,
response_tx,
));
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 0,
"method": "initialize",
}))
.expect("send initialize");
let initialize = recv_json(&mut response_rx).await;
assert_eq!(
initialize["result"]["authMethods"][0]["_meta"]["harn"]["scheme"],
"api_key"
);
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.expect("send unauthenticated session/new");
let blocked = recv_json(&mut response_rx).await;
assert_eq!(blocked["id"], 1);
assert_eq!(blocked["error"]["code"], ACP_AUTH_REQUIRED_CODE);
assert_eq!(blocked["error"]["message"], "auth_required");
assert_eq!(blocked["error"]["data"]["authMethods"][0]["id"], "apiKey");
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "authenticate",
"params": {
"methodId": "apiKey",
"_meta": {"harn": {"apiKey": "secret"}}
},
}))
.expect("send authenticate");
let authenticated = recv_json(&mut response_rx).await;
assert_eq!(authenticated["id"], 2);
assert_eq!(
authenticated["result"]["_meta"]["harn"]["principal"]["scheme"],
"api_key"
);
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/new",
"params": {"cwd": "."},
}))
.expect("send authenticated session/new");
let created = recv_json(&mut response_rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 4,
"method": "session/prompt",
"params": {
"sessionId": session_id,
"prompt": [{"type": "text", "text": "harness.stdio.println(\"allowed\")"}],
},
}))
.expect("send authenticated prompt");
let mut saw_allowed = false;
let mut saw_response = false;
for _ in 0..16 {
let message = recv_json(&mut response_rx).await;
match message.get("method").and_then(|value| value.as_str()) {
Some("host/capabilities") => {
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": message["id"].clone(),
"result": {},
}))
.expect("send host/capabilities response");
}
Some("session/update") => {
saw_allowed |= message["params"]["update"]["content"]["text"]
.as_str()
.is_some_and(|text| text == "allowed\n");
}
_ if message["id"] == 4 => {
saw_response = true;
assert_eq!(message["result"]["stopReason"], "end_turn");
break;
}
_ => {}
}
}
assert!(saw_allowed, "authenticated prompt should execute");
assert!(saw_response, "authenticated prompt should complete");
drop(request_tx);
server.await.expect("ACP channel server task");
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn acp_attach_server_advertises_and_authenticates_local_none_method() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
}))
.await;
let initialize = recv_json(&mut rx).await;
let methods = initialize["result"]["authMethods"]
.as_array()
.expect("authMethods array");
assert_eq!(methods.len(), 1);
assert_eq!(methods[0]["id"], "none");
assert_eq!(methods[0]["type"], "agent");
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "authenticate",
"params": {"methodId": "none"},
}))
.await;
let authenticated = recv_json(&mut rx).await;
assert_eq!(authenticated["id"], 2);
assert_eq!(
authenticated["result"]["_meta"]["harn"]["principal"]["scheme"],
"none"
);
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_new_advertises_session_mode_state_and_config_options() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let modes = &created["result"]["modes"];
assert_eq!(modes["currentModeId"], "ask");
let available = modes["availableModes"]
.as_array()
.expect("availableModes array");
let ids: Vec<&str> = available
.iter()
.map(|mode| mode["id"].as_str().expect("mode id"))
.collect();
assert_eq!(ids, vec!["ask", "architect", "code", "shadow"]);
for mode in available {
assert!(
mode["name"].as_str().is_some_and(|name| !name.is_empty()),
"mode {mode} missing name"
);
}
let config_options = created["result"]["configOptions"]
.as_array()
.expect("configOptions array");
let mode_option = config_options
.iter()
.find(|entry| entry["id"] == "mode")
.expect("mode config option");
assert_eq!(mode_option["category"], "mode");
assert_eq!(mode_option["type"], "select");
assert_eq!(mode_option["currentValue"], "ask");
let option_ids: Vec<&str> = mode_option["options"]
.as_array()
.expect("mode options")
.iter()
.map(|mode| mode["value"].as_str().expect("mode value"))
.collect();
assert_eq!(option_ids, vec!["ask", "architect", "code", "shadow"]);
let model_option = config_options
.iter()
.find(|entry| entry["id"] == "model")
.expect("model config option");
assert_eq!(model_option["category"], "model");
assert_eq!(model_option["type"], "select");
assert_eq!(model_option["currentValue"], "@inherit");
let thought_option = config_options
.iter()
.find(|entry| entry["id"] == "thought_level")
.expect("thought level config option");
assert_eq!(thought_option["category"], "model");
assert_eq!(thought_option["type"], "select");
assert_eq!(thought_option["currentValue"], "@inherit");
assert!(thought_option["options"]
.as_array()
.expect("thought options")
.iter()
.any(|entry| entry["value"] == "auto"));
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_load_includes_current_mode_state() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "architect"},
}))
.await;
let _ack = recv_json(&mut rx).await;
let _mode_notification = recv_json(&mut rx).await;
let _config_notification = recv_json(&mut rx).await;
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/load",
"params": {"sessionId": session_id},
}))
.await;
let loaded = recv_json(&mut rx).await;
assert_eq!(loaded["result"]["modes"]["currentModeId"], "architect");
assert_eq!(
loaded["result"]["configOptions"][0]["currentValue"],
"architect"
);
assert!(loaded["result"]["modes"]["availableModes"]
.as_array()
.expect("available modes")
.iter()
.any(|m| m["id"] == "architect"));
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_resume_includes_current_mode_state_without_replay() {
harn_vm::event_log::reset_active_event_log();
let log = harn_vm::event_log::install_memory_for_current_thread(64);
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
install_test_agent_event_log_sink(&log, &session_id);
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "architect"},
}))
.await;
let _ack = recv_json(&mut rx).await;
let _mode_notification = recv_json(&mut rx).await;
let _config_notification = recv_json(&mut rx).await;
harn_vm::agent_events::emit_event(&harn_vm::agent_events::AgentEvent::AgentMessageChunk {
session_id: session_id.clone(),
content: "do not replay me".to_string(),
});
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/resume",
"params": {"sessionId": session_id},
}))
.await;
let resumed = recv_json(&mut rx).await;
assert_eq!(resumed["id"], 3);
assert!(resumed.get("method").is_none(), "resume must respond first");
assert_eq!(resumed["result"]["modes"]["currentModeId"], "architect");
assert_eq!(
resumed["result"]["configOptions"][0]["currentValue"],
"architect"
);
assert!(
resumed["result"].get("replayed").is_none(),
"session/resume must not include replay metadata"
);
assert!(
rx.try_recv().is_err(),
"session/resume must not emit replay notifications"
);
harn_vm::agent_events::clear_session_sinks(created["result"]["sessionId"].as_str().unwrap());
harn_vm::event_log::reset_active_event_log();
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_restore_methods_reject_unknown_sessions() {
let project = tempfile::tempdir().expect("project root");
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
for (id, method) in [(1, "session/load"), (2, "session/resume")] {
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"method": method,
"params": {"sessionId": "missing-session", "cwd": project.path()},
}))
.await;
let response = recv_json(&mut rx).await;
assert_eq!(response["id"], id);
assert_eq!(response["error"]["code"], -32602);
assert!(
response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("unknown session")),
"unexpected error for {method}: {response}"
);
}
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_load_restores_persisted_session_unknown_to_server() {
harn_vm::event_log::reset_active_event_log();
let log = harn_vm::event_log::install_memory_for_current_thread(64);
let session_id = "burin-saved-session".to_string();
install_test_agent_event_log_sink(&log, &session_id);
harn_vm::agent_events::emit_event(&harn_vm::agent_events::AgentEvent::AgentMessageChunk {
session_id: session_id.clone(),
content: "restored history".to_string(),
});
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/load",
"params": {"sessionId": session_id, "cwd": "."},
}))
.await;
let replay = recv_json(&mut rx).await;
assert_eq!(replay["method"], "session/update");
assert_eq!(
replay["params"]["update"]["sessionUpdate"],
"agent_message_chunk"
);
assert_eq!(
replay["params"]["update"]["content"]["text"],
"restored history"
);
assert_eq!(
replay["params"]["update"]["_meta"]["harn"]["replayed"],
true
);
let loaded = recv_json(&mut rx).await;
assert_eq!(loaded["id"], 1);
assert_eq!(loaded["result"]["session"]["sessionId"], session_id);
assert_eq!(
loaded["result"]["session"]["liveState"], "live",
"the in-process server restores to a live, promptable session"
);
assert_eq!(loaded["result"]["replayed"].as_array().unwrap().len(), 1);
assert_eq!(
loaded["result"]["replayed"][0]["type"],
"agent_message_chunk"
);
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/load",
"params": {"sessionId": session_id},
}))
.await;
let _second_replay = recv_json(&mut rx).await;
let second = recv_json(&mut rx).await;
assert_eq!(second["id"], 2);
assert_eq!(second["result"]["session"]["sessionId"], session_id);
harn_vm::agent_events::clear_session_sinks(&session_id);
harn_vm::event_log::reset_active_event_log();
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_load_rejects_session_without_persisted_events() {
harn_vm::event_log::reset_active_event_log();
let _log = harn_vm::event_log::install_memory_for_current_thread(64);
let project = tempfile::tempdir().expect("project root");
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/load",
"params": {"sessionId": "never-existed", "cwd": project.path()},
}))
.await;
let response = recv_json(&mut rx).await;
assert_eq!(response["id"], 1);
assert_eq!(response["error"]["code"], -32602);
assert!(response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("unknown session")));
harn_vm::event_log::reset_active_event_log();
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_mode_emits_current_mode_update_notification() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "architect"},
}))
.await;
let ack = recv_json(&mut rx).await;
assert_eq!(ack["id"], 2);
assert!(ack["result"].is_object());
assert!(ack["error"].is_null());
let notification = recv_json(&mut rx).await;
assert_eq!(notification["method"], "session/update");
assert_eq!(notification["params"]["sessionId"], session_id);
assert_eq!(
notification["params"]["update"]["sessionUpdate"],
"current_mode_update"
);
assert_eq!(notification["params"]["update"]["modeId"], "architect");
let config_notification = recv_json(&mut rx).await;
assert_eq!(config_notification["method"], "session/update");
assert_eq!(
config_notification["params"]["update"]["sessionUpdate"],
"config_option_update"
);
assert_eq!(
config_notification["params"]["update"]["configOptions"][0]["currentValue"],
"architect"
);
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_mode_is_idempotent_when_mode_unchanged() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "ask"},
}))
.await;
let ack = recv_json(&mut rx).await;
assert_eq!(ack["id"], 2);
assert!(ack["result"].is_object());
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "architect"},
}))
.await;
let _ack2 = recv_json(&mut rx).await;
let notification = recv_json(&mut rx).await;
assert_eq!(notification["method"], "session/update");
assert_eq!(notification["params"]["update"]["modeId"], "architect");
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_updates_mode() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "mode",
"value": "shadow",
},
}))
.await;
let ack = recv_json(&mut rx).await;
assert_eq!(ack["id"], 2);
assert_eq!(ack["result"]["configOptions"][0]["currentValue"], "shadow");
let mode_notification = recv_json(&mut rx).await;
assert_eq!(
mode_notification["params"]["update"]["sessionUpdate"],
"current_mode_update"
);
assert_eq!(mode_notification["params"]["update"]["modeId"], "shadow");
let config_notification = recv_json(&mut rx).await;
assert_eq!(
config_notification["params"]["update"]["sessionUpdate"],
"config_option_update"
);
assert_eq!(
config_notification["params"]["update"]["configOptions"][0]["currentValue"],
"shadow"
);
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_mode_validates_inputs() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/set_mode",
"params": {"sessionId": "ghost", "modeId": "architect"},
}))
.await;
let unknown_session = recv_json(&mut rx).await;
assert_eq!(unknown_session["id"], 1);
assert_eq!(unknown_session["error"]["code"], -32602);
assert!(unknown_session["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("Unknown session"));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "not-a-mode"},
}))
.await;
let unknown_mode = recv_json(&mut rx).await;
assert_eq!(unknown_mode["id"], 3);
assert_eq!(unknown_mode["error"]["code"], -32602);
assert!(unknown_mode["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("Unknown mode"));
}
#[tokio::test(flavor = "current_thread")]
async fn acp_architect_mode_blocks_destructive_writes_in_prompt() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let dir = tempfile::tempdir().expect("tempdir");
let target = dir.path().join("forbidden.txt");
let target_str = target
.to_str()
.expect("temp path is utf-8")
.replace('\\', "\\\\");
let (request_tx, request_rx) = mpsc::unbounded_channel();
let (response_tx, mut response_rx) = mpsc::unbounded_channel();
let server = tokio::task::spawn_local(super::run_acp_channel_server(
AcpServerConfig::new(None),
request_rx,
response_tx,
));
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": dir.path().display().to_string()},
}))
.expect("send session/new");
let created = recv_json(&mut response_rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "architect"},
}))
.expect("send session/set_mode");
let _ack = recv_json(&mut response_rx).await;
let mode_notification = recv_json(&mut response_rx).await;
assert_eq!(
mode_notification["params"]["update"]["sessionUpdate"],
"current_mode_update"
);
assert_eq!(mode_notification["params"]["update"]["modeId"], "architect");
let prompt_source =
format!("harness.fs.write_text(\"{target_str}\", \"should not be written\")");
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/prompt",
"params": {
"sessionId": session_id,
"prompt": [{"type": "text", "text": prompt_source}],
},
}))
.expect("send session/prompt");
let mut saw_error = false;
for _ in 0..32 {
let message = recv_json(&mut response_rx).await;
if message["method"] == "host/capabilities" {
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": message["id"].clone(),
"result": {},
}))
.expect("send host capabilities response");
continue;
}
if message["id"] == 3 {
let error = &message["error"];
assert!(
!error.is_null(),
"architect mode should reject write_file but got result {:?}",
message["result"]
);
let message_text = error["message"].as_str().unwrap_or_default();
assert!(
message_text.contains("active effect ceiling")
&& message_text.contains("fs:write"),
"unexpected error message: {message_text}"
);
saw_error = true;
break;
}
}
assert!(saw_error, "prompt should produce a JSON-RPC error response");
assert!(
!target.exists(),
"architect mode must not allow write_file to mutate the workspace"
);
drop(request_tx);
server.await.expect("ACP channel server task");
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn acp_code_mode_allows_writes_in_prompt() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let dir = tempfile::tempdir().expect("tempdir");
let target = dir.path().join("allowed.txt");
let target_str = target
.to_str()
.expect("temp path is utf-8")
.replace('\\', "\\\\");
let (request_tx, request_rx) = mpsc::unbounded_channel();
let (response_tx, mut response_rx) = mpsc::unbounded_channel();
let server = tokio::task::spawn_local(super::run_acp_channel_server(
AcpServerConfig::new(None),
request_rx,
response_tx,
));
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": dir.path().display().to_string()},
}))
.expect("send session/new");
let created = recv_json(&mut response_rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": session_id, "modeId": "code"},
}))
.expect("send session/set_mode");
let _ack = recv_json(&mut response_rx).await;
let _notification = recv_json(&mut response_rx).await;
let prompt_source =
format!("harness.fs.write_text(\"{target_str}\", \"hello from code mode\")");
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/prompt",
"params": {
"sessionId": session_id,
"prompt": [{"type": "text", "text": prompt_source}],
},
}))
.expect("send session/prompt");
let mut saw_completed = false;
for _ in 0..32 {
let message = recv_json(&mut response_rx).await;
if message["method"] == "host/capabilities" {
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": message["id"].clone(),
"result": {},
}))
.expect("send host capabilities response");
continue;
}
if message["id"] == 3 {
assert_eq!(message["result"]["stopReason"], "end_turn");
saw_completed = true;
break;
}
}
assert!(saw_completed, "code-mode prompt should complete");
assert!(
target.exists(),
"code mode should allow write_file to mutate the workspace"
);
drop(request_tx);
server.await.expect("ACP channel server task");
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn acp_session_fork_inherits_parent_current_mode() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let parent_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_mode",
"params": {"sessionId": parent_id, "modeId": "architect"},
}))
.await;
let _ack = recv_json(&mut rx).await;
let _notification = recv_json(&mut rx).await;
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/fork",
"params": {"session_id": parent_id},
}))
.await;
let mut fork_response = None;
for _ in 0..6 {
let msg = recv_json(&mut rx).await;
if msg["id"] == 3 {
fork_response = Some(msg);
break;
}
}
let fork_response = fork_response.expect("fork response");
assert_eq!(fork_response["result"]["state"], "forked");
assert_eq!(
fork_response["result"]["modes"]["currentModeId"],
"architect"
);
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_pins_model_and_emits_update() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "model",
"value": "claude-sonnet-4-6",
},
}))
.await;
let ack = recv_json(&mut rx).await;
assert_eq!(ack["id"], 2);
let pinned_value = ack["result"]["configOptions"]
.as_array()
.expect("configOptions array")
.iter()
.find(|entry| entry["id"] == "model")
.expect("model config option")
.get("currentValue")
.and_then(|v| v.as_str())
.expect("currentValue string");
assert_eq!(pinned_value, "claude-sonnet-4-6");
let notification = recv_json(&mut rx).await;
assert_eq!(notification["method"], "session/update");
assert_eq!(
notification["params"]["update"]["sessionUpdate"],
"config_option_update"
);
let pin_in_notification = notification["params"]["update"]["configOptions"]
.as_array()
.expect("configOptions in notification")
.iter()
.find(|entry| entry["id"] == "model")
.expect("model entry in notification")
.get("currentValue")
.and_then(|v| v.as_str())
.expect("currentValue string");
assert_eq!(pin_in_notification, "claude-sonnet-4-6");
assert_eq!(
harn_vm::agent_sessions::pinned_model(&session_id).as_deref(),
Some("claude-sonnet-4-6"),
"vm session state must reflect the pin"
);
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "model",
"value": "@inherit",
},
}))
.await;
let clear_ack = recv_json(&mut rx).await;
assert_eq!(clear_ack["id"], 3);
let cleared_value = clear_ack["result"]["configOptions"]
.as_array()
.expect("configOptions array")
.iter()
.find(|entry| entry["id"] == "model")
.expect("model config option")
.get("currentValue")
.and_then(|v| v.as_str())
.expect("currentValue string");
assert_eq!(cleared_value, "@inherit");
let _clear_notification = recv_json(&mut rx).await;
assert!(
harn_vm::agent_sessions::pinned_model(&session_id).is_none(),
"clearing the pin should remove the vm-side selector"
);
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_pins_thought_level_and_emits_update() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "thought_level",
"value": "NO_THINK",
},
}))
.await;
let ack = recv_json(&mut rx).await;
assert_eq!(ack["id"], 2);
let pinned_value = ack["result"]["configOptions"]
.as_array()
.expect("configOptions array")
.iter()
.find(|entry| entry["id"] == "thought_level")
.expect("thought config option")
.get("currentValue")
.and_then(|v| v.as_str())
.expect("currentValue string");
assert_eq!(pinned_value, "off");
let notification = recv_json(&mut rx).await;
assert_eq!(notification["method"], "session/update");
assert_eq!(
notification["params"]["update"]["sessionUpdate"],
"config_option_update"
);
let pin_in_notification = notification["params"]["update"]["configOptions"]
.as_array()
.expect("configOptions in notification")
.iter()
.find(|entry| entry["id"] == "thought_level")
.expect("thought entry in notification")
.get("currentValue")
.and_then(|v| v.as_str())
.expect("currentValue string");
assert_eq!(pin_in_notification, "off");
assert_eq!(
harn_vm::agent_sessions::pinned_reasoning_policy(&session_id).as_deref(),
Some("off"),
);
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "thought_level",
"value": "@inherit",
},
}))
.await;
let clear_ack = recv_json(&mut rx).await;
assert_eq!(clear_ack["id"], 3);
let cleared_value = clear_ack["result"]["configOptions"]
.as_array()
.expect("configOptions array")
.iter()
.find(|entry| entry["id"] == "thought_level")
.expect("thought config option")
.get("currentValue")
.and_then(|v| v.as_str())
.expect("currentValue string");
assert_eq!(cleared_value, "@inherit");
let _clear_notification = recv_json(&mut rx).await;
assert!(harn_vm::agent_sessions::pinned_reasoning_policy(&session_id).is_none());
}
#[tokio::test(flavor = "current_thread")]
async fn acp_prompt_installs_default_llm_token_budget() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let config = AcpServerConfig::new(None).with_llm_token_budget(0);
let (request_tx, mut response_rx, server, session_id) =
start_acp_code_session_with_config(config, serde_json::json!(".")).await;
let response = run_mock_llm_prompt(&request_tx, &mut response_rx, &session_id, 3).await;
assert_eq!(response["id"], 3);
assert_eq!(response["error"]["code"], -32000);
let message = response["error"]["message"]
.as_str()
.unwrap_or_default()
.to_string();
assert!(
message.contains("LLM token budget exceeded"),
"prompt should fail under the inherited token budget: {message}"
);
drop(request_tx);
server.await.expect("ACP channel server task");
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_budget_rearms_next_prompt() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let (request_tx, mut response_rx, server, session_id) =
start_acp_code_session_with_config(
AcpServerConfig::new(None),
serde_json::json!("."),
)
.await;
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "budget",
"value": "{\"llm_tokens\":0}",
},
}))
.expect("send budget cap");
let ack = recv_response_with_id(&mut response_rx, 3).await;
assert_eq!(ack["id"], 3);
let budget_value = ack["result"]["configOptions"]
.as_array()
.expect("configOptions array")
.iter()
.find(|entry| entry["id"] == "budget")
.expect("budget option")
.get("currentValue")
.and_then(|value| value.as_str())
.expect("budget currentValue");
assert_eq!(budget_value, "{\"llm_tokens\":0}");
let _notification = recv_json(&mut response_rx).await;
let exhausted =
run_mock_llm_prompt(&request_tx, &mut response_rx, &session_id, 4).await;
assert_eq!(exhausted["error"]["code"], -32000);
assert!(exhausted["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("LLM token budget exceeded"));
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 5,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "budget",
"value": "{\"llm_tokens\":100}",
},
}))
.expect("send budget re-arm");
let rearm_ack = recv_response_with_id(&mut response_rx, 5).await;
assert_eq!(rearm_ack["id"], 5);
let _rearm_notification = recv_json(&mut response_rx).await;
let ok = run_mock_llm_prompt(&request_tx, &mut response_rx, &session_id, 6).await;
assert_eq!(ok["result"]["stopReason"], "end_turn");
drop(request_tx);
server.await.expect("ACP channel server task");
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_budget_off_disables_inherited_default() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let config = AcpServerConfig::new(None).with_llm_token_budget(0);
let (request_tx, mut response_rx, server, session_id) =
start_acp_code_session_with_config(config, serde_json::json!(".")).await;
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": 3,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "budget",
"value": "off",
},
}))
.expect("send budget off");
let ack = recv_response_with_id(&mut response_rx, 3).await;
let budget_value = ack["result"]["configOptions"]
.as_array()
.expect("configOptions array")
.iter()
.find(|entry| entry["id"] == "budget")
.expect("budget option")
.get("currentValue")
.and_then(|value| value.as_str())
.expect("budget currentValue");
assert_eq!(budget_value, "off");
let _notification = recv_json(&mut response_rx).await;
let ok = run_mock_llm_prompt(&request_tx, &mut response_rx, &session_id, 4).await;
assert_eq!(ok["result"]["stopReason"], "end_turn");
drop(request_tx);
server.await.expect("ACP channel server task");
})
.await;
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_rejects_unknown_model_provider() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "model",
"value": "nosuchprovider:nosuchmodel",
},
}))
.await;
let response = recv_json(&mut rx).await;
assert_eq!(response["id"], 2);
assert_eq!(response["error"]["code"], -32602);
let message = response["error"]["message"]
.as_str()
.unwrap_or_default()
.to_string();
assert!(
message.contains("invalid_model"),
"error message should be tagged invalid_model: {message}"
);
assert!(
harn_vm::agent_sessions::pinned_model(&session_id).is_none(),
"rejected pin must not mutate vm session state"
);
}
#[tokio::test(flavor = "current_thread")]
async fn acp_set_config_option_rejects_unknown_config_id() {
let (tx, mut rx) = mpsc::unbounded_channel();
let mut server = AcpServer::new_with_output(AcpServerConfig::new(None), AcpOutput::Channel(tx));
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": "."},
}))
.await;
let created = recv_json(&mut rx).await;
let session_id = created["result"]["sessionId"]
.as_str()
.expect("session id")
.to_string();
server
.handle_incoming_message(serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/set_config_option",
"params": {
"sessionId": session_id,
"configId": "temperature",
"value": "0.4",
},
}))
.await;
let response = recv_json(&mut rx).await;
assert_eq!(response["id"], 2);
assert_eq!(response["error"]["code"], -32602);
let message = response["error"]["message"]
.as_str()
.unwrap_or_default()
.to_string();
assert!(
message.contains("Unknown config option") && message.contains("model"),
"error message should advertise the registry's supported ids: {message}"
);
}