lattice_compilation/service.rs
1//! CM.1: the compilation process-lifecycle service.
2//!
3//! Runs a shell command **pipe-captured** off the actor thread and
4//! streams its stdout+stderr into the `*compilation*` buffer via
5//! [`CompilationOutputPushed`] events. The editor actor runs a
6//! `current_thread` runtime, so the blocking read loop must never
7//! land on `tokio::spawn`: the process + coordinator run on
8//! `spawn_blocking`, and the two pipe readers run on dedicated OS
9//! threads (mirroring the terminal reader task). Nothing here
10//! touches the UI/actor thread — paramount goal #1.
11//!
12//! Both stdout and stderr pipes are parsed for error locations —
13//! compile diagnostics (rustc/cargo) go to stderr while test failure
14//! output (thread panics) goes to stdout. A shared
15//! `Arc<Mutex<Vec<ErrorEntry>>>` accumulator merges entries from both
16//! streams into a single error list. Single-byte encoding errors in a
17//! pipe are logged and skipped (the reader does not stop).
18//!
19//! Lifecycle: `:recompile` reuses the last cmdline and kills the
20//! prior child before relaunching; the buffer is cleared via the
21//! `Reset` chunk and re-streamed. On exit a one-shot `info!`
22//! fires and a summary line is appended.
23//!
24//! # Safety
25//!
26//! The two `unsafe` blocks below (`libc::setpgid` in the spawn path
27//! and `libc::kill` in the kill path) are Unix-only and guarded by
28//! `#[cfg(unix)]`. Both are the only way to atomically terminate an
29//! entire process group (shell + all pipeline grandchildren). Without
30//! process-group kill, pipe grandchildren outlive the parent and keep
31//! the pipe readers blocking indefinitely.
32
33#![allow(unsafe_code)]
34
35use std::io::BufRead;
36use std::path::PathBuf;
37use std::process::{Child, Command, Stdio};
38use std::sync::{Arc, Mutex};
39
40#[cfg(unix)]
41use std::os::unix::process::CommandExt;
42
43use lattice_mode::inbound::InboundBus;
44use lattice_protocol::error_list::ErrorEntry;
45use lattice_runtime::EventBus;
46
47use crate::events::{CompilationOutputPushed, OutputChunk};
48use crate::parser::ParserRegistry;
49
50/// Process-lifecycle surface. Registered in the `ServiceRegistry`
51/// at boot (as [`CompilationServiceHandle`]); the `AppEffect::CompileRun`
52/// host arm looks it up and calls [`CompilationService::run`] (after
53/// creating the `*compilation*` buffer host-side).
54pub trait CompilationService: Send + Sync + std::fmt::Debug {
55 /// Launch (or relaunch) a compilation.
56 ///
57 /// `cmdline`: `Some(cmd)` runs and records `cmd` as the last
58 /// command; `None` reuses the last command (`:recompile` /
59 /// `:make` with no argument). With no prior command, publishes
60 /// a `Reset` explaining there is nothing to recompile and
61 /// returns without spawning.
62 ///
63 /// `cwd`: working directory for the child process.
64 fn run(&self, cmdline: Option<String>, cwd: Option<PathBuf>);
65 /// Kill the currently running compilation child process, if any.
66 /// No-op when no child is running. The reader pipes will EOF on
67 /// the closed child, and the drain will publish a Finished chunk
68 /// with the termination summary.
69 fn kill(&self);
70}
71
72/// Per the `ServiceRegistry` Arc/TypeId convention: register and
73/// look up under this exact alias.
74pub type CompilationServiceHandle = Arc<dyn CompilationService>;
75
76/// Publish one `Append` chunk per line for real-time streaming.
77/// The drain in `mode.rs` coalesces all events available per tick
78/// into a single `apply_edit_batch` — batching is its concern, not
79/// the pipe reader's. Publishing line-by-line gives the user
80/// immediate feedback without measurable overhead (event bus push is
81/// one `ArcSwap` store).
82const READER_BATCH_LINES: usize = 1;
83
84/// Mutable run state shared between `run` (which resolves the
85/// cmdline + kills the prior child) and the coordinator task
86/// (which stores the live child for kill-on-recompile and reaps
87/// it on exit).
88#[derive(Default)]
89struct RunState {
90 last_cmdline: Option<String>,
91 /// PR.4: the directory the last run was launched in, so
92 /// `:recompile` repeats it.
93 ///
94 /// Captured here rather than re-derived because by the time
95 /// `:recompile` fires, the active buffer is `*compilation*` itself —
96 /// which has no path, so re-resolving would silently fall back to
97 /// the working directory and rebuild the wrong project. "Do that
98 /// again" has to mean where, not just what.
99 ///
100 /// Lives beside `last_cmdline` for the same reason it does: this
101 /// state's lifetime is the service's, not the editor's.
102 last_cwd: Option<PathBuf>,
103 child: Option<Child>,
104}
105
106/// Default [`CompilationService`]: `sh -c <cmd>`, pipe-captured,
107/// streamed over the event bus.
108pub struct DefaultCompilationService {
109 events: Arc<EventBus>,
110 runtime: tokio::runtime::Handle,
111 state: Arc<Mutex<RunState>>,
112 /// CM.3a: the off-thread → host-state seam for parsed error
113 /// entries. The stderr reader accumulates parsed entries and sends
114 /// the FULL accumulated list; the inbound handler (in `install`)
115 /// maps it to `AppEffect::SetErrorList`. `send` wakes the editor so
116 /// the list reaches the screen off-keystroke.
117 qf_bus: InboundBus<Vec<ErrorEntry>>,
118 /// CM.5: the interned `compilation.ansi.*` elements captured
119 /// colour is painted with, filled by the mode during activation
120 /// (see [`crate::CompilationAnsiSlot`] for why it is late-bound).
121 ///
122 /// An empty slot leaves stripping in place and skips the spans —
123 /// the right degradation, because escape sequences must never
124 /// reach the buffer whether or not anyone can colour them.
125 ansi: Option<crate::CompilationAnsiSlot>,
126 /// CM.6b: plugin-contributed parser factories, snapshotted once per
127 /// run. `None` in a stripped harness that registered no handle;
128 /// empty in the common case where no `error-parser` plugin is
129 /// loaded. Either way every native parser still runs.
130 parser_factories: Option<crate::CompilationParserFactoriesHandle>,
131}
132
133impl std::fmt::Debug for DefaultCompilationService {
134 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
135 f.debug_struct("DefaultCompilationService")
136 .finish_non_exhaustive()
137 }
138}
139
140impl DefaultCompilationService {
141 pub fn new(
142 events: Arc<EventBus>,
143 runtime: tokio::runtime::Handle,
144 qf_bus: InboundBus<Vec<ErrorEntry>>,
145 ) -> Self {
146 Self {
147 events,
148 runtime,
149 state: Arc::new(Mutex::new(RunState::default())),
150 qf_bus,
151 ansi: None,
152 parser_factories: None,
153 }
154 }
155
156 /// CM.5: read the ANSI palette from `slot` when a run starts.
157 ///
158 /// Separate from [`Self::new`] because a stripped test harness
159 /// stands up neither a theme registry nor the slot — and a service
160 /// without one is still correct, just monochrome.
161 pub fn with_ansi_slot(mut self, slot: crate::CompilationAnsiSlot) -> Self {
162 self.ansi = Some(slot);
163 self
164 }
165
166 /// CM.6b: read plugin-contributed parser factories from `handle`
167 /// when a run starts.
168 ///
169 /// Separate from [`Self::new`] for the same reason the ANSI slot is:
170 /// a stripped harness registers no handle, and a service without one
171 /// is still correct — it just runs the built-in parsers only.
172 pub fn with_parser_factories(
173 mut self,
174 handle: crate::CompilationParserFactoriesHandle,
175 ) -> Self {
176 self.parser_factories = Some(handle);
177 self
178 }
179
180 fn publish(&self, chunk: OutputChunk) {
181 self.events.publish_typed(CompilationOutputPushed { chunk });
182 }
183}
184
185impl CompilationService for DefaultCompilationService {
186 fn run(&self, cmdline: Option<String>, cwd: Option<PathBuf>) {
187 // Resolve the cmdline AND the directory, record both for
188 // `:recompile`, and kill any prior child — all under one lock so
189 // a rapid recompile can't race two live children.
190 let (cmd, cwd) = {
191 let mut st = match self.state.lock() {
192 Ok(guard) => guard,
193 Err(_) => {
194 tracing::warn!("compilation: run-state lock poisoned; skipping run");
195 return;
196 }
197 };
198 let resolved = match cmdline {
199 Some(s) if !s.trim().is_empty() => {
200 st.last_cmdline = Some(s.clone());
201 // A fresh `:compile` re-binds the directory too: the
202 // caller resolved it from the buffer the command
203 // fired in, which is the answer we want to repeat.
204 st.last_cwd = cwd.clone();
205 (s, cwd)
206 }
207 // `:recompile` / bare `:make`. The caller's `cwd` is
208 // discarded on purpose — it was resolved from whatever
209 // is active NOW, and after the first run that is the
210 // pathless `*compilation*` buffer.
211 _ => match st.last_cmdline.clone() {
212 Some(prev) => (prev, st.last_cwd.clone()),
213 None => {
214 drop(st);
215 self.publish(OutputChunk::Reset {
216 header: "no previous compilation command\n\n".to_string(),
217 });
218 return;
219 }
220 },
221 };
222 if let Some(mut prior) = st.child.take() {
223 let _ = prior.kill();
224 let _ = prior.wait();
225 }
226 resolved
227 };
228
229 // Clear + seed the buffer with the run header (the drain's
230 // `Reset` path). Published before the spawn so it always
231 // precedes the streamed `Append`s.
232 self.publish(OutputChunk::Reset {
233 header: format!("$ {cmd}\n\n"),
234 });
235
236 // CM.3a: a new run clears the stale error list. Send an
237 // empty vec through the inbound seam so the host's
238 // replace-semantics `set_error_list` drops the prior run's
239 // entries before fresh ones stream in.
240 let _ = self.qf_bus.send(Vec::new());
241
242 let events = self.events.clone();
243 let state = self.state.clone();
244 let qf_bus = self.qf_bus.clone();
245 // Resolve the palette once per run rather than per line. An
246 // unfilled slot means the mode had no theme registry to intern
247 // against; stripping still happens, colouring does not.
248 let ansi: Option<crate::ansi::AnsiPalette> =
249 self.ansi.as_ref().and_then(|slot| slot.get().copied());
250 // CM.6b: snapshot the plugin factories once per run, not per
251 // reader and certainly not per line. A plugin loaded mid-build
252 // therefore joins the NEXT build — which is the honest
253 // behaviour, since a parser that starts halfway through a
254 // stream has no pending state for what it missed.
255 let factories: Option<Arc<crate::CompilationParserFactories>> = self
256 .parser_factories
257 .as_ref()
258 .map(|h| h.load_full())
259 .filter(|set| !set.is_empty());
260 self.runtime.spawn_blocking(move || {
261 let mut command = Command::new("sh");
262 command
263 .arg("-c")
264 .arg(&cmd)
265 .stdout(Stdio::piped())
266 .stderr(Stdio::piped());
267 // Unix: put the shell (and all its pipeline children) in
268 // their own process group so kill() terminates the entire
269 // tree, not just the shell PID. Without this, pipe
270 // grandchildren outlive the parent and keep the pipe
271 // readers blocking indefinitely.
272 #[cfg(unix)]
273 unsafe {
274 command.pre_exec(|| {
275 libc::setpgid(0, 0);
276 Ok(())
277 });
278 }
279 if let Some(dir) = cwd {
280 command.current_dir(dir);
281 }
282
283 let mut child = match command.spawn() {
284 Ok(child) => child,
285 Err(e) => {
286 publish(
287 &events,
288 OutputChunk::Finished {
289 summary: format!("\nCompilation failed to launch — {e}\n"),
290 },
291 );
292 return;
293 }
294 };
295
296 // Take the pipes out before parking the child in shared
297 // state, so kill-on-recompile and this task's reap never
298 // contend over the pipe fds.
299 let stdout = child.stdout.take();
300 let stderr = child.stderr.take();
301 if let Ok(mut st) = state.lock() {
302 st.child = Some(child);
303 }
304
305 // CM.3a+. Shared
306 // `Arc<Mutex<Vec<ErrorEntry>>>` so both stdout and
307 // stderr readers contribute to the same growing error
308 // list. Each reader locks, extends, clones, and sends the
309 // full state through the inbound seam — the host's
310 // replace-semantics `set_error_list` then grows the
311 // visible list regardless of which pipe delivered the
312 // entry. Each reader has its own `ParserRegistry` (the
313 // multi-line cargo parser is safe on both streams since
314 // cargo only emits diagnostics on stderr and the
315 // non-matching lines are no-ops).
316 //
317 // A shared list means the qf_bus always carries the
318 // complete state: stdout entries can't overwrite stderr
319 // entries (or vice versa) when the readers race.
320 let shared: Arc<Mutex<Vec<ErrorEntry>>> = Arc::new(Mutex::new(Vec::new()));
321
322 // Two dedicated reader threads so a large pipe can't
323 // deadlock the other. Both parse for error locations:
324 // compile diagnostics (rustc/cargo) go to stderr; test
325 // failure output (thread panics) goes to stdout.
326 let out_events = events.clone();
327 let out_qf = qf_bus.clone();
328 let out_shared = shared.clone();
329 let out_factories = factories.clone();
330 let out_reader = std::thread::spawn(move || {
331 read_parsed_pipe(
332 stdout,
333 &out_events,
334 &out_qf,
335 &out_shared,
336 ansi.as_ref(),
337 out_factories.as_deref(),
338 )
339 });
340 let err_events = events.clone();
341 let err_qf = qf_bus.clone();
342 let err_shared = shared.clone();
343 let err_factories = factories.clone();
344 let err_reader = std::thread::spawn(move || {
345 read_parsed_pipe(
346 stderr,
347 &err_events,
348 &err_qf,
349 &err_shared,
350 ansi.as_ref(),
351 err_factories.as_deref(),
352 )
353 });
354 let _ = out_reader.join();
355 let _ = err_reader.join();
356
357 // Reap the child — unless a concurrent recompile already
358 // took + killed it (then `child` is gone and the readers
359 // EOF'd on the closed pipes).
360 let waited = state.lock().ok().and_then(|mut st| st.child.take());
361 let summary = match waited {
362 Some(mut child) => match child.wait() {
363 Ok(status) => {
364 if status.success() {
365 tracing::info!("compilation finished");
366 format!("\nCompilation finished — {status}\n")
367 } else {
368 tracing::info!(%status, "compilation exited abnormally");
369 format!("\nCompilation exited abnormally — {status}\n")
370 }
371 }
372 Err(e) => format!("\nCompilation wait failed — {e}\n"),
373 },
374 None => "\nCompilation terminated\n".to_string(),
375 };
376 publish(&events, OutputChunk::Finished { summary });
377 });
378 }
379
380 fn kill(&self) {
381 if let Ok(mut st) = self.state.lock()
382 && let Some(mut child) = st.child.take()
383 {
384 // Unix: kill the entire process group, not just the
385 // shell PID. The shell was put in its own process
386 // group via pre_exec(setpgid(0,0)), so killpg()
387 // terminates the shell AND every pipeline grandchild.
388 // Without this, pipe grandchildren (seq, while, ...)
389 // survive the shell kill and keep stdout/stderr open.
390 #[cfg(unix)]
391 {
392 let pgid = child.id();
393 unsafe { libc::kill(-(pgid as i32), libc::SIGKILL) };
394 }
395 #[cfg(not(unix))]
396 {
397 let _ = child.kill();
398 }
399 let _ = child.wait();
400 }
401 }
402}
403
404/// Free helper so the coordinator closure (which owns `events` by
405/// move) can publish without borrowing `&self`.
406fn publish(events: &Arc<EventBus>, chunk: OutputChunk) {
407 events.publish_typed(CompilationOutputPushed { chunk });
408}
409
410/// CM.3a+. Blocking line-reader for one captured pipe.
411///
412/// Streams text into the `*compilation*` buffer (coalescing up to
413/// [`READER_BATCH_LINES`] lines per published `Append`) AND parses
414/// every line through a [`ParserRegistry`] for error locations. New
415/// [`ErrorEntry`]s are merged into the shared
416/// `Arc<Mutex<Vec<ErrorEntry>>>` accumulator and the full accumulated
417/// list is sent through the error inbound seam (replace-semantics
418/// `set_error_list`).
419///
420/// A per-line encoding error is logged at `debug!` and skipped — the
421/// reader does NOT stop on a single bad byte (prior behaviour lost the
422/// rest of the pipe). Flushes a partial batch on EOF.
423///
424/// CM.5: every line is passed through [`crate::ansi::clean_line`]
425/// **before** anything else sees it, so the escape sequences are gone
426/// from both the text that reaches the buffer and the text the
427/// parsers match against. That ordering is the point: a coloured
428/// `error[E0308]` carries `ESC[1m` in front of `error`, and matching
429/// the raw line would silently miss it.
430///
431/// `sgr` carries the active attributes across lines within this pipe
432/// (a producer may open a colour on one line and close it on the
433/// next). It is per-pipe, never shared — stdout and stderr are
434/// independent streams.
435///
436/// CM.6b: `factories` mints this reader's **own** plugin parsers. Each
437/// reader gets fresh instances for exactly the reason the `sgr` state
438/// above is per-pipe: the two streams carry independent pending state,
439/// and a WASM-backed parser could not be shared regardless (it owns a
440/// `Store`). They register ahead of the catch-all — see
441/// [`ParserRegistry::register_before_catch_all`].
442fn read_parsed_pipe<R: std::io::Read>(
443 pipe: Option<R>,
444 events: &Arc<EventBus>,
445 qf_bus: &InboundBus<Vec<ErrorEntry>>,
446 shared: &Mutex<Vec<ErrorEntry>>,
447 ansi: Option<&crate::ansi::AnsiPalette>,
448 factories: Option<&crate::CompilationParserFactories>,
449) {
450 let Some(pipe) = pipe else {
451 return;
452 };
453 let reader = std::io::BufReader::new(pipe);
454 let mut batch = String::new();
455 let mut batch_spans: Vec<Vec<lattice_cells::StyledSpan>> = Vec::new();
456 let mut lines_in_batch = 0usize;
457 let mut registry = ParserRegistry::with_builtins();
458 if let Some(factories) = factories {
459 for parser in factories.create_all() {
460 registry.register_before_catch_all(parser);
461 }
462 }
463 let mut sgr = crate::ansi::SgrState::default();
464 for line in reader.lines() {
465 let raw = match line {
466 Ok(l) => l,
467 Err(e) => {
468 tracing::debug!(error = %e, "compilation: pipe read error; skipping line");
469 continue;
470 }
471 };
472 let clean = crate::ansi::clean_line(&raw, &mut sgr, ansi);
473 let new_entries = registry.feed(&clean.text);
474 if !new_entries.is_empty() {
475 let mut guard = match shared.lock() {
476 Ok(g) => g,
477 Err(_) => return,
478 };
479 guard.extend(new_entries);
480 let _ = qf_bus.send(guard.clone());
481 }
482 batch.push_str(&clean.text);
483 batch.push('\n');
484 batch_spans.push(clean.spans);
485 lines_in_batch += 1;
486 if lines_in_batch >= READER_BATCH_LINES {
487 publish(
488 events,
489 OutputChunk::Append {
490 text: std::mem::take(&mut batch),
491 spans: std::mem::take(&mut batch_spans),
492 },
493 );
494 lines_in_batch = 0;
495 }
496 }
497 if !batch.is_empty() {
498 publish(
499 events,
500 OutputChunk::Append {
501 text: batch,
502 spans: batch_spans,
503 },
504 );
505 }
506}
507
508#[cfg(test)]
509mod tests {
510 #![allow(clippy::unwrap_used, clippy::panic)]
511 use super::*;
512 use lattice_protocol::error_list::ErrorSeverity;
513 use std::sync::Mutex;
514 use std::time::Duration;
515
516 /// Build a throwaway error inbound bus. The handler stashes the
517 /// LATEST full accumulated list into `latest` (the reader sends the
518 /// full list each time, so the last send is the complete set); the
519 /// returned drain must be run to flush queued sends through it.
520 fn qf_capture() -> (
521 InboundBus<Vec<ErrorEntry>>,
522 lattice_mode::tick_callback::TickCallback,
523 Arc<Mutex<Vec<ErrorEntry>>>,
524 ) {
525 let latest = Arc::new(Mutex::new(Vec::<ErrorEntry>::new()));
526 let latest_in = latest.clone();
527 let wake = Arc::new(tokio::sync::Notify::new());
528 let (bus, drain) = lattice_mode::inbound::make_inbound::<Vec<ErrorEntry>, _>(
529 wake,
530 move |entries: Vec<ErrorEntry>| {
531 *latest_in.lock().unwrap() = entries;
532 Vec::new()
533 },
534 );
535 (bus, drain, latest)
536 }
537
538 /// Collect chunks published to the bus for `dur`, running the
539 /// service on the ambient multi-thread test runtime.
540 fn collect_run(cmdline: Option<String>, dur: Duration) -> Vec<OutputChunk> {
541 collect_run_with_qf(cmdline, dur).0
542 }
543
544 /// Like [`collect_run`] but also returns the final parsed error
545 /// list captured off the inbound bus.
546 fn collect_run_with_qf(
547 cmdline: Option<String>,
548 dur: Duration,
549 ) -> (Vec<OutputChunk>, Vec<ErrorEntry>) {
550 collect_run_with_factories(cmdline, dur, None)
551 }
552
553 /// Like [`collect_run_with_qf`] but with CM.6b plugin parser
554 /// factories registered on the service.
555 fn collect_run_with_factories(
556 cmdline: Option<String>,
557 dur: Duration,
558 factories: Option<crate::CompilationParserFactoriesHandle>,
559 ) -> (Vec<OutputChunk>, Vec<ErrorEntry>) {
560 let bus = Arc::new(EventBus::new());
561 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<CompilationOutputPushed>();
562 bus.subscribe_typed::<CompilationOutputPushed>(tx);
563
564 let (qf_bus, mut qf_drain, latest) = qf_capture();
565 let mut svc =
566 DefaultCompilationService::new(bus.clone(), tokio::runtime::Handle::current(), qf_bus);
567 if let Some(handle) = factories {
568 svc = svc.with_parser_factories(handle);
569 }
570 svc.run(cmdline, None);
571
572 // Drain until quiescent (no new chunk within a short window)
573 // or the overall deadline elapses.
574 let mut chunks = Vec::new();
575 let deadline = std::time::Instant::now() + dur;
576 loop {
577 match rx.try_recv() {
578 Ok(ev) => chunks.push(ev.chunk),
579 Err(_) => {
580 if std::time::Instant::now() >= deadline {
581 break;
582 }
583 std::thread::sleep(Duration::from_millis(10));
584 }
585 }
586 }
587 // Flush any queued error sends through the capture handler.
588 let _ = qf_drain();
589 let entries = latest.lock().unwrap().clone();
590 (chunks, entries)
591 }
592
593 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
594 async fn run_streams_echo_output() {
595 let chunks = tokio::task::spawn_blocking(|| {
596 collect_run(Some("echo hello".to_string()), Duration::from_secs(3))
597 })
598 .await
599 .unwrap();
600
601 assert!(
602 matches!(chunks.first(), Some(OutputChunk::Reset { .. })),
603 "first chunk should be the run-header Reset, got {chunks:?}"
604 );
605 let joined: String = chunks
606 .iter()
607 .map(|c| match c {
608 OutputChunk::Reset { header } => header.clone(),
609 OutputChunk::Append { text, .. } => text.clone(),
610 OutputChunk::Finished { summary } => summary.clone(),
611 })
612 .collect();
613 assert!(
614 joined.contains("hello"),
615 "output should contain 'hello': {joined:?}"
616 );
617 assert!(
618 chunks
619 .iter()
620 .any(|c| matches!(c, OutputChunk::Finished { .. })),
621 "a Finished chunk should arrive, got {chunks:?}"
622 );
623 }
624
625 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
626 async fn stderr_diagnostics_populate_the_error_list() {
627 // A command whose stderr carries a gnu-style diagnostic must
628 // flow through the parser → inbound seam and land as a parsed
629 // error entry (0-based line/col).
630 let (_chunks, entries) = tokio::task::spawn_blocking(|| {
631 collect_run_with_qf(
632 Some("printf 'main.c:10:5: error: bad thing\\n' 1>&2".to_string()),
633 Duration::from_secs(3),
634 )
635 })
636 .await
637 .unwrap();
638
639 assert_eq!(
640 entries.len(),
641 1,
642 "expected one parsed entry, got {entries:?}"
643 );
644 let e = &entries[0];
645 assert_eq!(e.path, PathBuf::from("main.c"));
646 assert_eq!(e.line, 9, "1-based 10 → 0-based 9");
647 assert_eq!(e.col, 4, "1-based 5 → 0-based 4");
648 assert_eq!(e.severity, ErrorSeverity::Error);
649 assert_eq!(e.message, "bad thing");
650 }
651
652 /// CM.5, and the reason CM.5 is a correctness fix rather than a
653 /// cosmetic one: a colourised diagnostic carries `ESC[…m` in front
654 /// of `error`, so a parser fed the raw line matches nothing and the
655 /// entry silently never reaches the error list. This is the same
656 /// diagnostic as `stderr_diagnostics_populate_the_error_list`,
657 /// wearing the escapes a `--color=always` build would put on it.
658 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
659 async fn colourised_diagnostics_still_populate_the_error_list() {
660 let (chunks, entries) = tokio::task::spawn_blocking(|| {
661 collect_run_with_qf(
662 Some(
663 "printf '\\033[1m\\033[31mmain.c:10:5: error:\\033[0m bad thing\\n' 1>&2"
664 .to_string(),
665 ),
666 Duration::from_secs(3),
667 )
668 })
669 .await
670 .unwrap();
671
672 assert_eq!(
673 entries.len(),
674 1,
675 "expected one parsed entry from colourised output, got {entries:?}"
676 );
677 let e = &entries[0];
678 assert_eq!(e.path, PathBuf::from("main.c"));
679 assert_eq!(e.line, 9);
680 assert_eq!(e.col, 4);
681 assert_eq!(e.severity, ErrorSeverity::Error);
682
683 // And the text that reached the buffer carries no escapes.
684 let appended: String = chunks
685 .iter()
686 .filter_map(|c| match c {
687 OutputChunk::Append { text, .. } => Some(text.clone()),
688 _ => None,
689 })
690 .collect();
691 assert!(
692 !appended.contains('\u{1b}'),
693 "escape sequences must not reach the buffer, got {appended:?}"
694 );
695 assert!(appended.contains("main.c:10:5: error: bad thing"));
696 }
697
698 /// The palette slot is unfilled in this harness (no theme
699 /// registry), so colouring is off — but stripping is not
700 /// conditional on it.
701 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
702 async fn stripping_happens_without_a_palette() {
703 let chunks = tokio::task::spawn_blocking(|| {
704 collect_run(
705 Some("printf '\\033[32mgreen\\033[0m\\n'".to_string()),
706 Duration::from_secs(3),
707 )
708 })
709 .await
710 .unwrap();
711
712 let appended: String = chunks
713 .iter()
714 .filter_map(|c| match c {
715 OutputChunk::Append { text, spans } => {
716 assert!(
717 spans.iter().all(|l| l.is_empty()),
718 "no palette was interned, so no spans should be produced"
719 );
720 Some(text.clone())
721 }
722 _ => None,
723 })
724 .collect();
725 assert!(appended.contains("green"));
726 assert!(!appended.contains('\u{1b}'));
727 }
728
729 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
730 async fn recompile_with_no_prior_command_is_graceful() {
731 // `None` cmdline with no recorded last command must not
732 // panic; it publishes a single explanatory Reset.
733 let chunks = tokio::task::spawn_blocking(|| collect_run(None, Duration::from_millis(300)))
734 .await
735 .unwrap();
736
737 assert_eq!(
738 chunks.len(),
739 1,
740 "expected exactly one Reset, got {chunks:?}"
741 );
742 match &chunks[0] {
743 OutputChunk::Reset { header } => {
744 assert!(header.contains("no previous compilation command"));
745 }
746 other => panic!("expected Reset, got {other:?}"),
747 }
748 }
749
750 /// CM.6b: a factory registered on the service reaches the error
751 /// list for a line no built-in parser understands.
752 #[derive(Debug)]
753 struct QqFactory {
754 created: Arc<std::sync::atomic::AtomicUsize>,
755 }
756
757 #[derive(Debug)]
758 struct QqParser;
759
760 impl crate::CompilationParser for QqParser {
761 fn feed(&mut self, line: &str) -> Vec<ErrorEntry> {
762 line.strip_prefix("QQ ")
763 .map(|rest| {
764 vec![ErrorEntry {
765 path: std::path::PathBuf::from(rest),
766 line: 7,
767 col: 3,
768 severity: ErrorSeverity::Warning,
769 message: "from a plugin".to_string(),
770 }]
771 })
772 .unwrap_or_default()
773 }
774 }
775
776 impl crate::CompilationParserFactory for QqFactory {
777 fn plugin_id(&self) -> u64 {
778 42
779 }
780 fn create(&self) -> Option<Box<dyn crate::CompilationParser>> {
781 self.created
782 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
783 Some(Box::new(QqParser))
784 }
785 }
786
787 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
788 async fn a_registered_factory_contributes_entries() {
789 let created = Arc::new(std::sync::atomic::AtomicUsize::new(0));
790 let mut set = crate::CompilationParserFactories::new();
791 set.register(Arc::new(QqFactory {
792 created: created.clone(),
793 }));
794 let handle: crate::CompilationParserFactoriesHandle =
795 Arc::new(arc_swap::ArcSwap::from_pointee(set));
796
797 let created_probe = created.clone();
798 let (_chunks, entries) = tokio::task::spawn_blocking(move || {
799 collect_run_with_factories(
800 Some("echo 'QQ src/plugin.rs'".to_string()),
801 Duration::from_secs(3),
802 Some(handle),
803 )
804 })
805 .await
806 .unwrap();
807
808 assert!(
809 entries
810 .iter()
811 .any(|e| e.path == std::path::PathBuf::from("src/plugin.rs")
812 && e.severity == ErrorSeverity::Warning
813 && e.message == "from a plugin"),
814 "the plugin parser's entry should reach the error list: {entries:?}"
815 );
816 // One instance per reader — the property the factory exists for.
817 assert_eq!(
818 created_probe.load(std::sync::atomic::Ordering::SeqCst),
819 2,
820 "stdout and stderr each mint their own parser"
821 );
822 }
823
824 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
825 async fn no_registered_factory_leaves_the_builtins_alone() {
826 // The common case: no `error-parser` plugin loaded. The
827 // built-in parsers must behave exactly as before — a
828 // gnu-style line still lands.
829 let (_chunks, entries) = tokio::task::spawn_blocking(|| {
830 collect_run_with_factories(
831 Some("echo 'src/a.rs:3:5: error: boom'".to_string()),
832 Duration::from_secs(3),
833 None,
834 )
835 })
836 .await
837 .unwrap();
838
839 assert!(
840 entries
841 .iter()
842 .any(|e| e.path == std::path::PathBuf::from("src/a.rs")),
843 "built-in parsing is unaffected by the factory seam: {entries:?}"
844 );
845 }
846
847 /// PR.4: `:recompile` repeats WHERE, not just what.
848 ///
849 /// By the time it fires, the active buffer is `*compilation*`,
850 /// which has no path — so the host resolves a project root from it
851 /// and gets the working directory. If the service took that value,
852 /// a recompile would rebuild the wrong project while looking like
853 /// it worked.
854 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
855 async fn recompile_reuses_the_directory_of_the_last_real_run() {
856 let dir = std::env::temp_dir().join(format!(
857 "lattice-recompile-cwd-{}",
858 std::time::SystemTime::now()
859 .duration_since(std::time::UNIX_EPOCH)
860 .map(|d| d.as_nanos())
861 .unwrap_or(0)
862 ));
863 std::fs::create_dir_all(&dir).unwrap();
864 let canonical = std::fs::canonicalize(&dir).unwrap();
865 // What `pwd` is matched on: the directory's unique leaf, not the whole
866 // path. The shell and `canonicalize` spell a path differently on
867 // Windows (`/c/Users/…` against `\\?\C:\Users\…`), and the leaf is
868 // the part both agree on — and the part the second run's directory,
869 // its parent, does not contain.
870 let leaf = canonical
871 .file_name()
872 .and_then(|n| n.to_str())
873 .unwrap()
874 .to_string();
875
876 let bus = Arc::new(EventBus::new());
877 let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<CompilationOutputPushed>();
878 bus.subscribe_typed::<CompilationOutputPushed>(tx);
879 let (qf_bus, _drain, _latest) = qf_capture();
880 let svc =
881 DefaultCompilationService::new(bus.clone(), tokio::runtime::Handle::current(), qf_bus);
882
883 let collect = |rx: &mut tokio::sync::mpsc::UnboundedReceiver<CompilationOutputPushed>| {
884 let deadline = std::time::Instant::now() + Duration::from_secs(3);
885 let mut text = String::new();
886 let mut saw_finish = false;
887 while std::time::Instant::now() < deadline && !saw_finish {
888 match rx.try_recv() {
889 Ok(ev) => match ev.chunk {
890 OutputChunk::Append { text: t, .. } => text.push_str(&t),
891 OutputChunk::Finished { .. } => saw_finish = true,
892 OutputChunk::Reset { .. } => {}
893 },
894 Err(_) => std::thread::sleep(Duration::from_millis(10)),
895 }
896 }
897 text
898 };
899
900 // A real run, in `dir`.
901 svc.run(Some("pwd".to_string()), Some(canonical.clone()));
902 let first = tokio::task::spawn_blocking({
903 let mut rx = rx;
904 move || {
905 let t = collect(&mut rx);
906 (t, rx)
907 }
908 })
909 .await
910 .unwrap();
911 let (first_text, mut rx) = first;
912 assert!(
913 first_text.contains(&leaf),
914 "the first run should be in {canonical:?}, got {first_text:?}"
915 );
916
917 // `:recompile` — no cmdline, and a DIFFERENT cwd, standing in for
918 // the pathless `*compilation*` buffer resolving to somewhere else.
919 svc.run(None, Some(std::env::temp_dir()));
920 let second = tokio::task::spawn_blocking(move || collect(&mut rx))
921 .await
922 .unwrap();
923 assert!(
924 second.contains(&leaf),
925 "recompile must reuse the first run's directory, got {second:?}"
926 );
927
928 let _ = std::fs::remove_dir_all(&dir);
929 }
930
931 #[test]
932 fn service_is_debug() {
933 // Trait bound sanity: the handle stays object-safe + Debug.
934 fn assert_debug<T: std::fmt::Debug + Send + Sync>() {}
935 assert_debug::<DefaultCompilationService>();
936 }
937}