harn-serve 0.10.151

Shared outbound workflow server core for Harn adapters
//! A `json`-format tool call is hidden from the reply a host shows, and a call
//! to a tool the agent does not have still reaches the host as a failed tool
//! row. Together they mean hiding the fence never makes a call disappear.
//!
//! Driven over the real ACP wire with the provider mock armed outside the
//! governed Harn source, the same way `served_agent_turn` does.

use super::*;

struct MockModeGuard;

impl Drop for MockModeGuard {
    fn drop(&mut self) {
        harn_vm::llm::clear_cli_llm_mock_mode();
    }
}

#[tokio::test(flavor = "current_thread")]
async fn a_json_tool_fence_is_hidden_and_an_undeclared_call_still_shows_as_failed() {
    let local = tokio::task::LocalSet::new();
    local
        .run_until(async {
            let mocks = [
                "Checking first.\n\n```tool\n{\"name\":\"nuke_repo\",\"args\":{\"path\":\"a\"}}\n```",
                "Done.",
            ]
            .into_iter()
            .map(|text| {
                harn_vm::llm::parse_llm_mock_value(&serde_json::json!({
                    "text": text,
                    "model": "fence-proof",
                    "provider": "mock",
                }))
                .expect("mock fixture")
            })
            .collect();
            harn_vm::llm::install_cli_llm_mocks(mocks);
            let _mock_mode = MockModeGuard;

            let dir = tempfile::tempdir().expect("tempdir");
            let pipeline = dir.path().join("fence.harn");
            std::fs::write(
                &pipeline,
                r#"import { agent_loop } from "std/agent/loop"
pipeline default(harness: Harness) {
  let tools = tool_registry()
  tools = tool_define(tools, "inspect", "Inspect a path.", {
    handler: { args -> "observed " + (args?.path ?? "") },
    parameters: {path: {type: "string"}},
    returns: {type: "string"},
    annotations: {kind: "read", side_effect_level: "read_only"},
  })
  agent_loop(harness, prompt, nil, {
    provider: "mock",
    model: "fence-proof",
    tool_format: "json",
    loop_until_done: true,
    max_iterations: 3,
    max_nudges: 2,
    tools: tools,
  })
}
"#,
            )
            .expect("write 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;
            request_tx
                .send(serde_json::json!({
                    "jsonrpc": "2.0",
                    "id": 2,
                    "method": "session/prompt",
                    "params": {
                        "sessionId": session_id,
                        "prompt": [{"type": "text", "text": "Inspect a."}],
                    },
                }))
                .expect("send session/prompt");

            let mut updates = Vec::new();
            let response = tokio::time::timeout(std::time::Duration::from_secs(30), async {
                loop {
                    let line = response_rx.recv().await.expect("ACP channel closed");
                    let message: serde_json::Value = serde_json::from_str(&line).expect("JSON");
                    if message["method"] == "host/capabilities" {
                        request_tx
                            .send(serde_json::json!({
                                "jsonrpc": "2.0",
                                "id": message["id"].clone(),
                                "result": {},
                            }))
                            .expect("answer host capabilities");
                    } else if message["id"] == 2 {
                        return message;
                    } else if message["method"] == "session/update" {
                        updates.push(message["params"]["update"].clone());
                    }
                }
            })
            .await
            .expect("session/prompt response");
            assert!(response["error"].is_null(), "{response:#}");

            let shown: Vec<String> = updates
                .iter()
                .filter(|update| update["sessionUpdate"] == "agent_message_chunk")
                .filter_map(|update| {
                    let content = &update["content"];
                    content["_meta"]["harn"]["visible_text"]
                        .as_str()
                        .or_else(|| content["text"].as_str())
                        .map(str::to_string)
                })
                .collect();
            assert!(
                shown.iter().any(|text| text.contains("Checking first.")),
                "the narration around the call must stay visible: {shown:?}"
            );
            assert!(
                shown
                    .iter()
                    .all(|text| !text.contains("```tool") && !text.contains("nuke_repo")),
                "the call must not be shown again as raw JSON: {shown:?}"
            );

            let rows: Vec<&serde_json::Value> = updates
                .iter()
                .filter(|update| {
                    update["sessionUpdate"] == "tool_call"
                        || update["sessionUpdate"] == "tool_call_update"
                })
                .collect();
            assert!(
                rows.iter()
                    .any(|row| row.to_string().contains("nuke_repo") && row["status"] == "failed"),
                "the undeclared call must reach the host as a failed tool row: {rows:#?}"
            );
            drop(request_tx);
        })
        .await;
}