Skip to main content

rightkit_qa/
engine.rs

1//! Black-box driving of an app's engine binary: one-shot CLI words and the
2//! framed stdio sidecar protocol (newline-delimited JSON through
3//! `rightkit-framed-sidecar`). Nothing mocks the engine.
4use 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    /// Env applied to every engine process (e.g. `GENRIGHT_NO_KEYCHAIN=1`).
22    pub env: Vec<(String, String)>,
23    /// Env var that receives the per-scenario data dir (e.g. `GENRIGHT_APP_DATA`).
24    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    /// Appended as `--input <json>` the way the engine CLI contract expects.
32    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    /// Every JSON object line the engine wrote to stderr (progress events etc.).
43    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    /// `engine run('job status', input)`: words split, non-zero exit or `json.error` fails.
109    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
141/// A running `--stdio` sidecar speaking
142/// `{"type":"request","id":N,"method":{"method":"command","params":{"name","input"}}}`
143/// and answering `{"type":"response","status":"ok|err","id":N,"result|error":...}`.
144pub 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    /// Frames that were neither hello nor the awaited response (events).
153    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, // non-JSON line (log noise): skip, stay on the frame boundary
212                }
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    /// Send an arbitrary frame and return the response with the same id.
263    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    /// A command; `Ok(result)` for `status: ok`, `Err("code: message")` for `err`.
277    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    /// Close stdin, give the engine a moment to exit, then kill the owned tree.
301    /// Returns the exit code if it exited on its own.
302    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
324/// Build the engine through the managed wrapper and return the executable it
325/// produced (from cargo's `compiler-artifact` message), never a stale copy.
326pub 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}