Skip to main content

agentic_server/
agentic_process.rs

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
10/// Reasoning effort passed to Claude Code unless `AGENTIC_CLAUDE_EFFORT` overrides it.
11///
12/// Qwen chat templates served by vLLM accept `low`, `medium`, and `xhigh`; Claude Code's
13/// default of `high` is rejected by the template, so the CLI always pins a compatible value.
14pub 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
97/// Resolve the harness model: the explicit `--model`, or the first model the upstream serves.
98///
99/// # Errors
100///
101/// Returns a configuration error when no model is given and the upstream lists none.
102pub 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
142/// Wait until the gateway is live and, unless skipped, its upstream is ready.
143///
144/// # Errors
145///
146/// Returns a configuration error when the timeout expires.
147pub 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
179/// Run one gateway-plus-harness session and return the harness exit status.
180///
181/// # Errors
182///
183/// Returns an error when a child cannot start or readiness fails.
184pub 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        // Claude Code gives CLAUDE_CODE_EFFORT_LEVEL precedence over --effort, so set both
307        // to keep an inherited `high` from reaching the Qwen chat template.
308        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}