1use crate::process::{run_owned, RunOptions, Tracker};
5use crate::util::{err, Error, Result};
6use rightkit_framed_sidecar::{read_json_line, write_json_line, FrameError};
7use rightkit_process::{OwnedChild, OwnedCommand};
8use serde_json::{json, Value};
9use std::io::{BufReader, Write};
10use std::path::{Path, PathBuf};
11use std::process::{ChildStdin, Command, Stdio};
12use std::sync::mpsc::{channel, Receiver, RecvTimeoutError};
13use std::sync::{Arc, Mutex};
14use std::time::{Duration, Instant};
15
16pub const FRAME_LIMIT: usize = 16 * 1024 * 1024;
17
18#[derive(Debug, Clone)]
19pub struct EngineTarget {
20 pub binary: PathBuf,
21 pub env: Vec<(String, String)>,
23 pub data_dir_env: Option<String>,
25}
26
27#[derive(Debug, Clone, Default)]
28pub struct CliOptions {
29 pub data_dir: Option<PathBuf>,
30 pub env: Vec<(String, String)>,
31 pub input: Option<Value>,
33 pub timeout: Option<Duration>,
34}
35
36#[derive(Debug, Clone)]
37pub struct CliResult {
38 pub code: Option<i32>,
39 pub json: Value,
40 pub stdout: String,
41 pub stderr: String,
42 pub events: Vec<Value>,
44 pub timed_out: bool,
45}
46
47impl CliResult {
48 pub fn to_value(&self) -> Value {
49 json!({"code": self.code, "json": self.json, "stdout": self.stdout, "stderr": self.stderr, "events": self.events, "timed_out": self.timed_out})
50 }
51}
52
53impl EngineTarget {
54 fn env_for(
55 &self,
56 data_dir: Option<&Path>,
57 extra: &[(String, String)],
58 ) -> Vec<(String, String)> {
59 let mut env = self.env.clone();
60 if let (Some(var), Some(dir)) = (&self.data_dir_env, data_dir) {
61 env.push((var.clone(), dir.to_string_lossy().into()));
62 }
63 env.extend(extra.iter().cloned());
64 env
65 }
66
67 pub fn cli(&self, args: &[String], opts: &CliOptions, tracker: &Tracker) -> Result<CliResult> {
68 if !self.binary.exists() {
69 return err(format!(
70 "engine binary not found: {}",
71 self.binary.display()
72 ));
73 }
74 let mut argv = args.to_vec();
75 if let Some(input) = &opts.input {
76 argv.push("--input".into());
77 argv.push(input.to_string());
78 }
79 let out = run_owned(
80 &self.binary.to_string_lossy(),
81 &argv,
82 &RunOptions {
83 env: self.env_for(opts.data_dir.as_deref(), &opts.env),
84 timeout: opts.timeout,
85 label: format!("engine {}", args.join(" ")),
86 ..Default::default()
87 },
88 tracker,
89 )?;
90 let json = serde_json::from_str(out.stdout.trim()).unwrap_or(Value::Null);
91 let events = out
92 .stderr
93 .lines()
94 .map(str::trim)
95 .filter(|l| l.starts_with('{'))
96 .filter_map(|l| serde_json::from_str(l).ok())
97 .collect();
98 Ok(CliResult {
99 code: out.code,
100 json,
101 stdout: out.stdout,
102 stderr: out.stderr,
103 events,
104 timed_out: out.timed_out,
105 })
106 }
107
108 pub fn run(
110 &self,
111 words: &str,
112 input: Value,
113 opts: &CliOptions,
114 tracker: &Tracker,
115 ) -> Result<Value> {
116 let args: Vec<String> = words.split_whitespace().map(String::from).collect();
117 let r = self.cli(
118 &args,
119 &CliOptions {
120 input: Some(input),
121 ..opts.clone()
122 },
123 tracker,
124 )?;
125 if r.code != Some(0) || r.json.get("error").map(|e| !e.is_null()).unwrap_or(false) {
126 return err(format!(
127 "engine command failed ({:?}): {words}\nstdout: {}\nstderr: {}",
128 r.code,
129 r.stdout,
130 crate::util::tail(&r.stderr, 2000)
131 ));
132 }
133 Ok(r.json)
134 }
135
136 pub fn start_sidecar(&self, opts: &CliOptions, tracker: &Tracker) -> Result<Sidecar> {
137 Sidecar::start(self, opts, tracker)
138 }
139}
140
141pub struct Sidecar {
145 child: Option<OwnedChild>,
146 stdin: Option<ChildStdin>,
147 rx: Receiver<Value>,
148 seq: u64,
149 pub hello: Value,
150 pub pid: u32,
151 stderr: Arc<Mutex<String>>,
152 pub events: Vec<Value>,
154 tracker: Tracker,
155}
156
157impl Sidecar {
158 fn start(engine: &EngineTarget, opts: &CliOptions, tracker: &Tracker) -> Result<Self> {
159 let mut cmd = Command::new(&engine.binary);
160 cmd.arg("--stdio")
161 .stdin(Stdio::piped())
162 .stdout(Stdio::piped())
163 .stderr(Stdio::piped());
164 for (k, v) in engine.env_for(opts.data_dir.as_deref(), &opts.env) {
165 cmd.env(k, v);
166 }
167 let mut owned = OwnedCommand::from_command(cmd);
168 owned.windows_hide();
169 let mut child = owned.spawn().map_err(|e| {
170 Error(format!(
171 "failed to start sidecar {}: {e}",
172 engine.binary.display()
173 ))
174 })?;
175 let pid = child.id();
176 tracker.register(pid, "engine --stdio");
177 let stdin = child.take_stdin();
178 let stdout = child
179 .take_stdout()
180 .ok_or_else(|| Error("sidecar stdout unavailable".into()))?;
181 let mut stderr_pipe = child.take_stderr();
182 let stderr = Arc::new(Mutex::new(String::new()));
183 let sink = stderr.clone();
184 std::thread::spawn(move || {
185 use std::io::Read;
186 let mut buf = [0u8; 4096];
187 if let Some(p) = stderr_pipe.as_mut() {
188 while let Ok(n) = p.read(&mut buf) {
189 if n == 0 {
190 break;
191 }
192 sink.lock()
193 .unwrap()
194 .push_str(&String::from_utf8_lossy(&buf[..n]));
195 }
196 }
197 });
198 let (tx, rx) = channel();
199 std::thread::spawn(move || {
200 let mut reader = BufReader::new(stdout);
201 let mut buf = Vec::new();
202 loop {
203 match read_json_line::<_, Value>(&mut reader, &mut buf, FRAME_LIMIT) {
204 Ok(Some(v)) => {
205 if tx.send(v).is_err() {
206 break;
207 }
208 }
209 Ok(None) => break,
210 Err(FrameError::Io(_)) => break,
211 Err(_) => continue, }
213 }
214 });
215 let mut me = Self {
216 child: Some(child),
217 stdin,
218 rx,
219 seq: 0,
220 hello: Value::Null,
221 pid,
222 stderr,
223 events: vec![],
224 tracker: tracker.clone(),
225 };
226 me.hello = me.await_frame(Duration::from_secs(30), |f| f["type"] == "hello")?;
227 Ok(me)
228 }
229
230 fn await_frame(
231 &mut self,
232 timeout: Duration,
233 mut want: impl FnMut(&Value) -> bool,
234 ) -> Result<Value> {
235 let deadline = Instant::now() + timeout;
236 loop {
237 let left = deadline.saturating_duration_since(Instant::now());
238 match self.rx.recv_timeout(left.max(Duration::from_millis(1))) {
239 Ok(f) if want(&f) => return Ok(f),
240 Ok(f) => self.events.push(f),
241 Err(RecvTimeoutError::Timeout) => {
242 return err(format!(
243 "sidecar timeout after {}ms; stderr tail: {}",
244 timeout.as_millis(),
245 self.stderr_tail()
246 ))
247 }
248 Err(RecvTimeoutError::Disconnected) => {
249 return err(format!(
250 "sidecar closed its stdout; stderr tail: {}",
251 self.stderr_tail()
252 ))
253 }
254 }
255 }
256 }
257
258 pub fn stderr_tail(&self) -> String {
259 crate::util::tail(&self.stderr.lock().unwrap(), 1500)
260 }
261
262 pub fn raw(&mut self, mut frame: Value, timeout: Duration) -> Result<Value> {
264 self.seq += 1;
265 let id = self.seq;
266 frame["id"] = json!(id);
267 let stdin = self
268 .stdin
269 .as_mut()
270 .ok_or_else(|| Error("sidecar stdin closed".into()))?;
271 write_json_line(stdin, &frame, FRAME_LIMIT).map_err(|e| Error(e.to_string()))?;
272 stdin.flush()?;
273 self.await_frame(timeout, |f| f["type"] == "response" && f["id"] == json!(id))
274 }
275
276 pub fn command(&mut self, name: &str, input: Value, timeout: Duration) -> Result<Value> {
278 let f = self.raw(json!({"type": "request", "method": {"method": "command", "params": {"name": name, "input": input}}}), timeout)?;
279 if f["status"] == "err" {
280 return err(format!(
281 "{}: {}",
282 f["error"]["code"].as_str().unwrap_or("error"),
283 f["error"]["message"].as_str().unwrap_or("")
284 ));
285 }
286 Ok(f["result"].clone())
287 }
288
289 pub fn shutdown(&mut self) {
290 self.seq += 1;
291 if let Some(stdin) = self.stdin.as_mut() {
292 let _ = write_json_line(
293 stdin,
294 &json!({"type":"request","id":self.seq,"method":{"method":"shutdown"}}),
295 FRAME_LIMIT,
296 );
297 }
298 }
299
300 pub fn close(&mut self) -> Option<i32> {
303 self.stdin.take();
304 let mut code = None;
305 if let Some(mut child) = self.child.take() {
306 match child.wait_timeout(Duration::from_secs(5)) {
307 Ok(Some(s)) => code = s.code(),
308 _ => {
309 let _ = child.terminate_tree();
310 }
311 }
312 }
313 self.tracker.forget(self.pid);
314 code
315 }
316}
317
318impl Drop for Sidecar {
319 fn drop(&mut self) {
320 let _ = self.close();
321 }
322}
323
324pub fn build_engine(
327 cwd: &Path,
328 command: &[String],
329 package: &str,
330 env: &[(String, String)],
331 tracker: &Tracker,
332) -> Result<PathBuf> {
333 let Some((program, rest)) = command.split_first() else {
334 return err("engine.build is empty");
335 };
336 let mut args = rest.to_vec();
337 if !args.iter().any(|a| a.starts_with("--message-format")) {
338 args.push("--message-format=json-render-diagnostics".into());
339 }
340 let out = run_owned(
341 program,
342 &args,
343 &RunOptions {
344 cwd: Some(cwd.to_path_buf()),
345 env: env.to_vec(),
346 timeout: Some(Duration::from_secs(45 * 60)),
347 label: "engine build".into(),
348 ..Default::default()
349 },
350 tracker,
351 )?;
352 if out.code != Some(0) {
353 return err(format!(
354 "engine build failed ({:?}):\n{}",
355 out.code,
356 crate::util::tail(&out.stderr, 4000)
357 ));
358 }
359 out.stdout
360 .lines()
361 .filter(|l| l.starts_with('{'))
362 .filter_map(|l| serde_json::from_str::<Value>(l).ok())
363 .filter(|m| m["reason"] == "compiler-artifact" && m["target"]["name"] == package)
364 .filter_map(|m| m["executable"].as_str().map(PathBuf::from))
365 .next_back()
366 .ok_or_else(|| {
367 Error(format!(
368 "build produced no executable for package {package}"
369 ))
370 })
371}