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, clap::Args)]
42#[clap(hide = true, verbatim_doc_comment)]
43pub struct LogSink {
44    /// Qualified id of the daemon whose output this is
45    #[clap(long)]
46    daemon_id: String,
47
48    /// Log format to parse lines with (`json`, `logfmt`, `auto`, or `text`)
49    #[clap(long, default_value = "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    #[clap(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    #[clap(long, default_value_t = 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    #[clap(long)]
73    report_output: bool,
74
75    /// Only report lines containing this substring
76    #[clap(long)]
77    output_filter: Option<String>,
78
79    /// Only report lines matching this regex
80    #[clap(long)]
81    output_regex: Option<String>,
82
83    /// Shortest gap between reported lines, in milliseconds
84    #[clap(long, default_value_t = 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 = String::from_utf8_lossy(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.
385///
386/// Decided by inspecting the final bytes rather than by asking `from_utf8` where
387/// the string stops being valid: that reports the *first* problem, so a single
388/// invalid byte earlier in the line would hide an unfinished character at the
389/// end, and the character would be split after all.
390fn split_before_incomplete_char(bytes: &[u8]) -> usize {
391    let len = bytes.len();
392    // A character is at most four bytes, so only the last few can be unfinished.
393    for i in (len.saturating_sub(4)..len).rev() {
394        let byte = bytes[i];
395        if byte & 0b1100_0000 == 0b1000_0000 {
396            continue; // a continuation byte; keep looking back for its lead
397        }
398        let expected = match byte {
399            0x00..=0x7f => 1,
400            b if b >> 5 == 0b110 => 2,
401            b if b >> 4 == 0b1110 => 3,
402            b if b >> 3 == 0b11110 => 4,
403            // Not a valid lead byte at all, so nothing is pending; the lossy
404            // conversion will render it.
405            _ => 1,
406        };
407        return if i + expected > len && i > 0 { i } else { len };
408    }
409    len
410}
411
412/// Parse `line` and hand it to the writer, clearing it either way.
413///
414/// A closed queue means the writer task is gone, which is a failure rather than
415/// the end of the stream: reporting it as success would tell the supervisor this
416/// sink had reached end of file, and it would stop replacing it while the daemon
417/// was still writing.
418async fn queue(
419    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
420    line: &mut Vec<u8>,
421    log_format: &str,
422    matcher: &mut ReadyMatcher,
423) -> std::io::Result<()> {
424    // Convert lossily: a daemon emitting a stray non-UTF-8 byte must not be able
425    // to stop its own logging.
426    let text = String::from_utf8_lossy(line);
427    let text = text.trim_end_matches('\r');
428    let parsed = crate::log_parse::parse(text, log_format);
429    // Strip ANSI before matching so a pattern works whether or not the daemon
430    // colours its output, matching what in-process capture did.
431    let report = matcher.consider(text, false);
432    line.clear();
433
434    tx.send(SinkEvent::Line(parsed))
435        .await
436        .map_err(|_| std::io::Error::other("log writer stopped"))?;
437    if let Some(report) = report {
438        tx.send(SinkEvent::Report(report))
439            .await
440            .map_err(|_| std::io::Error::other("log writer stopped"))?;
441    }
442    Ok(())
443}
444
445/// Write queued lines in batches until the queue closes.
446async fn write_batches(
447    id: DaemonId,
448    relay_token: u64,
449    mut rx: tokio::sync::mpsc::Receiver<SinkEvent>,
450) {
451    let mut events: Vec<SinkEvent> = Vec::with_capacity(BATCH_SIZE);
452    let mut batch: Vec<ParsedLog> = Vec::with_capacity(BATCH_SIZE);
453    let mut flush_interval = tokio::time::interval(FLUSH_INTERVAL);
454    flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
455
456    loop {
457        let closed = tokio::select! {
458            received = rx.recv_many(&mut events, BATCH_SIZE) => received == 0,
459            _ = flush_interval.tick() => false,
460        };
461        for event in events.drain(..) {
462            match event {
463                SinkEvent::Line(parsed) => batch.push(parsed),
464                SinkEvent::Report(report) => {
465                    // Everything up to and including the reported line goes to
466                    // the store before the supervisor is told, so whatever it
467                    // does next can already read it.
468                    flush(&id, &mut batch).await;
469                    report_line(&id, relay_token, report).await;
470                }
471            }
472        }
473        flush(&id, &mut batch).await;
474        if closed {
475            break;
476        }
477    }
478}
479
480/// Hand the supervisor a line it needs to act on.
481///
482/// Failure is logged and dropped. There is no supervisor to retry against if it
483/// has crashed — and if it has, nothing is waiting on this daemon's readiness or
484/// hooks — while the daemon's output keeps being captured either way.
485async fn report_line(id: &DaemonId, relay_token: u64, report: Report) {
486    // `autostart: false` — a sink must never bring a supervisor into being.
487    match crate::ipc::client::IpcClient::connect(false).await {
488        Ok(client) => {
489            if let Err(e) = client
490                .sink_output_line(id.clone(), relay_token, report.fires_hook, report.text)
491                .await
492            {
493                warn!("log sink for {id} could not report a line of output: {e}");
494            }
495        }
496        Err(e) => {
497            warn!("log sink for {id} could not reach the supervisor to report output: {e}");
498        }
499    }
500}
501
502/// Write one batch, off the runtime so the SQLite call cannot stall other tasks.
503async fn flush(id: &DaemonId, batch: &mut Vec<ParsedLog>) {
504    if batch.is_empty() {
505        return;
506    }
507    let daemon_id = id.clone();
508    let entries = std::mem::take(batch);
509    let written = tokio::task::spawn_blocking(move || {
510        LOG_STORE.append_structured_batch(&daemon_id, &entries)
511    })
512    .await;
513    if let Ok(Err(e)) = written {
514        // Nothing useful to do but report it: the supervisor is not necessarily
515        // alive to be told, and dropping a batch is preferable to stalling the
516        // daemon behind a pipe nobody is draining.
517        error!("log sink failed to write batch for {id}: {e}");
518    }
519}
520
521#[cfg(test)]
522mod tests {
523    use super::{HookMatcher, MATCH_CARRY_BYTES, ReadyMatcher};
524
525    fn matcher(pattern: &str) -> ReadyMatcher {
526        ReadyMatcher::new(Some(regex::Regex::new(pattern).unwrap()), None)
527    }
528
529    fn hook_matcher(filter: Option<&str>, debounce_ms: u64) -> ReadyMatcher {
530        ReadyMatcher::new(
531            None,
532            Some(HookMatcher {
533                filter: filter.map(str::to_string),
534                regex: None,
535                debounce: std::time::Duration::from_millis(debounce_ms),
536                last_reported: None,
537            }),
538        )
539    }
540
541    #[test]
542    fn reports_the_first_match_and_then_stops_looking() {
543        let mut m = matcher("READY");
544        assert_eq!(m.consider("starting up", false), None);
545        assert_eq!(
546            m.consider("READY to serve", false)
547                .map(|r| r.text)
548                .as_deref(),
549            Some("READY to serve")
550        );
551        // Readiness happens once; later matches are somebody else's business.
552        assert_eq!(m.consider("READY again", false), None);
553    }
554
555    #[test]
556    fn matches_a_pattern_split_across_the_line_cap() {
557        // The daemon emitted one enormous line whose announcement straddles the
558        // point where the sink had to cut it.
559        let mut m = matcher("SERVER READY");
560        assert_eq!(m.consider("....SERVER ", true), None);
561        let report = m
562            .consider("READY....", false)
563            .expect("should match across the split");
564        assert!(
565            report.text.contains("SERVER READY"),
566            "reported {:?}",
567            report.text
568        );
569    }
570
571    #[test]
572    fn does_not_match_across_a_completed_line() {
573        // Two separate lines are not one line: a pattern spanning them must not
574        // match, or "SERVER" at the end of one line plus "READY" at the start of
575        // the next would look like an announcement.
576        let mut m = matcher("SERVER READY");
577        assert_eq!(m.consider("SERVER ", false), None);
578        assert_eq!(m.consider("READY", false), None);
579    }
580
581    #[test]
582    fn carries_a_bounded_amount_of_a_capped_line() {
583        let mut m = matcher("nothing-matches-this");
584        m.consider(&"x".repeat(MATCH_CARRY_BYTES * 3), true);
585        assert!(m.carried.len() <= MATCH_CARRY_BYTES);
586    }
587
588    #[test]
589    fn carrying_never_splits_a_character() {
590        // A carry boundary landing mid-character would panic on the slice, or
591        // corrupt the text a pattern is matched against.
592        let mut m = matcher("nothing-matches-this");
593        m.consider(&"é".repeat(MATCH_CARRY_BYTES), true);
594        assert!(m.carried.chars().all(|c| c == 'é'));
595    }
596
597    #[test]
598    fn hook_reports_matching_lines_within_the_debounce_window() {
599        let mut m = hook_matcher(Some("ALERT"), 0);
600        assert_eq!(m.consider("nothing here", false), None);
601        let report = m.consider("ALERT disk full", false).expect("should report");
602        assert_eq!(report.text, "ALERT disk full");
603        assert!(report.fires_hook);
604        // Unlike readiness, a hook keeps firing.
605        assert!(m.consider("ALERT again", false).is_some());
606    }
607
608    #[test]
609    fn hook_debounce_suppresses_a_second_line_in_the_same_window() {
610        let mut m = hook_matcher(None, 60_000);
611        assert!(m.consider("first", false).is_some());
612        // Every line matches a hook with no filter, so only the window stops it.
613        assert_eq!(m.consider("second", false), None);
614    }
615
616    #[test]
617    fn a_readiness_only_line_does_not_fire_the_hook() {
618        // A daemon can have both, filtered differently. The line that announces
619        // readiness is nothing to do with the hook, and reporting it must not
620        // be taken as permission to run the hook's command.
621        let mut m = ReadyMatcher::new(
622            Some(regex::Regex::new("READY").unwrap()),
623            Some(HookMatcher {
624                filter: Some("ALERT".to_string()),
625                regex: None,
626                debounce: std::time::Duration::from_millis(0),
627                last_reported: None,
628            }),
629        );
630        let report = m.consider("READY to serve", false).expect("should report");
631        assert!(!report.fires_hook, "readiness line must not fire the hook");
632        // ...and a line matching both says so.
633        let mut m = ReadyMatcher::new(
634            Some(regex::Regex::new("READY").unwrap()),
635            Some(HookMatcher {
636                filter: Some("READY".to_string()),
637                regex: None,
638                debounce: std::time::Duration::from_millis(0),
639                last_reported: None,
640            }),
641        );
642        assert!(
643            m.consider("READY", false)
644                .expect("should report")
645                .fires_hook
646        );
647    }
648
649    #[test]
650    fn strips_ansi_before_matching() {
651        let mut m = matcher("^READY$");
652        assert!(m.consider("\x1b[32mREADY\x1b[0m", false).is_some());
653    }
654}