Skip to main content

pitchfork_cli/cli/
log_sink.rs

1use crate::Result;
2use crate::daemon_id::DaemonId;
3use crate::log_parse::ParsedLog;
4use crate::log_store::LogStore;
5use crate::log_store::sqlite::LOG_STORE;
6use tokio::io::AsyncReadExt;
7
8/// Number of parsed lines to accumulate before writing them as one batch.
9const BATCH_SIZE: usize = 100;
10
11/// Longest a parsed line waits in the batch before being written.
12const FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
13
14/// How many parsed lines may be queued for writing before reading slows down.
15///
16/// Reading and writing run separately so a slow write cannot stop the pipe being
17/// drained, but the queue is bounded: output that sustainably outpaces the store
18/// has to push back on the daemon eventually, which is preferable to growing
19/// without limit.
20const QUEUE_DEPTH: usize = 8192;
21
22/// Bytes read from the pipe at a time.
23const READ_CHUNK: usize = 8192;
24
25/// Longest run of bytes treated as a single line.
26///
27/// Output containing no newline must not accumulate indefinitely: a daemon
28/// emitting an endless stream, or binary data, would otherwise grow the buffer
29/// until the sink was killed for using too much memory — whereupon the
30/// supervisor would start another sink and repeat it.
31const MAX_LINE_BYTES: usize = 64 * 1024;
32
33/// Reads a daemon's output on stdin and writes it to the log store
34///
35/// Spawned by the supervisor as a sibling of the daemon, holding the read end
36/// of the daemon's output pipe. Keeping the reader in its own process is what
37/// makes logging survive a supervisor crash: the pipe still has a reader, so
38/// the daemon is neither killed by SIGPIPE nor blocked, and no output is lost.
39/// Exits when the pipe reaches end of file, which happens once the daemon and
40/// every descendant holding the write end have gone.
41#[derive(Debug, usage_rs::Args)]
42#[usage(verbatim_doc_comment)]
43pub struct LogSink {
44    /// Qualified id of the daemon whose output this is
45    #[usage(long)]
46    daemon_id: String,
47
48    /// Log format to parse lines with (`json`, `logfmt`, `auto`, or `text`)
49    #[usage(long, default = "text")]
50    log_format: String,
51
52    /// Regex whose first match means the daemon is ready
53    ///
54    /// Set for a daemon configured with `ready_output`. The supervisor cannot
55    /// match it itself — this process holds the output — so the match is
56    /// reported back over IPC.
57    #[usage(long)]
58    ready_pattern: Option<String>,
59
60    /// Token identifying the start attempt this sink belongs to
61    ///
62    /// Quoted back when reporting a match. The supervisor drops reports whose
63    /// token is no longer current, so a sink still draining a failed attempt
64    /// cannot mark that daemon's retry ready.
65    #[usage(long, default_value_t = 0, default = "0")]
66    relay_token: u64,
67
68    /// Report lines so the supervisor can fire the daemon's `on_output` hook
69    ///
70    /// Without `--output-filter` or `--output-regex` every line qualifies,
71    /// which is what a hook with no pattern asks for.
72    #[usage(long)]
73    report_output: bool,
74
75    /// Only report lines containing this substring
76    #[usage(long)]
77    output_filter: Option<String>,
78
79    /// Only report lines matching this regex
80    #[usage(long)]
81    output_regex: Option<String>,
82
83    /// Shortest gap between reported lines, in milliseconds
84    #[usage(long, default_value_t = 1000, default = "1000")]
85    output_debounce_ms: u64,
86}
87
88impl LogSink {
89    pub async fn run(&self) -> Result<()> {
90        let id = DaemonId::parse(&self.daemon_id)?;
91
92        // A pattern that does not compile is reported and then ignored, rather
93        // than failing the sink: refusing to start would leave the daemon's
94        // output unread, which is far worse than a readiness check that never
95        // fires. The supervisor validates patterns too, so this is a backstop.
96        let compile = |what: &str, pattern: &str| {
97            regex::Regex::new(pattern)
98                .map_err(|e| error!("log sink for {id} ignoring unparsable {what}: {e}"))
99                .ok()
100        };
101        let ready_pattern = self
102            .ready_pattern
103            .as_deref()
104            .and_then(|p| compile("ready pattern", p));
105        let hook = self.report_output.then(|| HookMatcher {
106            filter: self.output_filter.clone(),
107            regex: self
108                .output_regex
109                .as_deref()
110                .and_then(|p| compile("output pattern", p)),
111            debounce: std::time::Duration::from_millis(self.output_debounce_ms),
112            last_reported: None,
113        });
114
115        // Reading and writing are separate tasks. A write to SQLite can block —
116        // for as long as the store's busy timeout, if another writer holds the
117        // lock — and this process is the only reader of the daemon's pipe, so a
118        // write must never stop it being drained.
119        let (tx, rx) = tokio::sync::mpsc::channel::<SinkEvent>(QUEUE_DEPTH);
120        let writer = tokio::spawn(write_batches(id.clone(), self.relay_token, rx));
121
122        let read_result =
123            read_lines(tx, &self.log_format, ReadyMatcher::new(ready_pattern, hook)).await;
124
125        // The sender has been dropped by now, so the writer drains its queue and
126        // returns; wait for it so nothing queued is lost on exit.
127        let _ = writer.await;
128
129        read_result.map_err(|e| {
130            miette::miette!("log sink for {id} could not read the daemon's output: {e}")
131        })
132    }
133}
134
135/// Something for the writer task to do, in the order the reader saw it.
136///
137/// Reporting a readiness match travels the same queue as the lines rather than
138/// jumping ahead of them, so the line that triggered the match is always in the
139/// log store by the time the supervisor hears about it — `collect_startup_logs`
140/// and `pitchfork logs` would otherwise be able to miss it.
141enum SinkEvent {
142    Line(ParsedLog),
143    Report(Report),
144}
145
146/// A line the supervisor needs to see, and why.
147#[derive(Debug, PartialEq, Eq)]
148struct Report {
149    text: String,
150    /// Whether it passed the `on_output` hook's filter and debounce. False for
151    /// a line reported only because it matched the readiness pattern — firing a
152    /// hook that filters for something else would be wrong.
153    fires_hook: bool,
154}
155
156/// How much of a capped line is carried forward for matching.
157///
158/// A line longer than [`MAX_LINE_BYTES`] is emitted in pieces, and a readiness
159/// pattern straddling a split would match none of them — the daemon would then
160/// be killed at its readiness timeout despite having announced itself. Keeping
161/// the tail of the previous piece closes that for any pattern shorter than this
162/// while still bounding what is held.
163const MATCH_CARRY_BYTES: usize = 4 * 1024;
164
165/// Watches a daemon's output for the things the supervisor would look for if it
166/// could still read the stream: the readiness pattern, and whatever fires the
167/// `on_output` hook.
168///
169/// What the supervisor does about a reported line — mark the daemon ready, run
170/// the hook — remains its own business.
171struct ReadyMatcher {
172    /// Readiness pattern, cleared once it has matched: readiness happens once.
173    pattern: Option<regex::Regex>,
174    /// The `on_output` hook's filter and rate limit, if the daemon has one.
175    hook: Option<HookMatcher>,
176    /// Tail of the previous piece of a line split at the length cap. Empty
177    /// whenever the last piece ended at a real newline.
178    carried: String,
179}
180
181/// The `on_output` hook's line filter and its rate limit.
182struct HookMatcher {
183    filter: Option<String>,
184    regex: Option<regex::Regex>,
185    debounce: std::time::Duration,
186    last_reported: Option<std::time::Instant>,
187}
188
189impl HookMatcher {
190    /// Whether this line should fire the hook, consuming the debounce window.
191    ///
192    /// The debounce is applied here rather than by the supervisor because this
193    /// process is the one that sees every line: enforcing it here means one IPC
194    /// message per window instead of one per line, which matters for a hook
195    /// with no filter, where every line qualifies.
196    fn matches(&mut self, clean: &str) -> bool {
197        let matched = match (&self.filter, &self.regex) {
198            // Mutually exclusive, and a hook setting both is rejected before it
199            // ever reaches this process.
200            (Some(substr), _) => clean.contains(substr.as_str()),
201            (None, Some(re)) => re.is_match(clean),
202            // A hook with neither fires on every line.
203            (None, None) => true,
204        };
205        if !matched {
206            return false;
207        }
208        let now = std::time::Instant::now();
209        if self
210            .last_reported
211            .is_some_and(|last| now.duration_since(last) < self.debounce)
212        {
213            return false;
214        }
215        self.last_reported = Some(now);
216        true
217    }
218}
219
220impl ReadyMatcher {
221    fn new(pattern: Option<regex::Regex>, hook: Option<HookMatcher>) -> Self {
222        Self {
223            pattern,
224            hook,
225            carried: String::new(),
226        }
227    }
228
229    /// Whether anything is still being watched for. Once nothing is, matching is
230    /// skipped: stripping ANSI from every line of a chatty daemon is not free.
231    fn is_watching(&self) -> bool {
232        self.pattern.is_some() || self.hook.is_some()
233    }
234
235    /// Test `text` — one whole line, or one piece of an over-long one — and
236    /// return what matched, which is what the supervisor is told about.
237    ///
238    /// `split_at_cap` says another piece of the same logical line follows.
239    fn consider(&mut self, text: &str, split_at_cap: bool) -> Option<Report> {
240        if !self.is_watching() {
241            self.carried.clear();
242            return None;
243        }
244        let clean = console::strip_ansi_codes(text);
245        let candidate = if self.carried.is_empty() {
246            clean.into_owned()
247        } else {
248            format!("{}{clean}", self.carried)
249        };
250
251        self.carried = if split_at_cap {
252            let start = candidate.len().saturating_sub(MATCH_CARRY_BYTES);
253            // Never split a character in half; walk forward to a boundary.
254            let start = (start..candidate.len())
255                .find(|i| candidate.is_char_boundary(*i))
256                .unwrap_or(candidate.len());
257            candidate[start..].to_string()
258        } else {
259            String::new()
260        };
261
262        // A line can qualify for either reason, or both. The supervisor
263        // re-checks the readiness pattern on what it is sent, but it cannot
264        // re-check the hook without redoing the debounce, so that answer is
265        // carried explicitly.
266        let mut ready_matched = false;
267        if self
268            .pattern
269            .as_ref()
270            .is_some_and(|re| re.is_match(&candidate))
271        {
272            self.pattern = None;
273            ready_matched = true;
274        }
275        let fires_hook = self
276            .hook
277            .as_mut()
278            .is_some_and(|hook| hook.matches(&candidate));
279
280        // Send the text the patterns actually matched against, not just this
281        // piece of it: the supervisor re-matches it, and hands it to the hook.
282        (ready_matched || fires_hook).then_some(Report {
283            text: candidate,
284            fires_hook,
285        })
286    }
287}
288
289/// Split the daemon's output into lines and queue them for writing.
290///
291/// Returns once the pipe reaches end of file. A read error is propagated so the
292/// process can exit non-zero: exiting cleanly would tell the supervisor the
293/// stream had finished and it would stop replacing this sink.
294async fn read_lines(
295    tx: tokio::sync::mpsc::Sender<SinkEvent>,
296    log_format: &str,
297    mut matcher: ReadyMatcher,
298) -> std::io::Result<()> {
299    let mut stdin = tokio::io::stdin();
300    let mut chunk = vec![0u8; READ_CHUNK];
301    let mut line: Vec<u8> = Vec::with_capacity(256);
302    // Whether the last line was emitted because it reached the cap rather than
303    // because it ended. A newline arriving straight afterwards terminates the
304    // line already written, so it must not produce an empty one.
305    let mut split_at_cap = false;
306
307    loop {
308        let read = stdin.read(&mut chunk).await?;
309        if read == 0 {
310            break;
311        }
312        for &byte in &chunk[..read] {
313            if byte == b'\n' {
314                if split_at_cap && line.is_empty() {
315                    split_at_cap = false;
316                    continue;
317                }
318                split_at_cap = false;
319                queue(&tx, &mut line, log_format, &mut matcher).await?;
320            } else {
321                line.push(byte);
322                split_at_cap = false;
323                // Emit an over-long run as its own line rather than letting the
324                // buffer grow without bound.
325                if line.len() >= MAX_LINE_BYTES {
326                    queue_capped(&tx, &mut line, log_format, &mut matcher).await?;
327                    split_at_cap = true;
328                }
329            }
330        }
331    }
332
333    // Anything written without a trailing newline is still output.
334    if !line.is_empty() {
335        queue(&tx, &mut line, log_format, &mut matcher).await?;
336    }
337    Ok(())
338}
339
340/// Emit a line that has reached the length cap, keeping any trailing bytes that
341/// form an incomplete character.
342///
343/// Splitting purely by byte count would cut a multi-byte character in half, and
344/// converting each half on its own turns one valid character into two
345/// replacement characters.
346async fn queue_capped(
347    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
348    line: &mut Vec<u8>,
349    log_format: &str,
350    matcher: &mut ReadyMatcher,
351) -> std::io::Result<()> {
352    let split = split_before_incomplete_char(line);
353    let tail = line.split_off(split);
354    let result = queue_piece(tx, line, log_format, matcher).await;
355    *line = tail;
356    result
357}
358
359/// Queue one piece of a line that hit the length cap, telling the matcher that
360/// the rest of the logical line is still to come.
361async fn queue_piece(
362    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
363    line: &mut Vec<u8>,
364    log_format: &str,
365    matcher: &mut ReadyMatcher,
366) -> std::io::Result<()> {
367    let text = decode_line(line);
368    let text = text.trim_end_matches('\r');
369    let parsed = crate::log_parse::parse(text, log_format);
370    let report = matcher.consider(text, true);
371    line.clear();
372
373    tx.send(SinkEvent::Line(parsed))
374        .await
375        .map_err(|_| std::io::Error::other("log writer stopped"))?;
376    if let Some(report) = report {
377        tx.send(SinkEvent::Report(report))
378            .await
379            .map_err(|_| std::io::Error::other("log writer stopped"))?;
380    }
381    Ok(())
382}
383
384/// Length to cut `bytes` at so no character is left half-written, in whichever
385/// encoding [`decode_line`] will read the piece in.
386fn split_before_incomplete_char(bytes: &[u8]) -> usize {
387    #[cfg(windows)]
388    return split_before_incomplete_char_in(bytes, console_code_page());
389    #[cfg(not(windows))]
390    split_before_incomplete_utf8(bytes)
391}
392
393/// [`split_before_incomplete_char`] for a console using `code_page`.
394#[cfg(windows)]
395fn split_before_incomplete_char_in(bytes: &[u8], code_page: u32) -> usize {
396    let split = split_before_incomplete_utf8(bytes);
397    if uses_code_page(&bytes[..split]) {
398        // The piece is read in the code page, so the cut must fall between
399        // its characters. Within that, also hold back an unfinished UTF-8
400        // character: a line read in the code page can still carry UTF-8 — a
401        // stray byte decides how it is read, not what it contains — and a
402        // byte carried over to the next piece is decoded there, so holding
403        // one back loses nothing. A UTF-8 cut inside a double-byte character
404        // the code page reads as whole is not taken.
405        let dbcs_split = split_before_incomplete_dbcs(bytes, code_page);
406        if split < dbcs_split && is_dbcs_boundary(bytes, split, code_page) {
407            return split;
408        }
409        return dbcs_split;
410    }
411    split
412}
413
414/// Length to cut `bytes` at so no UTF-8 character is left half-written.
415///
416/// Decided by inspecting the final bytes rather than by asking `from_utf8` where
417/// the string stops being valid: that reports the *first* problem, so a single
418/// invalid byte earlier in the line would hide an unfinished character at the
419/// end, and the character would be split after all.
420fn split_before_incomplete_utf8(bytes: &[u8]) -> usize {
421    let len = bytes.len();
422    // A character is at most four bytes, so only the last few can be unfinished.
423    for i in (len.saturating_sub(4)..len).rev() {
424        let byte = bytes[i];
425        if byte & 0b1100_0000 == 0b1000_0000 {
426            continue; // a continuation byte; keep looking back for its lead
427        }
428        let expected = match byte {
429            0x00..=0x7f => 1,
430            b if b >> 5 == 0b110 => 2,
431            b if b >> 4 == 0b1110 => 3,
432            b if b >> 3 == 0b11110 => 4,
433            // Not a valid lead byte at all, so nothing is pending; the lossy
434            // conversion will render it.
435            _ => 1,
436        };
437        return if i + expected > len && i > 0 { i } else { len };
438    }
439    len
440}
441
442/// Text of one line of daemon output.
443///
444/// UTF-8 is taken as it is. Anything else is not assumed to be damaged UTF-8:
445/// on Windows, console programs — `cmd` among them — write in the console's
446/// code page, so a Japanese system hands over Shift_JIS, and reading that as
447/// UTF-8 would store nothing but replacement characters.
448pub(crate) fn decode_line(bytes: &[u8]) -> std::borrow::Cow<'_, str> {
449    if let Ok(text) = std::str::from_utf8(bytes) {
450        return std::borrow::Cow::Borrowed(text);
451    }
452    #[cfg(windows)]
453    if uses_code_page(bytes)
454        && let Some(text) = decode_in_code_page(bytes, console_code_page())
455    {
456        return std::borrow::Cow::Owned(text);
457    }
458    String::from_utf8_lossy(bytes)
459}
460
461/// Whether [`decode_line`] reads `bytes` in the console code page.
462///
463/// Only when they are not UTF-8, and not mostly UTF-8 either: a UTF-8 line with
464/// a stray invalid byte — or with another stream's output spliced into it —
465/// would have all of its text garbled by decoding it in the code page, where
466/// the lossy conversion loses just the stray bytes.
467///
468/// Mostly means more characters read as multi-byte UTF-8 than runs of bytes
469/// that cannot. Code page text does sometimes form a valid UTF-8 character by
470/// chance, but around it are far more bytes that do not.
471#[cfg(windows)]
472fn uses_code_page(bytes: &[u8]) -> bool {
473    if std::str::from_utf8(bytes).is_ok() {
474        return false;
475    }
476    let (mut multi_byte, mut invalid) = (0usize, 0usize);
477    for chunk in bytes.utf8_chunks() {
478        multi_byte += chunk.valid().chars().filter(|c| !c.is_ascii()).count();
479        invalid += usize::from(!chunk.invalid().is_empty());
480    }
481    multi_byte <= invalid
482}
483
484/// Code page the daemon's console programs write in.
485///
486/// This process is started the same way as the daemon, so it shares the
487/// daemon's console or gets one set up alike. Without a console at all, the
488/// OEM code page is what a new console would have used.
489#[cfg(windows)]
490fn console_code_page() -> u32 {
491    match unsafe { windows_sys::Win32::System::Console::GetConsoleOutputCP() } {
492        0 => unsafe { windows_sys::Win32::Globalization::GetOEMCP() },
493        code_page => code_page,
494    }
495}
496
497/// Decode `bytes` from `code_page`, or `None` when Windows cannot.
498///
499/// A UTF-8 code page gets `None` too: the bytes are already known not to be
500/// valid UTF-8, and the lossy conversion handles them just as well.
501#[cfg(windows)]
502fn decode_in_code_page(bytes: &[u8], code_page: u32) -> Option<String> {
503    use windows_sys::Win32::Globalization::{CP_UTF8, MultiByteToWideChar};
504
505    if code_page == CP_UTF8 {
506        return None;
507    }
508    let len = i32::try_from(bytes.len()).ok()?;
509    // The first call measures, the second converts.
510    let wide_len =
511        unsafe { MultiByteToWideChar(code_page, 0, bytes.as_ptr(), len, std::ptr::null_mut(), 0) };
512    if wide_len <= 0 {
513        return None;
514    }
515    let mut wide = vec![0u16; wide_len as usize];
516    let written = unsafe {
517        MultiByteToWideChar(
518            code_page,
519            0,
520            bytes.as_ptr(),
521            len,
522            wide.as_mut_ptr(),
523            wide_len,
524        )
525    };
526    if written <= 0 {
527        return None;
528    }
529    wide.truncate(written as usize);
530    Some(String::from_utf16_lossy(&wide))
531}
532
533/// Length to cut `bytes` at so no double-byte character of `code_page` is left
534/// half-written.
535///
536/// A trail byte can have the same value as a lead byte, so looking at the last
537/// byte alone cannot tell which it is; only walking from the start of the line
538/// can. Single-byte code pages have no lead bytes, so they never cut short.
539#[cfg(windows)]
540fn split_before_incomplete_dbcs(bytes: &[u8], code_page: u32) -> usize {
541    use windows_sys::Win32::Globalization::IsDBCSLeadByteEx;
542
543    let mut i = 0;
544    while i < bytes.len() {
545        let lead = unsafe { IsDBCSLeadByteEx(code_page, bytes[i]) } != 0;
546        i += if lead { 2 } else { 1 };
547    }
548    // Stepping past the end means the last byte opened a character.
549    if i > bytes.len() && bytes.len() > 1 {
550        bytes.len() - 1
551    } else {
552        bytes.len()
553    }
554}
555
556/// Whether `at` falls between two characters of `bytes` read in `code_page`,
557/// rather than inside a double-byte one. Walks from the start, for the same
558/// reason as [`split_before_incomplete_dbcs`].
559#[cfg(windows)]
560fn is_dbcs_boundary(bytes: &[u8], at: usize, code_page: u32) -> bool {
561    use windows_sys::Win32::Globalization::IsDBCSLeadByteEx;
562
563    let mut i = 0;
564    while i < at {
565        let lead = unsafe { IsDBCSLeadByteEx(code_page, bytes[i]) } != 0;
566        i += if lead { 2 } else { 1 };
567    }
568    i == at
569}
570
571/// Parse `line` and hand it to the writer, clearing it either way.
572///
573/// A closed queue means the writer task is gone, which is a failure rather than
574/// the end of the stream: reporting it as success would tell the supervisor this
575/// sink had reached end of file, and it would stop replacing it while the daemon
576/// was still writing.
577async fn queue(
578    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
579    line: &mut Vec<u8>,
580    log_format: &str,
581    matcher: &mut ReadyMatcher,
582) -> std::io::Result<()> {
583    // Never fails: a daemon emitting a stray non-UTF-8 byte must not be able to
584    // stop its own logging.
585    let text = decode_line(line);
586    let text = text.trim_end_matches('\r');
587    let parsed = crate::log_parse::parse(text, log_format);
588    // Strip ANSI before matching so a pattern works whether or not the daemon
589    // colours its output, matching what in-process capture did.
590    let report = matcher.consider(text, false);
591    line.clear();
592
593    tx.send(SinkEvent::Line(parsed))
594        .await
595        .map_err(|_| std::io::Error::other("log writer stopped"))?;
596    if let Some(report) = report {
597        tx.send(SinkEvent::Report(report))
598            .await
599            .map_err(|_| std::io::Error::other("log writer stopped"))?;
600    }
601    Ok(())
602}
603
604/// Write queued lines in batches until the queue closes.
605async fn write_batches(
606    id: DaemonId,
607    relay_token: u64,
608    mut rx: tokio::sync::mpsc::Receiver<SinkEvent>,
609) {
610    let mut events: Vec<SinkEvent> = Vec::with_capacity(BATCH_SIZE);
611    let mut batch: Vec<ParsedLog> = Vec::with_capacity(BATCH_SIZE);
612    let mut flush_interval = tokio::time::interval(FLUSH_INTERVAL);
613    flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
614
615    loop {
616        let closed = tokio::select! {
617            received = rx.recv_many(&mut events, BATCH_SIZE) => received == 0,
618            _ = flush_interval.tick() => false,
619        };
620        for event in events.drain(..) {
621            match event {
622                SinkEvent::Line(parsed) => batch.push(parsed),
623                SinkEvent::Report(report) => {
624                    // Everything up to and including the reported line goes to
625                    // the store before the supervisor is told, so whatever it
626                    // does next can already read it.
627                    flush(&id, &mut batch).await;
628                    report_line(&id, relay_token, report).await;
629                }
630            }
631        }
632        flush(&id, &mut batch).await;
633        if closed {
634            break;
635        }
636    }
637}
638
639/// Hand the supervisor a line it needs to act on.
640///
641/// Failure is logged and dropped. There is no supervisor to retry against if it
642/// has crashed — and if it has, nothing is waiting on this daemon's readiness or
643/// hooks — while the daemon's output keeps being captured either way.
644async fn report_line(id: &DaemonId, relay_token: u64, report: Report) {
645    // `autostart: false` — a sink must never bring a supervisor into being.
646    match crate::ipc::client::IpcClient::connect(false).await {
647        Ok(client) => {
648            if let Err(e) = client
649                .sink_output_line(id.clone(), relay_token, report.fires_hook, report.text)
650                .await
651            {
652                warn!("log sink for {id} could not report a line of output: {e}");
653            }
654        }
655        Err(e) => {
656            warn!("log sink for {id} could not reach the supervisor to report output: {e}");
657        }
658    }
659}
660
661/// Write one batch, off the runtime so the SQLite call cannot stall other tasks.
662async fn flush(id: &DaemonId, batch: &mut Vec<ParsedLog>) {
663    if batch.is_empty() {
664        return;
665    }
666    let daemon_id = id.clone();
667    let entries = std::mem::take(batch);
668    let written = tokio::task::spawn_blocking(move || {
669        LOG_STORE.append_structured_batch(&daemon_id, &entries)
670    })
671    .await;
672    if let Ok(Err(e)) = written {
673        // Nothing useful to do but report it: the supervisor is not necessarily
674        // alive to be told, and dropping a batch is preferable to stalling the
675        // daemon behind a pipe nobody is draining.
676        error!("log sink failed to write batch for {id}: {e}");
677    }
678}
679
680#[cfg(test)]
681mod tests {
682    use super::{HookMatcher, MATCH_CARRY_BYTES, ReadyMatcher};
683
684    fn matcher(pattern: &str) -> ReadyMatcher {
685        ReadyMatcher::new(Some(regex::Regex::new(pattern).unwrap()), None)
686    }
687
688    fn hook_matcher(filter: Option<&str>, debounce_ms: u64) -> ReadyMatcher {
689        ReadyMatcher::new(
690            None,
691            Some(HookMatcher {
692                filter: filter.map(str::to_string),
693                regex: None,
694                debounce: std::time::Duration::from_millis(debounce_ms),
695                last_reported: None,
696            }),
697        )
698    }
699
700    #[test]
701    fn reports_the_first_match_and_then_stops_looking() {
702        let mut m = matcher("READY");
703        assert_eq!(m.consider("starting up", false), None);
704        assert_eq!(
705            m.consider("READY to serve", false)
706                .map(|r| r.text)
707                .as_deref(),
708            Some("READY to serve")
709        );
710        // Readiness happens once; later matches are somebody else's business.
711        assert_eq!(m.consider("READY again", false), None);
712    }
713
714    #[test]
715    fn matches_a_pattern_split_across_the_line_cap() {
716        // The daemon emitted one enormous line whose announcement straddles the
717        // point where the sink had to cut it.
718        let mut m = matcher("SERVER READY");
719        assert_eq!(m.consider("....SERVER ", true), None);
720        let report = m
721            .consider("READY....", false)
722            .expect("should match across the split");
723        assert!(
724            report.text.contains("SERVER READY"),
725            "reported {:?}",
726            report.text
727        );
728    }
729
730    #[test]
731    fn does_not_match_across_a_completed_line() {
732        // Two separate lines are not one line: a pattern spanning them must not
733        // match, or "SERVER" at the end of one line plus "READY" at the start of
734        // the next would look like an announcement.
735        let mut m = matcher("SERVER READY");
736        assert_eq!(m.consider("SERVER ", false), None);
737        assert_eq!(m.consider("READY", false), None);
738    }
739
740    #[test]
741    fn carries_a_bounded_amount_of_a_capped_line() {
742        let mut m = matcher("nothing-matches-this");
743        m.consider(&"x".repeat(MATCH_CARRY_BYTES * 3), true);
744        assert!(m.carried.len() <= MATCH_CARRY_BYTES);
745    }
746
747    #[test]
748    fn carrying_never_splits_a_character() {
749        // A carry boundary landing mid-character would panic on the slice, or
750        // corrupt the text a pattern is matched against.
751        let mut m = matcher("nothing-matches-this");
752        m.consider(&"é".repeat(MATCH_CARRY_BYTES), true);
753        assert!(m.carried.chars().all(|c| c == 'é'));
754    }
755
756    #[test]
757    fn hook_reports_matching_lines_within_the_debounce_window() {
758        let mut m = hook_matcher(Some("ALERT"), 0);
759        assert_eq!(m.consider("nothing here", false), None);
760        let report = m.consider("ALERT disk full", false).expect("should report");
761        assert_eq!(report.text, "ALERT disk full");
762        assert!(report.fires_hook);
763        // Unlike readiness, a hook keeps firing.
764        assert!(m.consider("ALERT again", false).is_some());
765    }
766
767    #[test]
768    fn hook_debounce_suppresses_a_second_line_in_the_same_window() {
769        let mut m = hook_matcher(None, 60_000);
770        assert!(m.consider("first", false).is_some());
771        // Every line matches a hook with no filter, so only the window stops it.
772        assert_eq!(m.consider("second", false), None);
773    }
774
775    #[test]
776    fn a_readiness_only_line_does_not_fire_the_hook() {
777        // A daemon can have both, filtered differently. The line that announces
778        // readiness is nothing to do with the hook, and reporting it must not
779        // be taken as permission to run the hook's command.
780        let mut m = ReadyMatcher::new(
781            Some(regex::Regex::new("READY").unwrap()),
782            Some(HookMatcher {
783                filter: Some("ALERT".to_string()),
784                regex: None,
785                debounce: std::time::Duration::from_millis(0),
786                last_reported: None,
787            }),
788        );
789        let report = m.consider("READY to serve", false).expect("should report");
790        assert!(!report.fires_hook, "readiness line must not fire the hook");
791        // ...and a line matching both says so.
792        let mut m = ReadyMatcher::new(
793            Some(regex::Regex::new("READY").unwrap()),
794            Some(HookMatcher {
795                filter: Some("READY".to_string()),
796                regex: None,
797                debounce: std::time::Duration::from_millis(0),
798                last_reported: None,
799            }),
800        );
801        assert!(
802            m.consider("READY", false)
803                .expect("should report")
804                .fires_hook
805        );
806    }
807
808    #[test]
809    fn strips_ansi_before_matching() {
810        let mut m = matcher("^READY$");
811        assert!(m.consider("\x1b[32mREADY\x1b[0m", false).is_some());
812    }
813
814    #[test]
815    fn decodes_utf8_unchanged() {
816        assert_eq!(
817            super::decode_line("起動しました".as_bytes()),
818            "起動しました"
819        );
820    }
821
822    /// What `cmd /C exec` writes on a Japanese system: "'exec' は、内部コマンド".
823    const CP932_CMD_ERROR: &[u8] =
824        b"'exec' \x82\xcd\x81\x41\x93\xe0\x95\x94\x83\x52\x83\x7d\x83\x93\x83\x68";
825
826    #[cfg(windows)]
827    #[test]
828    fn decodes_the_console_code_page() {
829        assert_eq!(
830            super::decode_in_code_page(CP932_CMD_ERROR, 932).as_deref(),
831            Some("'exec' は、内部コマンド")
832        );
833    }
834
835    #[cfg(windows)]
836    #[test]
837    fn leaves_a_utf8_code_page_to_the_lossy_conversion() {
838        assert_eq!(super::decode_in_code_page(b"\xff", 65001), None);
839    }
840
841    #[cfg(not(windows))]
842    #[test]
843    fn replaces_non_utf8_outside_windows() {
844        assert_eq!(super::decode_line(b"ok \xff"), "ok \u{fffd}");
845    }
846
847    /// The second line of the same message: "操作可能なプログラムまたはバッチ
848    /// ファイルとして認識されていません。" Only the Windows tests decode it.
849    #[cfg(windows)]
850    const CP932_CMD_ERROR_2: &[u8] = b"\x91\x80\x8d\xec\x89\xc2\x94\x5c\x82\xc8\x83\x76\x83\x8d\x83\x4f\x83\x89\x83\x80\x82\xdc\x82\xbd\x82\xcd\x83\x6f\x83\x62\x83\x60\x20\x83\x74\x83\x40\x83\x43\x83\x8b\x82\xc6\x82\xb5\x82\xc4\x94\x46\x8e\xaf\x82\xb3\x82\xea\x82\xc4\x82\xa2\x82\xdc\x82\xb9\x82\xf1\x81\x42";
851
852    #[test]
853    fn keeps_utf8_text_around_a_stray_byte() {
854        // Decoding this in a code page would garble every character; only the
855        // stray byte should be lost.
856        let mut line = "起動しました ".as_bytes().to_vec();
857        line.push(0xff);
858        assert_eq!(super::decode_line(&line), "起動しました \u{fffd}");
859    }
860
861    #[cfg(windows)]
862    #[test]
863    fn reads_code_page_text_in_the_code_page() {
864        // The first line happens to contain a valid UTF-8 character (0xCD 0x81).
865        assert!(super::uses_code_page(CP932_CMD_ERROR));
866        assert!(super::uses_code_page(CP932_CMD_ERROR_2));
867        assert_eq!(
868            super::decode_in_code_page(CP932_CMD_ERROR_2, 932).as_deref(),
869            Some("操作可能なプログラムまたはバッチ ファイルとして認識されていません。")
870        );
871    }
872
873    #[cfg(windows)]
874    #[test]
875    fn splits_before_an_unfinished_double_byte_character() {
876        // "aaaあ" in CP932, cut by the length cap between あ's two bytes.
877        let line = b"aaa\x82\xa0";
878        let split = super::split_before_incomplete_dbcs(&line[..4], 932);
879        assert_eq!(split, 3);
880        let first = super::decode_in_code_page(&line[..split], 932).unwrap();
881        let rest = super::decode_in_code_page(&line[split..], 932).unwrap();
882        assert_eq!(first + &rest, "aaaあ");
883    }
884
885    #[cfg(windows)]
886    #[test]
887    fn a_line_read_in_the_code_page_keeps_an_unfinished_utf8_character() {
888        // A stray byte sends the piece to the code page, but the line ends in
889        // the first byte of "é" (0xC3 0xA9), which is a whole character in
890        // CP932. Cutting there would still split the UTF-8 character.
891        let line = b"ab\xffcd\xc3";
892        assert!(super::uses_code_page(&line[..5]));
893        assert_eq!(super::split_before_incomplete_char_in(line, 932), 5);
894        // A double-byte character left unfinished is still held back too.
895        assert_eq!(
896            super::split_before_incomplete_char_in(b"ab\xffcd\x82", 932),
897            5
898        );
899    }
900
901    #[cfg(windows)]
902    #[test]
903    fn a_complete_double_byte_character_is_not_cut_for_utf8() {
904        // "づ" is 0x82 0xC3 in CP932. Its second byte looks like the start of
905        // a UTF-8 character, but cutting before it would split a character
906        // the code page reads as whole.
907        let line = b"ab\xffcd\x82\xc3";
908        assert_eq!(super::split_before_incomplete_char_in(line, 932), 7);
909        let line = b"aa\xe0\xc3";
910        assert_eq!(super::split_before_incomplete_char_in(line, 932), 4);
911    }
912
913    #[cfg(windows)]
914    #[test]
915    fn a_trail_byte_that_looks_like_a_lead_byte_is_not_cut() {
916        // 0x82 is a lead byte, but here the second one completes the first.
917        assert_eq!(super::split_before_incomplete_dbcs(b"a\x82\x82", 932), 3);
918        // ...and here the first pair is whole, leaving the last 0x82 unfinished.
919        assert_eq!(super::split_before_incomplete_dbcs(b"\x82\x9f\x82", 932), 2);
920        // A single-byte code page has nothing to finish.
921        assert_eq!(super::split_before_incomplete_dbcs(b"aa\x82", 1252), 3);
922    }
923
924    #[test]
925    fn never_fails_on_non_utf8() {
926        // Whatever the code page, the line is still stored.
927        assert!(super::decode_line(CP932_CMD_ERROR).starts_with("'exec' "));
928    }
929}