Skip to main content

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}