1use super::command::{EngineCommand, ImageBuild};
14use super::{Engine, EngineKind};
15use qframe::i18n::I18n;
16use qframe::runtime::{Line, Process, ProcessOutcome};
17use std::io::Read;
18use std::process::{Child, Command};
19use std::sync::{Arc, OnceLock, mpsc};
20use std::time::{Duration, Instant};
21
22#[derive(Debug, Clone, PartialEq, Eq)]
25pub struct Failure {
26 pub command: EngineCommand,
28 pub code: Option<i32>,
30 pub output: String,
32}
33
34#[derive(Debug)]
39pub enum EngineError {
40 NotRunnable {
42 command: EngineCommand,
44 error: std::io::Error,
46 },
47 Failed(Failure),
49 Cancelled {
51 command: EngineCommand,
53 },
54 TimedOut {
58 command: EngineCommand,
60 after: Duration,
62 },
63}
64
65type Words = Box<dyn Fn() -> Arc<I18n> + Send + Sync>;
67
68static WORDS: OnceLock<Words> = OnceLock::new();
70
71pub fn speak_with(words: impl Fn() -> Arc<I18n> + Send + Sync + 'static) {
76 let _ = WORDS.set(Box::new(words));
77}
78
79#[must_use]
82pub fn timed_out(command: &EngineCommand, after: Duration) -> String {
83 let binary = command.program.file_name().map(|name| name.to_string_lossy().into_owned()).unwrap_or_default();
84 let say = || {
85 let engine = match EngineKind::from_name(&binary) {
86 Some(EngineKind::Podman) => qframe::t!("settings.podman"),
87 Some(EngineKind::Docker) => qframe::t!("settings.docker"),
88 None => binary.clone(),
89 };
90 qframe::t!("engine.timed-out", engine = engine, seconds = after.as_secs().to_string())
91 };
92 let said = say();
94 if !said.starts_with('⟦') {
95 return said;
96 }
97 let words = WORDS.get().map_or_else(|| Arc::new(crate::service::translator(Some("en"))), |words| words());
98 qframe::i18n::scope(words, say)
99}
100
101pub fn capture(command: &EngineCommand) -> Result<String, EngineError> {
110 use std::process::Stdio;
111
112 let child = start(
113 Command::new(&command.program)
114 .args(&command.args)
115 .stdin(Stdio::null())
116 .stdout(Stdio::piped())
117 .stderr(Stdio::piped()),
118 )
119 .map_err(|error| EngineError::NotRunnable { command: command.clone(), error })?;
120 finish(command, child)
121}
122
123pub fn feed(command: &EngineCommand, input: &[u8]) -> Result<String, EngineError> {
133 use std::io::Write;
134 use std::process::Stdio;
135
136 let mut child = start(
137 Command::new(&command.program)
138 .args(&command.args)
139 .stdin(Stdio::piped())
140 .stdout(Stdio::piped())
141 .stderr(Stdio::piped()),
142 )
143 .map_err(|error| EngineError::NotRunnable { command: command.clone(), error })?;
144 if let Some(mut stdin) = child.stdin.take() {
145 let input = input.to_vec();
146 std::thread::spawn(move || {
150 let _ = stdin.write_all(&input);
151 });
152 }
153 finish(command, child)
154}
155
156const BUSY_TRIES: u32 = 20;
158const BUSY_PAUSE: Duration = Duration::from_millis(50);
159
160fn start(command: &mut Command) -> std::io::Result<Child> {
168 let mut tries = 0;
169 loop {
170 match command.spawn() {
171 Err(error) if error.kind() == std::io::ErrorKind::ExecutableFileBusy && tries < BUSY_TRIES => {
172 tries += 1;
173 std::thread::sleep(BUSY_PAUSE);
174 }
175 answer => return answer,
176 }
177 }
178}
179
180const GRACE: Duration = Duration::from_secs(5);
183
184fn finish(command: &EngineCommand, mut child: Child) -> Result<String, EngineError> {
191 let (sender, said) = mpsc::channel();
192 for (index, stream) in [child.stdout.take().map(boxed), child.stderr.take().map(boxed)].into_iter().enumerate() {
193 let sender = sender.clone();
194 std::thread::spawn(move || {
195 let mut bytes = Vec::new();
196 if let Some(mut stream) = stream {
197 let _ = stream.read_to_end(&mut bytes);
198 }
199 let _ = sender.send((index, bytes));
200 });
201 }
202 drop(sender);
203 let started = Instant::now();
204 let left = || command.deadline.map(|deadline| deadline.saturating_sub(started.elapsed()));
205 let mut streams = [Vec::new(), Vec::new()];
206 for _ in 0..2 {
207 let received = match left() {
208 Some(left) => said.recv_timeout(left).map_err(|_| ()),
209 None => said.recv().map_err(|_| ()),
210 };
211 match received {
212 Ok((index, bytes)) => streams[index] = bytes,
213 Err(()) => return Err(stuck(command, child)),
214 }
215 }
216 let status = loop {
217 if command.deadline.is_none() {
218 break child.wait().map_err(|error| EngineError::NotRunnable { command: command.clone(), error })?;
219 }
220 match child.try_wait() {
221 Ok(Some(status)) => break status,
222 Ok(None) if left().is_some_and(|left| left.is_zero()) => return Err(stuck(command, child)),
223 Ok(None) => std::thread::sleep(Duration::from_millis(5)),
224 Err(error) => return Err(EngineError::NotRunnable { command: command.clone(), error }),
225 }
226 };
227 let [out, err] = streams;
228 if status.success() {
229 return Ok(String::from_utf8_lossy(&out).into_owned());
230 }
231 Err(EngineError::Failed(Failure {
232 command: command.clone(),
233 code: status.code(),
234 output: both(&String::from_utf8_lossy(&out), &String::from_utf8_lossy(&err)),
235 }))
236}
237
238fn boxed(stream: impl Read + Send + 'static) -> Box<dyn Read + Send> {
240 Box::new(stream)
241}
242
243fn stuck(command: &EngineCommand, mut child: Child) -> EngineError {
247 #[cfg(unix)]
248 {
249 let _ = rustix::process::kill_process(rustix::process::Pid::from_child(&child), rustix::process::Signal::TERM);
250 let asked = Instant::now();
251 while asked.elapsed() < GRACE {
252 if matches!(child.try_wait(), Ok(Some(_))) {
253 break;
254 }
255 std::thread::sleep(Duration::from_millis(20));
256 }
257 }
258 let _ = child.kill();
259 let _ = child.wait();
260 EngineError::TimedOut { command: command.clone(), after: command.deadline.unwrap_or_default() }
261}
262
263pub fn pipe(from: &EngineCommand, to: &EngineCommand) -> Result<(), EngineError> {
276 use std::process::Stdio;
277
278 let not_runnable = |command: &EngineCommand| {
279 let command = command.clone();
280 move |error| EngineError::NotRunnable { command, error }
281 };
282 let mut reader = Command::new(&from.program)
283 .args(&from.args)
284 .stdin(Stdio::null())
285 .stdout(Stdio::piped())
286 .stderr(Stdio::piped())
287 .spawn()
288 .map_err(not_runnable(from))?;
289 let Some(packed) = reader.stdout.take() else {
290 let _ = reader.kill();
291 return Err(EngineError::NotRunnable { command: from.clone(), error: std::io::Error::other("no output") });
292 };
293 let writer = Command::new(&to.program)
294 .args(&to.args)
295 .stdin(Stdio::from(packed))
296 .stdout(Stdio::piped())
297 .stderr(Stdio::piped())
298 .spawn();
299 let writer = match writer {
300 Ok(writer) => writer,
301 Err(error) => {
302 let _ = reader.kill();
303 let _ = reader.wait();
304 return Err(EngineError::NotRunnable { command: to.clone(), error });
305 }
306 };
307 let said = reader.stderr.take().map(|mut stream| {
308 std::thread::spawn(move || {
309 let mut text = String::new();
310 let _ = stream.read_to_string(&mut text);
311 text
312 })
313 });
314 let written = writer.wait_with_output().map_err(not_runnable(to))?;
315 let read = reader.wait().map_err(not_runnable(from))?;
316 let said = said.and_then(|thread| thread.join().ok()).unwrap_or_default();
317 if !read.success() {
318 return Err(EngineError::Failed(Failure {
319 command: from.clone(),
320 code: read.code(),
321 output: said.trim_end().to_owned(),
322 }));
323 }
324 if !written.status.success() {
325 return Err(EngineError::Failed(Failure {
326 command: to.clone(),
327 code: written.status.code(),
328 output: both(&String::from_utf8_lossy(&written.stdout), &String::from_utf8_lossy(&written.stderr)),
329 }));
330 }
331 Ok(())
332}
333
334fn both(out: &str, err: &str) -> String {
336 match (out.trim_end(), err.trim_end()) {
337 ("", err) => err.to_owned(),
338 (out, "") => out.to_owned(),
339 (out, err) => format!("{out}\n{err}"),
340 }
341}
342
343pub fn stream(
355 command: &EngineCommand,
356 cancel: &dyn Fn() -> bool,
357 line: &mut dyn FnMut(&str),
358) -> Result<(), EngineError> {
359 let mut collected = String::new();
360 let mut on_line = |text: Line| {
361 let text = match text {
362 Line::Out(text) | Line::Err(text) => qframe::text::printable(&text).into_owned(),
363 };
364 collected.push_str(&text);
365 collected.push('\n');
366 line(&text);
367 };
368 let outcome = Process::new(&command.program).args(&command.args).no_stdin().run(cancel, &mut on_line);
369 let code = match outcome {
370 Ok(ProcessOutcome::Finished { code }) => code,
371 Ok(ProcessOutcome::Cancelled) => return Err(EngineError::Cancelled { command: command.clone() }),
372 Err(error) => return Err(EngineError::NotRunnable { command: command.clone(), error }),
373 };
374 if code == Some(0) {
375 return Ok(());
376 }
377 Err(EngineError::Failed(Failure { command: command.clone(), code, output: collected.trim_end().to_owned() }))
378}
379
380pub fn build_image(
392 engine: &Engine,
393 request: &ImageBuild<'_>,
394 cancel: &dyn Fn() -> bool,
395 line: &mut dyn FnMut(&str),
396) -> Result<(), EngineError> {
397 let result = stream(&engine.build_image(request), cancel, line);
398 if result.is_err() {
399 let _ = capture(&engine.remove_image(request.image));
400 }
401 result
402}
403
404#[cfg(test)]
405mod tests {
406
407 #[cfg(unix)]
408 #[test]
409 fn a_program_still_open_for_writing_a_moment_is_started_once_it_is_closed() {
410 use std::io::Write as _;
411 use std::os::unix::fs::PermissionsExt as _;
412 let stamp = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_nanos();
413 let dir = std::env::temp_dir().join(format!("qcode-busy-{}-{stamp}", std::process::id()));
414 std::fs::create_dir_all(&dir).expect("a folder");
415 let program = dir.join("engine");
416 let mut open = std::fs::File::create(&program).expect("the program is written");
417 open.write_all(b"#!/bin/sh\necho started\n").expect("written");
418 open.flush().expect("flushed");
419 std::fs::set_permissions(&program, std::fs::Permissions::from_mode(0o755)).expect("runnable");
420 let closer = std::thread::spawn(move || {
422 std::thread::sleep(Duration::from_millis(200));
423 drop(open);
424 });
425 let command =
426 EngineCommand { program: program.clone(), args: Vec::new(), deadline: Some(Duration::from_secs(30)) };
427 let said = capture(&command);
428 closer.join().expect("closed");
429 let _ = std::fs::remove_dir_all(&dir);
430 assert_eq!(said.expect("started once the file was closed").trim(), "started");
431 }
432 use super::{EngineError, capture, feed, pipe, stream};
433 use crate::engine::EngineCommand;
434 use std::cell::Cell;
435 use std::ffi::OsString;
436 use std::path::PathBuf;
437 use std::time::Duration;
438
439 #[cfg(unix)]
442 fn shell(script: &str) -> EngineCommand {
443 EngineCommand {
444 program: PathBuf::from("/bin/sh"),
445 args: vec![OsString::from("-c"), OsString::from(script)],
446 deadline: Some(crate::engine::ANSWER_WITHIN),
447 }
448 }
449
450 #[cfg(unix)]
453 fn stuck_engine(name: &str, script: &str, deadline: Duration) -> (EngineCommand, crate::engine::scratch::Scratch) {
454 use std::os::unix::fs::PermissionsExt as _;
455 let folder = crate::engine::scratch::Scratch::new(&format!("stuck-{name}")).expect("a folder");
456 let bin = folder.path().join("podman");
457 std::fs::write(&bin, format!("#!/bin/sh\n{script}\n")).expect("the stand-in is written");
458 std::fs::set_permissions(&bin, std::fs::Permissions::from_mode(0o755)).expect("runnable");
459 (EngineCommand { program: bin, args: Vec::new(), deadline: Some(deadline) }, folder)
460 }
461
462 fn within_bound<T: Send + 'static>(run: impl FnOnce() -> T + Send + 'static) -> T {
466 let (sender, answer) = std::sync::mpsc::channel();
467 std::thread::spawn(move || {
468 let _ = sender.send(run());
469 });
470 answer.recv_timeout(Duration::from_secs(30)).expect("the command was never given up on")
471 }
472
473 #[cfg(target_os = "linux")]
476 fn ends(pid: &str) -> bool {
477 let stat = format!("/proc/{}/stat", pid.trim());
478 let started = std::time::Instant::now();
479 while started.elapsed() < Duration::from_secs(20) {
480 match std::fs::read_to_string(&stat) {
481 Ok(stat) if !stat[stat.rfind(')').expect("name") + 2..].starts_with('Z') => {
482 std::thread::sleep(Duration::from_millis(20));
483 }
484 _ => return true,
485 }
486 }
487 false
488 }
489
490 #[cfg(target_os = "linux")]
491 #[test]
492 fn an_engine_that_does_not_answer_is_stopped_at_its_deadline_and_said_to_be_stuck() {
493 let (command, folder) =
494 stuck_engine("capture", "echo $$ > \"$(dirname \"$0\")/pid\"; exec sleep 600", Duration::from_millis(500));
495 let started = std::time::Instant::now();
496 let asked = command.clone();
497 let error = within_bound(move || capture(&asked)).expect_err("it never answers");
498 assert!(started.elapsed() < Duration::from_secs(5), "{:?}", started.elapsed());
499 let EngineError::TimedOut { command: stopped, after } = error else { panic!("{error:?}") };
500 assert_eq!((stopped, after), (command, Duration::from_millis(500)));
501 let pid = std::fs::read_to_string(folder.path().join("pid")).expect("it started");
502 assert!(ends(&pid), "the stuck engine was ended");
503 }
504
505 #[cfg(target_os = "linux")]
506 #[test]
507 fn an_engine_that_will_not_stop_when_asked_is_ended_after_the_grace() {
508 let script = "trap '' TERM; echo $$ > \"$(dirname \"$0\")/pid\"; while :; do sleep 1; done";
509 let (command, folder) = stuck_engine("stubborn", script, Duration::from_millis(300));
510 let error = within_bound(move || capture(&command)).expect_err("it never answers");
511 assert!(matches!(error, EngineError::TimedOut { .. }), "{error:?}");
512 let pid = std::fs::read_to_string(folder.path().join("pid")).expect("it started");
513 assert!(ends(&pid), "the engine that ignored the request was ended");
514 }
515
516 #[cfg(unix)]
517 #[test]
518 fn an_engine_fed_input_it_never_reads_is_given_up_on_too() {
519 let (command, _folder) = stuck_engine("feed", "exec sleep 600", Duration::from_millis(500));
520 let input = vec![b'x'; 300_000];
522 let error = within_bound(move || feed(&command, &input)).expect_err("it never answers");
523 assert!(matches!(error, EngineError::TimedOut { .. }), "{error:?}");
524 }
525
526 #[cfg(unix)]
527 #[test]
528 fn work_without_a_deadline_is_waited_for_however_long_it_takes() {
529 let (mut command, _folder) = stuck_engine("long", "sleep 1; echo done", Duration::from_millis(100));
530 command.deadline = None;
531 assert_eq!(within_bound(move || capture(&command)).expect("it finishes"), "done\n");
532 }
533
534 #[test]
535 fn a_stuck_engine_is_named_with_how_long_it_was_given() {
536 let env = qframe::env::Env::load(&qframe::env::AssetDirs {
537 locale_sources: crate::locales(),
538 ..qframe::env::AssetDirs::default()
539 })
540 .expect("the built-in files load");
541 let mut i18n = env.i18n().clone();
542 let command = crate::engine::Engine::new(crate::engine::EngineKind::Docker, "/usr/bin/docker").list_volumes();
543 for (language, said) in [
544 (
545 "en",
546 "Docker did not answer in 60 seconds, so QCode stopped waiting. It may be stuck on a lock of its own.",
547 ),
548 (
549 "tr",
550 "Docker 60 saniye içinde cevap vermedi, QCode da beklemeyi bıraktı. Kendi kilitlerinden birinde takılı kalmış olabilir.",
551 ),
552 ] {
553 assert!(i18n.set_active(language), "{language}");
554 let sentence = qframe::i18n::scope(std::sync::Arc::new(i18n.clone()), || {
555 super::timed_out(&command, Duration::from_secs(60))
556 });
557 assert_eq!(sentence, said);
558 }
559 }
560
561 #[cfg(unix)]
562 #[test]
563 fn captures_what_a_command_prints() {
564 assert_eq!(capture(&shell("echo qcode")).expect("the shell runs"), "qcode\n");
565 }
566
567 #[cfg(unix)]
568 #[test]
569 fn a_failure_carries_the_code_and_every_word_of_the_output() {
570 let error = capture(&shell("echo out; echo err >&2; exit 3")).expect_err("the script fails");
571 let EngineError::Failed(failure) = error else { panic!("the command ran and refused") };
572 assert_eq!(failure.code, Some(3));
573 assert_eq!(failure.output, "out\nerr");
574 }
575
576 #[cfg(unix)]
577 #[test]
578 fn what_is_fed_reaches_the_command_whole_even_past_the_pipe_size() {
579 let input = "x".repeat(300_000);
581 assert_eq!(feed(&shell("wc -c"), input.as_bytes()).expect("the shell runs").trim(), "300000");
582 let error = feed(&shell("cat > /dev/null; exit 4"), b"text").expect_err("the script fails");
583 let EngineError::Failed(failure) = error else { panic!("the command ran and refused") };
584 assert_eq!(failure.code, Some(4));
585 }
586
587 #[test]
588 fn a_missing_binary_is_not_the_same_as_a_failing_one() {
589 let command =
590 EngineCommand { program: PathBuf::from("/qcode/no/such/engine"), args: Vec::new(), deadline: None };
591 let error = capture(&command).expect_err("there is no such binary");
592 assert!(matches!(error, EngineError::NotRunnable { .. }), "{error:?}");
593 }
594
595 #[cfg(unix)]
596 #[test]
597 fn streams_both_streams_line_by_line() {
598 let mut lines = Vec::new();
599 stream(&shell("echo one; echo two >&2; echo three"), &|| false, &mut |line| lines.push(line.to_owned()))
600 .expect("the shell runs");
601 lines.sort();
602 assert_eq!(lines, ["one", "three", "two"]);
603 }
604
605 #[cfg(unix)]
606 #[test]
607 fn a_long_command_stops_when_it_is_cancelled() {
608 let seen = Cell::new(0_usize);
609 let error = stream(&shell("echo started; sleep 30"), &|| seen.get() > 0, &mut |_| seen.set(seen.get() + 1))
610 .expect_err("it was cancelled");
611 assert!(matches!(error, EngineError::Cancelled { .. }), "{error:?}");
612 }
613
614 #[cfg(target_os = "linux")]
615 #[test]
616 fn a_streamed_command_reads_no_input_and_cancelling_ends_what_it_started() {
617 let seen = std::cell::RefCell::new(Vec::new());
619 let error = stream(
620 &shell("readlink /proc/$$/fd/0; sleep 60 & echo $!; wait"),
621 &|| seen.borrow().len() == 2,
622 &mut |line| seen.borrow_mut().push(line.to_owned()),
623 )
624 .expect_err("it was cancelled");
625 assert!(matches!(error, EngineError::Cancelled { .. }), "{error:?}");
626 let seen = seen.into_inner();
627 assert_eq!(seen[0], "/dev/null");
628 let stat = format!("/proc/{}/stat", seen[1]);
629 let started = std::time::Instant::now();
630 while std::fs::read_to_string(&stat)
631 .is_ok_and(|stat| !stat[stat.rfind(')').expect("name") + 2..].starts_with('Z'))
632 {
633 assert!(started.elapsed() < std::time::Duration::from_secs(20), "the step still runs");
634 std::thread::sleep(std::time::Duration::from_millis(20));
635 }
636 }
637
638 #[cfg(unix)]
639 #[test]
640 fn a_pipe_hands_every_byte_from_one_command_to_the_other() {
641 let scratch = std::env::temp_dir().join(format!("qcode-pipe-{}", std::process::id()));
642 let _ = std::fs::remove_file(&scratch);
643 pipe(
646 &shell("head -c 300000 /dev/urandom; head -c 1000 /dev/zero"),
647 &shell(&format!("cat > {}", scratch.display())),
648 )
649 .expect("both ran");
650 let got = std::fs::read(&scratch).expect("the copy was written");
651 let _ = std::fs::remove_file(&scratch);
652 assert!(got.len() == 301_000 && got.ends_with(&[0_u8; 1000]), "{}", got.len());
653 }
654
655 #[cfg(unix)]
656 #[test]
657 fn a_pipe_says_which_side_failed_and_in_its_own_words() {
658 let reading = pipe(&shell("echo gone >&2; exit 3"), &shell("cat > /dev/null")).expect_err("the reader fails");
659 let EngineError::Failed(failure) = reading else { panic!("{reading:?}") };
660 assert_eq!((failure.code, failure.output.as_str()), (Some(3), "gone"));
661 let writing =
662 pipe(&shell("echo data"), &shell("cat > /dev/null; echo full >&2; exit 4")).expect_err("the writer fails");
663 let EngineError::Failed(failure) = writing else { panic!("{writing:?}") };
664 assert_eq!((failure.code, failure.output.as_str()), (Some(4), "full"));
665 }
666
667 #[cfg(unix)]
668 #[test]
669 fn what_a_stream_hands_on_has_no_control_character_left() {
670 let mut lines = Vec::new();
671 stream(&shell("printf 'context 2kB\\r\\r\\n\\033[1mbold\\033[0m\\n'"), &|| false, &mut |line| {
672 lines.push(line.to_owned());
673 })
674 .expect("the shell runs");
675 assert_eq!(lines, ["context 2kB", "bold"]);
676 }
677}