1use std::{ffi::OsString, path::Path, time::Duration};
2
3use agentic_core::error::Error;
4use reqwest::Client;
5use serde::Deserialize;
6use tokio::time::{Instant, sleep};
7
8use crate::agentic_cli::{CommonOptions, SourceOptions};
9
10pub const DEFAULT_CLAUDE_EFFORT: &str = "medium";
15const CLAUDE_EFFORT_ENV: &str = "AGENTIC_CLAUDE_EFFORT";
16const PLACEHOLDER_MODEL: &str = "agentic-api";
17
18#[must_use]
19pub fn server_args(source: &SourceOptions, common: &CommonOptions) -> Vec<OsString> {
20 let mut args = Vec::new();
21 if let Some(upstream) = &source.upstream {
22 args.extend([OsString::from("--llm-api-base"), OsString::from(upstream)]);
23 } else if let Some(model) = &source.model {
24 args.extend([OsString::from("serve"), OsString::from(model)]);
25 args.extend([OsString::from("--port"), OsString::from(source.llm_port.to_string())]);
26 }
27 args.extend([
28 OsString::from("--gateway-host"),
29 OsString::from(&common.gateway_host),
30 OsString::from("--gateway-port"),
31 OsString::from(common.gateway_port.to_string()),
32 OsString::from("--db-url"),
33 OsString::from(&common.database_url),
34 OsString::from("--llm-ready-timeout-s"),
35 OsString::from(common.llm_ready_timeout_s.to_string()),
36 OsString::from("--llm-ready-interval-s"),
37 OsString::from(common.llm_ready_interval_s.to_string()),
38 ]);
39 if let Some(api_key) = &common.api_key {
40 args.extend([OsString::from("--openai-api-key"), OsString::from(api_key)]);
41 }
42 if common.skip_llm_ready_check {
43 args.push(OsString::from("--skip-llm-ready-check"));
44 }
45 args
46}
47
48#[must_use]
49pub fn server_binary_path(current_exe: &Path) -> std::path::PathBuf {
50 current_exe.with_file_name("agentic-server")
51}
52
53#[must_use]
54pub fn claude_effort() -> String {
55 std::env::var(CLAUDE_EFFORT_ENV)
56 .ok()
57 .map(|value| value.trim().to_owned())
58 .filter(|value| !value.is_empty())
59 .unwrap_or_else(|| DEFAULT_CLAUDE_EFFORT.to_owned())
60}
61
62fn harness_launch_args(
63 harness: crate::agentic_cli::Harness,
64 yolo: bool,
65 claude_effort: &str,
66 passthrough: &[String],
67) -> Vec<String> {
68 let mut args = Vec::with_capacity(passthrough.len() + 3);
69 match harness {
70 crate::agentic_cli::Harness::Codex => {
71 if yolo {
72 args.push("--dangerously-bypass-approvals-and-sandbox".to_owned());
73 }
74 }
75 crate::agentic_cli::Harness::Claude => {
76 if yolo {
77 args.push("--dangerously-skip-permissions".to_owned());
78 }
79 args.extend(["--effort".to_owned(), claude_effort.to_owned()]);
80 }
81 }
82 args.extend_from_slice(passthrough);
83 args
84}
85
86#[derive(Debug, Deserialize)]
87struct ModelList {
88 #[serde(default)]
89 data: Vec<ModelEntry>,
90}
91
92#[derive(Debug, Deserialize)]
93struct ModelEntry {
94 id: String,
95}
96
97pub async fn resolve_model(client: &Client, source: &SourceOptions, api_key: Option<&str>) -> Result<String, Error> {
103 if let Some(model) = &source.model {
104 return Ok(model.clone());
105 }
106 let Some(upstream) = &source.upstream else {
107 return Ok(PLACEHOLDER_MODEL.to_owned());
108 };
109 let models_url = format!("{}/v1/models", agentic_core::config::normalize_base_url(upstream));
110 let mut request = client.get(&models_url);
111 if let Some(api_key) = api_key {
112 request = request.bearer_auth(api_key);
113 }
114 let response = request
115 .send()
116 .await
117 .map_err(|error| Error::Config(format!("failed to list upstream models at {models_url}: {error}")))?
118 .error_for_status()
119 .map_err(|error| Error::Config(format!("upstream model listing at {models_url} failed: {error}")))?;
120 let body = response
121 .text()
122 .await
123 .map_err(|error| Error::Config(format!("failed to read model listing from {models_url}: {error}")))?;
124 let list: ModelList = agentic_core::utils::common::deserialize_from_str(&body)
125 .map_err(|error| Error::Config(format!("invalid model listing from {models_url}: {error}")))?;
126 let mut ids = list.data.into_iter().map(|entry| entry.id);
127 let Some(model) = ids.next() else {
128 return Err(Error::Config(format!(
129 "upstream {upstream} serves no models; pass --model explicitly"
130 )));
131 };
132 let remaining = ids.count();
133 if remaining > 0 {
134 eprintln!(
135 "upstream serves {} models; using {model}. Pass --model to choose another.",
136 remaining + 1
137 );
138 }
139 Ok(model)
140}
141
142pub async fn wait_for_gateway(
148 client: &Client,
149 gateway_url: &str,
150 timeout: Duration,
151 interval: Duration,
152 skip_llm_ready_check: bool,
153) -> Result<(), Error> {
154 let deadline = Instant::now() + timeout;
155 let health_url = format!("{}/health", gateway_url.trim_end_matches('/'));
156 let ready_url = format!("{}/ready", gateway_url.trim_end_matches('/'));
157 loop {
158 if Instant::now() >= deadline {
159 return Err(Error::Config(format!("gateway did not become ready at {gateway_url}")));
160 }
161 let health_ok = client
162 .get(&health_url)
163 .send()
164 .await
165 .is_ok_and(|response| response.status().is_success());
166 let ready_ok = skip_llm_ready_check
167 || client
168 .get(&ready_url)
169 .send()
170 .await
171 .is_ok_and(|response| response.status().is_success());
172 if health_ok && ready_ok {
173 return Ok(());
174 }
175 sleep(interval).await;
176 }
177}
178
179pub async fn run_session(
185 current_exe: &Path,
186 harness: crate::agentic_cli::Harness,
187 options: crate::agentic_cli::HarnessOptions,
188) -> Result<std::process::ExitStatus, Error> {
189 let gateway_url = format!("http://{}:{}", options.common.gateway_host, options.common.gateway_port);
190 let session_root = std::env::temp_dir().join(format!(
191 "agentic-api-session-{}-{}",
192 std::process::id(),
193 std::time::SystemTime::now()
194 .duration_since(std::time::UNIX_EPOCH)
195 .map_or(0, |duration| duration.as_nanos())
196 ));
197 tokio::fs::create_dir_all(&session_root).await?;
198
199 let mut server = start_server(current_exe, &options)?;
200
201 let client = match Client::builder()
202 .timeout(Duration::from_secs(2))
203 .build()
204 .map_err(Error::HttpClient)
205 {
206 Ok(client) => client,
207 Err(error) => {
208 cleanup(&mut server, &session_root).await;
209 return Err(error);
210 }
211 };
212 if let Err(error) = wait_for_gateway(
213 &client,
214 &gateway_url,
215 Duration::from_secs_f64(options.common.llm_ready_timeout_s),
216 Duration::from_secs_f64(options.common.llm_ready_interval_s),
217 options.common.skip_llm_ready_check,
218 )
219 .await
220 {
221 cleanup(&mut server, &session_root).await;
222 return Err(error);
223 }
224
225 let model = match resolve_model(&client, &options.source, options.common.api_key.as_deref()).await {
226 Ok(model) => model,
227 Err(error) => {
228 cleanup(&mut server, &session_root).await;
229 return Err(error);
230 }
231 };
232 let harness_env = match harness_environment(harness, &gateway_url, &model, &options, &session_root) {
233 Ok(environment) => environment,
234 Err(error) => {
235 cleanup(&mut server, &session_root).await;
236 return Err(error);
237 }
238 };
239 if !options.common.quiet {
240 println!("{}", harness_env.summary);
241 }
242
243 let mut harness_child = match spawn_harness(harness, &options, &harness_env) {
244 Ok(child) => child,
245 Err(error) => {
246 cleanup(&mut server, &session_root).await;
247 return Err(error);
248 }
249 };
250
251 let harness_status = tokio::select! {
252 status = harness_child.wait() => status?,
253 signal = tokio::signal::ctrl_c() => {
254 signal?;
255 let _ = harness_child.kill().await;
256 harness_child.wait().await?
257 }
258 };
259 cleanup(&mut server, &session_root).await;
260 Ok(harness_status)
261}
262
263fn start_server(
264 current_exe: &Path,
265 options: &crate::agentic_cli::HarnessOptions,
266) -> Result<tokio::process::Child, Error> {
267 let server_path = server_binary_path(current_exe);
268 if !server_path.is_file() {
269 return Err(Error::Config(format!(
270 "agentic-server binary not found beside {}; run cargo build -p agentic-server --bins first",
271 current_exe.display()
272 )));
273 }
274 let mut server = tokio::process::Command::new(server_path);
275 server.args(server_args(&options.source, &options.common));
276 server.stdout(std::process::Stdio::inherit());
277 server.stderr(std::process::Stdio::inherit());
278 Ok(server.spawn()?)
279}
280
281fn harness_environment(
282 harness: crate::agentic_cli::Harness,
283 gateway_url: &str,
284 model: &str,
285 options: &crate::agentic_cli::HarnessOptions,
286 session_root: &Path,
287) -> Result<crate::agentic_harness::HarnessEnv, Error> {
288 if matches!(harness, crate::agentic_cli::Harness::Claude) {
289 crate::agentic_harness::validate_claude_model(model).map_err(Error::Config)?;
290 }
291 let mut environment = match harness {
292 crate::agentic_cli::Harness::Codex => crate::agentic_harness::prepare_codex_home(
293 session_root,
294 gateway_url,
295 model,
296 options.common.api_key.as_deref(),
297 )
298 .map_err(Error::from),
299 crate::agentic_cli::Harness::Claude => Ok(crate::agentic_harness::prepare_claude_env(
300 gateway_url,
301 model,
302 options.common.api_key.as_deref(),
303 )),
304 }?;
305 if matches!(harness, crate::agentic_cli::Harness::Claude) {
306 environment
309 .environment
310 .insert("CLAUDE_CODE_EFFORT_LEVEL".to_owned(), claude_effort());
311 }
312 Ok(environment)
313}
314
315fn spawn_harness(
316 harness: crate::agentic_cli::Harness,
317 options: &crate::agentic_cli::HarnessOptions,
318 harness_env: &crate::agentic_harness::HarnessEnv,
319) -> Result<tokio::process::Child, Error> {
320 let binary_name = match harness {
321 crate::agentic_cli::Harness::Codex => "codex",
322 crate::agentic_cli::Harness::Claude => "claude",
323 };
324 let override_name = match harness {
325 crate::agentic_cli::Harness::Codex => "AGENTIC_CODEX_BIN",
326 crate::agentic_cli::Harness::Claude => "AGENTIC_CLAUDE_BIN",
327 };
328 let binary = std::env::var_os(override_name).unwrap_or_else(|| binary_name.into());
329 let mut harness_command = tokio::process::Command::new(binary);
330 harness_command.args(harness_launch_args(
331 harness,
332 options.common.yolo,
333 &claude_effort(),
334 &options.harness_args,
335 ));
336 harness_command.envs(&harness_env.environment);
337 harness_command.stdin(std::process::Stdio::inherit());
338 harness_command.stdout(std::process::Stdio::inherit());
339 harness_command.stderr(std::process::Stdio::inherit());
340 harness_command
341 .spawn()
342 .map_err(|error| Error::Config(format!("failed to launch {binary_name} ({override_name}): {error}")))
343}
344
345async fn cleanup(server: &mut tokio::process::Child, session_root: &Path) {
346 let _ = server.kill().await;
347 let _ = server.wait().await;
348 let _ = tokio::fs::remove_dir_all(session_root).await;
349}
350
351#[cfg(test)]
352mod tests {
353 use std::ffi::OsString;
354
355 use super::{DEFAULT_CLAUDE_EFFORT, harness_launch_args, server_args};
356 use crate::agentic_cli::{CommonOptions, Harness, SourceOptions};
357
358 #[test]
359 fn integrated_mode_builds_server_arguments() {
360 let args = server_args(
361 &SourceOptions {
362 upstream: None,
363 model: Some("Qwen/test".to_owned()),
364 llm_port: 8000,
365 },
366 &CommonOptions::default(),
367 );
368 let args: Vec<_> = args.iter().map(OsString::as_os_str).collect();
369
370 assert_eq!(args[0], "serve");
371 assert_eq!(args[1], "Qwen/test");
372 assert!(
373 args.windows(2)
374 .any(|pair| pair == ["--db-url", "sqlite://./agentic_api.db"])
375 );
376 }
377
378 #[test]
379 fn standalone_mode_builds_upstream_arguments() {
380 let args = server_args(
381 &SourceOptions {
382 upstream: Some("http://127.0.0.1:8000".to_owned()),
383 model: None,
384 llm_port: 8000,
385 },
386 &CommonOptions::default(),
387 );
388 let args: Vec<_> = args.iter().map(OsString::as_os_str).collect();
389
390 assert_eq!(args[0], "--llm-api-base");
391 assert_eq!(args[1], "http://127.0.0.1:8000");
392 }
393
394 #[test]
395 fn explicit_upstream_wins_when_model_names_the_harness_model() {
396 let args = server_args(
397 &SourceOptions {
398 upstream: Some("http://127.0.0.1:8000".to_owned()),
399 model: Some("Qwen/test".to_owned()),
400 llm_port: 8000,
401 },
402 &CommonOptions::default(),
403 );
404 let args: Vec<_> = args.iter().map(OsString::as_os_str).collect();
405
406 assert_eq!(args[0], "--llm-api-base");
407 assert!(!args.iter().any(|arg| *arg == "serve"));
408 }
409
410 #[test]
411 fn yolo_mode_uses_native_codex_bypass_flag() {
412 assert_eq!(
413 harness_launch_args(Harness::Codex, true, DEFAULT_CLAUDE_EFFORT, &["exec".to_owned()]),
414 ["--dangerously-bypass-approvals-and-sandbox", "exec"]
415 );
416 }
417
418 #[test]
419 fn yolo_mode_uses_native_claude_bypass_and_compatible_effort() {
420 assert_eq!(
421 harness_launch_args(Harness::Claude, true, DEFAULT_CLAUDE_EFFORT, &[]),
422 ["--dangerously-skip-permissions", "--effort", "medium"]
423 );
424 }
425
426 #[test]
427 fn claude_always_receives_a_compatible_effort() {
428 assert_eq!(
429 harness_launch_args(Harness::Claude, false, "low", &["-p".to_owned(), "hi".to_owned()]),
430 ["--effort", "low", "-p", "hi"]
431 );
432 assert_eq!(
433 harness_launch_args(Harness::Codex, false, DEFAULT_CLAUDE_EFFORT, &[]),
434 Vec::<String>::new()
435 );
436 }
437
438 #[test]
439 fn claude_environment_pins_effort_without_yolo() {
440 let options = crate::agentic_cli::HarnessOptions {
441 source: SourceOptions {
442 upstream: Some("http://127.0.0.1:8000".to_owned()),
443 model: None,
444 llm_port: 8000,
445 },
446 common: CommonOptions::default(),
447 harness_args: Vec::new(),
448 };
449 let root = std::env::temp_dir().join(format!("agentic-api-effort-test-{}", std::process::id()));
450 let environment = super::harness_environment(
451 Harness::Claude,
452 "http://127.0.0.1:3000",
453 "served-discovered",
454 &options,
455 &root,
456 )
457 .expect("Claude environment");
458
459 assert_eq!(
460 environment.environment.get("CLAUDE_CODE_EFFORT_LEVEL"),
461 Some(&DEFAULT_CLAUDE_EFFORT.to_owned())
462 );
463 assert_eq!(
464 environment.environment.get("ANTHROPIC_MODEL"),
465 Some(&"served-discovered".to_owned())
466 );
467 }
468
469 #[tokio::test]
470 async fn resolve_model_does_not_double_the_v1_suffix() {
471 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
472 let address = listener.local_addr().unwrap();
473 tokio::spawn(async move {
474 use tokio::io::{AsyncReadExt, AsyncWriteExt};
475 let (mut socket, _) = listener.accept().await.unwrap();
476 let mut buffer = [0_u8; 1024];
477 let read = socket.read(&mut buffer).await.unwrap();
478 let request_line = String::from_utf8_lossy(&buffer[..read])
479 .lines()
480 .next()
481 .unwrap_or_default()
482 .to_owned();
483 let body = r#"{"data":[{"id":"Qwen/served"}]}"#;
484 let status = if request_line.starts_with("GET /v1/models ") {
485 "200 OK"
486 } else {
487 "404 Not Found"
488 };
489 let response = format!(
490 "HTTP/1.1 {status}\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
491 body.len()
492 );
493 socket.write_all(response.as_bytes()).await.unwrap();
494 });
495 let client = reqwest::Client::new();
496 let source = SourceOptions {
497 upstream: Some(format!("http://{address}/v1")),
498 model: None,
499 llm_port: 8000,
500 };
501 assert_eq!(
502 super::resolve_model(&client, &source, None).await.unwrap(),
503 "Qwen/served"
504 );
505 }
506
507 #[tokio::test]
508 async fn resolve_model_prefers_explicit_model() {
509 let client = reqwest::Client::new();
510 let source = SourceOptions {
511 upstream: Some("http://127.0.0.1:9".to_owned()),
512 model: Some("Qwen/test".to_owned()),
513 llm_port: 8000,
514 };
515 assert_eq!(super::resolve_model(&client, &source, None).await.unwrap(), "Qwen/test");
516 }
517
518 #[tokio::test]
519 async fn resolve_model_discovers_first_upstream_model() {
520 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
521 let address = listener.local_addr().unwrap();
522 tokio::spawn(async move {
523 use tokio::io::{AsyncReadExt, AsyncWriteExt};
524 let (mut socket, _) = listener.accept().await.unwrap();
525 let mut buffer = [0_u8; 1024];
526 let _ = socket.read(&mut buffer).await;
527 let body = r#"{"object":"list","data":[{"id":"Qwen/served"},{"id":"other"}]}"#;
528 let response = format!(
529 "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
530 body.len()
531 );
532 socket.write_all(response.as_bytes()).await.unwrap();
533 });
534 let client = reqwest::Client::new();
535 let source = SourceOptions {
536 upstream: Some(format!("http://{address}")),
537 model: None,
538 llm_port: 8000,
539 };
540 assert_eq!(
541 super::resolve_model(&client, &source, None).await.unwrap(),
542 "Qwen/served"
543 );
544 }
545
546 #[test]
547 fn yolo_claude_environment_overrides_inherited_effort() {
548 let options = crate::agentic_cli::HarnessOptions {
549 source: SourceOptions {
550 upstream: Some("http://127.0.0.1:8000".to_owned()),
551 model: Some("served-test".to_owned()),
552 llm_port: 8000,
553 },
554 common: CommonOptions {
555 yolo: true,
556 ..CommonOptions::default()
557 },
558 harness_args: Vec::new(),
559 };
560 let root = std::env::temp_dir().join(format!("agentic-api-yolo-test-{}", std::process::id()));
561 let environment =
562 super::harness_environment(Harness::Claude, "http://127.0.0.1:3000", "served-test", &options, &root)
563 .expect("Claude environment");
564
565 assert_eq!(
566 environment.environment.get("CLAUDE_CODE_EFFORT_LEVEL"),
567 Some(&"medium".to_owned())
568 );
569 }
570}