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 = 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}