Skip to main content

mail4agent_messenger_shell/provider/
spawn.rs

1//! Last resort: start a NEW provider process that resumes the session for
2//! one turn. Only for a session marked headless (no client holds it open).
3//! A session that has a live client is never spawned into: two writers on
4//! one transcript is how sessions get corrupted.
5//!
6//! | Kind | Command (prompt never in argv when the CLI reads stdin) |
7//! |---|---|
8//! | grok/local | `grok --resume <id> -p <prompt>` |
9//! | codex/local | `codex exec resume <id> -` (prompt on stdin) |
10//! | kimi_code/local | `kimi -S <id> -p <prompt>` |
11//! | claude_code/local | `claude --resume <id> -p` (prompt on stdin) |
12//! | cursor/local | `agent --resume <id> -p <prompt>` |
13//! | codex/web | `codex cloud exec --env <env> <prompt>` (new cloud task) |
14//!
15//! The binary is `M4A_<PROVIDER>_BIN` or the vendor name on `PATH`.
16
17use std::io::Write;
18use std::path::PathBuf;
19use std::process::{Command, Stdio};
20
21use super::{
22    wake_prompt, ProviderKind, ProviderSession, SessionKind, Surface, WakeAdapter, WakeError,
23    WakeLetter, WakeOutcome,
24};
25
26/// Env var naming the cloud environment for `codex cloud exec --env`.
27pub const CODEX_CLOUD_ENV_ENV: &str = "M4A_CODEX_CLOUD_ENV";
28
29/// Env override for a provider binary (`M4A_CODEX_BIN`, ...).
30pub fn bin_env(provider: ProviderKind) -> &'static str {
31    match provider {
32        ProviderKind::Grok => "M4A_GROK_BIN",
33        ProviderKind::KimiCode => "M4A_KIMI_BIN",
34        ProviderKind::ClaudeCode => "M4A_CLAUDE_BIN",
35        ProviderKind::Codex => "M4A_CODEX_BIN",
36        ProviderKind::Cursor => "M4A_CURSOR_AGENT_BIN",
37    }
38}
39
40/// Default program name on `PATH`.
41pub fn default_program(provider: ProviderKind) -> &'static str {
42    match provider {
43        ProviderKind::Grok => "grok",
44        ProviderKind::KimiCode => "kimi",
45        ProviderKind::ClaudeCode => "claude",
46        ProviderKind::Codex => "codex",
47        ProviderKind::Cursor => "agent",
48    }
49}
50
51/// Argv (after the program) and the stdin payload for one resume turn.
52pub fn resume_command(
53    kind: SessionKind,
54    session_id: &str,
55    prompt: &str,
56    cloud_env: Option<&str>,
57) -> Result<(Vec<String>, Option<String>), WakeError> {
58    let id = session_id.to_string();
59    let p = prompt.to_string();
60    Ok(match (kind.surface, kind.provider) {
61        (Surface::Local, ProviderKind::Grok) => (vec!["--resume".into(), id, "-p".into(), p], None),
62        (Surface::Local, ProviderKind::Codex) => (
63            vec!["exec".into(), "resume".into(), id, "-".into()],
64            Some(p),
65        ),
66        (Surface::Local, ProviderKind::KimiCode) => (vec!["-S".into(), id, "-p".into(), p], None),
67        (Surface::Local, ProviderKind::ClaudeCode) => {
68            (vec!["--resume".into(), id, "-p".into()], Some(p))
69        }
70        (Surface::Local, ProviderKind::Cursor) => {
71            (vec!["--resume".into(), id, "-p".into(), p], None)
72        }
73        (Surface::Web, ProviderKind::Codex) => {
74            let Some(env) = cloud_env.filter(|env| !env.is_empty()) else {
75                return Err(WakeError::Unavailable(format!(
76                    "{CODEX_CLOUD_ENV_ENV} unset"
77                )));
78            };
79            (
80                vec![
81                    "cloud".into(),
82                    "exec".into(),
83                    "--env".into(),
84                    env.to_string(),
85                    p,
86                ],
87                None,
88            )
89        }
90        _ => return Err(WakeError::NoInboundTrigger(kind)),
91    })
92}
93
94/// Spawns one resume turn for a headless session. Does not wait for it.
95pub struct ResumeSpawnAdapter {
96    kind: SessionKind,
97    program: Option<PathBuf>,
98    cloud_env: Option<String>,
99}
100
101impl ResumeSpawnAdapter {
102    /// Adapter for `kind`. `program` overrides the binary.
103    pub fn new(kind: SessionKind, program: Option<PathBuf>, cloud_env: Option<String>) -> Self {
104        Self {
105            kind,
106            program,
107            cloud_env,
108        }
109    }
110
111    /// Program from `M4A_<PROVIDER>_BIN` or `PATH`.
112    pub fn from_env(kind: SessionKind) -> Self {
113        let program = std::env::var_os(bin_env(kind.provider))
114            .filter(|value| !value.is_empty())
115            .map(PathBuf::from);
116        let cloud_env = std::env::var(CODEX_CLOUD_ENV_ENV)
117            .ok()
118            .filter(|v| !v.is_empty());
119        Self::new(kind, program, cloud_env)
120    }
121
122    fn program(&self) -> PathBuf {
123        self.program
124            .clone()
125            .unwrap_or_else(|| PathBuf::from(default_program(self.kind.provider)))
126    }
127}
128
129impl WakeAdapter for ResumeSpawnAdapter {
130    fn kind(&self) -> SessionKind {
131        self.kind
132    }
133
134    fn probe(&self, session: &ProviderSession) -> Result<(), WakeError> {
135        if !session.headless {
136            return Err(WakeError::Unavailable(
137                "session has a live client; resume spawn is headless-only".into(),
138            ));
139        }
140        resume_command(
141            self.kind,
142            &session.session_id,
143            "",
144            self.cloud_env.as_deref(),
145        )
146        .map(|_| ())
147    }
148
149    fn wake(
150        &mut self,
151        session: &ProviderSession,
152        letter: &WakeLetter<'_>,
153    ) -> Result<WakeOutcome, WakeError> {
154        self.probe(session)?;
155        let prompt = wake_prompt(session, letter);
156        let (args, stdin) = resume_command(
157            self.kind,
158            &session.session_id,
159            &prompt,
160            self.cloud_env.as_deref(),
161        )?;
162        let mut command = Command::new(self.program());
163        command
164            .args(&args)
165            .stdin(if stdin.is_some() {
166                Stdio::piped()
167            } else {
168                Stdio::null()
169            })
170            .stdout(Stdio::null())
171            .stderr(Stdio::null());
172        if let Some(cwd) = session.cwd.as_deref().filter(|cwd| cwd.is_dir()) {
173            command.current_dir(cwd);
174        }
175        let mut child = command
176            .spawn()
177            .map_err(|err| WakeError::Unavailable(format!("spawn: {}", err.kind())))?;
178        if let (Some(text), Some(mut pipe)) = (stdin, child.stdin.take()) {
179            pipe.write_all(text.as_bytes())
180                .map_err(|err| WakeError::Transport(format!("spawn stdin: {}", err.kind())))?;
181        }
182        // Reap in the background; the turn runs on its own.
183        std::thread::spawn(move || {
184            let _ = child.wait();
185        });
186        Ok(WakeOutcome::Delivered)
187    }
188}
189
190#[cfg(test)]
191mod tests {
192    use super::*;
193    use crate::provider::tests::{letter, session};
194
195    #[test]
196    fn commands_keep_prompt_off_argv_where_the_cli_reads_stdin() {
197        let (args, stdin) =
198            resume_command(SessionKind::local(ProviderKind::Codex), "t1", "hi", None).unwrap();
199        assert_eq!(args, ["exec", "resume", "t1", "-"]);
200        assert_eq!(stdin.as_deref(), Some("hi"));
201        let (args, stdin) = resume_command(
202            SessionKind::local(ProviderKind::ClaudeCode),
203            "c1",
204            "hi",
205            None,
206        )
207        .unwrap();
208        assert!(!args.iter().any(|a| a == "hi") && stdin.is_some());
209        assert!(resume_command(SessionKind::web(ProviderKind::Codex), "x", "hi", None).is_err());
210        assert!(resume_command(SessionKind::web(ProviderKind::KimiCode), "x", "hi", None).is_err());
211    }
212
213    #[test]
214    fn never_spawns_into_a_session_with_a_live_client() {
215        let kind = SessionKind::local(ProviderKind::KimiCode);
216        let mut adapter = ResumeSpawnAdapter::new(kind, Some("/bin/false".into()), None);
217        let s = session(kind);
218        assert!(!s.headless);
219        let err = adapter.wake(&s, &letter("x")).unwrap_err().to_string();
220        assert!(err.contains("headless-only"), "{err}");
221    }
222
223    #[cfg(unix)]
224    #[test]
225    fn headless_session_gets_one_resume_process() {
226        let dir = std::env::temp_dir().join(format!("m4a-spawn-{}", std::process::id()));
227        let _ = std::fs::create_dir_all(&dir);
228        let out = dir.join("argv");
229        let script = dir.join("fake-codex");
230        std::fs::write(
231            &script,
232            format!(
233                "#!/bin/sh\necho \"$@\" > {0}.tmp\ncat >> {0}.tmp\nmv {0}.tmp {0}\n",
234                out.display()
235            ),
236        )
237        .unwrap();
238        use std::os::unix::fs::PermissionsExt;
239        std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap();
240        let kind = SessionKind::local(ProviderKind::Codex);
241        let mut adapter = ResumeSpawnAdapter::new(kind, Some(script), None);
242        let mut s = session(kind);
243        s.headless = true;
244        assert_eq!(
245            adapter.wake(&s, &letter("ping")).unwrap(),
246            WakeOutcome::Delivered
247        );
248        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
249        while !out.exists() && std::time::Instant::now() < deadline {
250            std::thread::sleep(std::time::Duration::from_millis(20));
251        }
252        let text = std::fs::read_to_string(&out).unwrap();
253        assert!(text.starts_with("exec resume s-1 -"));
254        assert!(text.trim_end().ends_with("ping"));
255        let _ = std::fs::remove_dir_all(&dir);
256    }
257}