use super::*;
struct MockModeGuard;
impl Drop for MockModeGuard {
fn drop(&mut self) {
harn_vm::llm::clear_cli_llm_mock_mode();
}
}
fn answer_host_capabilities(
request_tx: &mpsc::UnboundedSender<serde_json::Value>,
message: &serde_json::Value,
) -> bool {
if message["method"] != "host/capabilities" {
return false;
}
request_tx
.send(serde_json::json!({
"jsonrpc": "2.0",
"id": message["id"].clone(),
"result": {},
}))
.expect("send host capabilities response");
true
}
async fn prompt(
request_tx: &mpsc::UnboundedSender<serde_json::Value>,
response_rx: &mut mpsc::UnboundedReceiver<String>,
session_id: &str,
id: u64,
source: &str,
) -> 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": source}],
},
}))
.expect("send session/prompt");
let mut seen_methods = Vec::new();
tokio::time::timeout(std::time::Duration::from_secs(15), async {
loop {
let line = response_rx
.recv()
.await
.expect("ACP response channel closed");
let message: serde_json::Value = serde_json::from_str(&line).expect("ACP JSON line");
if answer_host_capabilities(request_tx, &message) {
seen_methods.push("host/capabilities".to_string());
continue;
}
if message["id"] == id {
return message;
}
seen_methods.push(
message["method"]
.as_str()
.unwrap_or("<response-with-another-id>")
.to_string(),
);
}
})
.await
.unwrap_or_else(|_| {
panic!("timed out waiting for session/prompt response {id}; saw {seen_methods:?}")
})
}
#[tokio::test(flavor = "current_thread")]
async fn default_mode_reaches_a_model_turn_and_still_refuses_a_workspace_write() {
let local = tokio::task::LocalSet::new();
local
.run_until(async {
let mock = harn_vm::llm::parse_llm_mock_value(&serde_json::json!({
"text": "all done",
"model": "served-proof",
"provider": "mock",
}))
.expect("mock fixture");
harn_vm::llm::install_cli_llm_mocks(vec![mock]);
let _mock_mode = MockModeGuard;
let dir = tempfile::tempdir().expect("tempdir");
let pipeline = dir.path().join("served-agent-turn.harn");
std::fs::write(
&pipeline,
r#"import { agent_loop } from "std/agent/loop"
pipeline default(harness: Harness) {
if prompt == "runtime-control" {
agent_loop(harness, prompt, nil, {provider: "mock", model: "served-proof"})
assert(
len(harness.llm.mock_calls()) == 1,
"the default served task must consume exactly one model call",
)
harness.stdio.println("control-plane-ok")
} else {
harness.fs.write_text("ceiling-probe.txt", "must-not-write")
}
}
"#,
)
.expect("write served pipeline");
let (request_tx, mut response_rx, _server, session_id) =
start_acp_channel_session_with_config(
AcpServerConfig::for_pipeline(pipeline.to_string_lossy().to_string()),
serde_json::json!(dir.path()),
)
.await;
let admitted = prompt(
&request_tx,
&mut response_rx,
&session_id,
2,
"runtime-control",
)
.await;
assert!(
admitted["error"].is_null(),
"the agent loop's own session control plane must survive the ceiling \
every non-`code` session mode installs: {}",
admitted["error"]
);
assert_eq!(
admitted["result"]["stopReason"], "end_turn",
"default-mode control-plane turn should complete; got {}",
admitted["result"]
);
let durable = prompt(
&request_tx,
&mut response_rx,
&session_id,
3,
"workspace-write",
)
.await;
let rejection = format!(
"{}{}",
durable["error"]["message"].as_str().unwrap_or_default(),
durable["result"]["stopReason"].as_str().unwrap_or_default()
);
assert!(
rejection.contains("exceeds the active effect ceiling"),
"this session must really be carrying the `read_only` ceiling, or the \
admitted control-plane turn above proves nothing; got error={} \
result={}",
durable["error"],
durable["result"]
);
assert!(
!dir.path().join("ceiling-probe.txt").exists(),
"the rejected workspace write reached the filesystem"
);
drop(request_tx);
})
.await;
}