mail4agent_messenger_shell/provider/
spawn.rs1use 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
26pub const CODEX_CLOUD_ENV_ENV: &str = "M4A_CODEX_CLOUD_ENV";
28
29pub 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
40pub 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
51pub 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
94pub struct ResumeSpawnAdapter {
96 kind: SessionKind,
97 program: Option<PathBuf>,
98 cloud_env: Option<String>,
99}
100
101impl ResumeSpawnAdapter {
102 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 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 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}