use luft::adapters::{AcpAdapter, AcpConfig};
use luft::core::contract::backend::{AgentBackend, RunContext};
use luft::core::{BackendRegistry, Scheduler, SchedulerConfig};
use luft::runtime::{ExecLimits, Runtime};
use std::sync::Arc;
use std::time::Duration;
use tokio_util::sync::CancellationToken;
fn fake_acp_backend() -> Arc<dyn AgentBackend> {
let mut env_passthrough: Vec<String> = AcpConfig::DEFAULT_ENV_PASSTHROUGH
.iter()
.map(|s| s.to_string())
.collect();
env_passthrough.push("FAKE_ACP_RAW_INPUT".into());
let binary = std::env::var("CARGO_BIN_EXE_fake_acp")
.or_else(|_| std::env::var("CARGO_BIN_EXE_fake-acp"))
.expect("CARGO_BIN_EXE_fake_acp not set; fake-acp binary must be built")
.into();
let config = AcpConfig {
id: "fake-acp",
binary,
acp_args: vec![],
log_level: None,
connect_timeout: Duration::from_secs(10),
emit_raw_events: true,
env_passthrough,
env: Default::default(),
model: None,
luft_binary: None,
};
Arc::new(AcpAdapter::new(config))
}
async fn run_with_fake_acp(
schema: serde_json::Value,
raw_input: serde_json::Value,
) -> serde_json::Value {
std::env::set_var("FAKE_ACP_RAW_INPUT", raw_input.to_string());
let backend = fake_acp_backend();
let registry = BackendRegistry::new().with(backend);
let scheduler = Scheduler::new(SchedulerConfig::default(), registry, None);
let (tx, _rx) = tokio::sync::broadcast::channel(256);
let run_id = uuid::Uuid::now_v7();
let run_ctx = RunContext {
run_id,
cancel: CancellationToken::new(),
events: tx,
};
scheduler.init_run_with(run_id, run_ctx.events.clone());
let handle = tokio::runtime::Handle::current();
let rt = Runtime::new(
scheduler,
run_ctx,
serde_json::json!({}),
ExecLimits::default(),
None,
handle,
)
.expect("runtime init");
let schema_json = serde_json::to_string(&schema).unwrap();
let script = format!(
r#"
function main()
local result = agent({{
prompt = "analyze result_collector.rs",
model = "fake-acp",
backend = "fake-acp",
timeout_ms = 10000,
schema = '{schema_json}'
}})
report({{
ok = result.ok,
status = result.status,
output = result.output
}})
end
"#
);
let result = tokio::task::spawn_blocking(move || rt.execute(&script))
.await
.expect("join")
.expect("script ok");
std::env::remove_var("FAKE_ACP_RAW_INPUT");
result
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn e2e_workflow_validate_schema_tool_call_returns_valid_schema_data() {
let schema = serde_json::json!({
"type": "object",
"properties": {
"file": {"type": "string"},
"kind": {"type": "string"},
"summary": {"type": "string"}
},
"required": ["file", "kind", "summary"]
});
let raw_input = serde_json::json!({
"file": "src/adapters/result_collector.rs",
"kind": "rust",
"summary": "collects agent results"
});
let report = run_with_fake_acp(schema, raw_input.clone()).await;
assert_eq!(report["ok"], true, "report: {report}");
assert_eq!(report["status"], "ok");
assert_eq!(report["output"], raw_input);
}