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
8const BATCH_SIZE: usize = 100;
10
11const FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
13
14const QUEUE_DEPTH: usize = 8192;
21
22const READ_CHUNK: usize = 8192;
24
25const MAX_LINE_BYTES: usize = 64 * 1024;
32
33#[derive(Debug, usage_rs::Args)]
42#[usage(verbatim_doc_comment)]
43pub struct LogSink {
44 #[usage(long)]
46 daemon_id: String,
47
48 #[usage(long, default = "text")]
50 log_format: String,
51
52 #[usage(long)]
58 ready_pattern: Option<String>,
59
60 #[usage(long, default_value_t = 0, default = "0")]
66 relay_token: u64,
67
68 #[usage(long)]
73 report_output: bool,
74
75 #[usage(long)]
77 output_filter: Option<String>,
78
79 #[usage(long)]
81 output_regex: Option<String>,
82
83 #[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 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 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 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
135enum SinkEvent {
142 Line(ParsedLog),
143 Report(Report),
144}
145
146#[derive(Debug, PartialEq, Eq)]
148struct Report {
149 text: String,
150 fires_hook: bool,
154}
155
156const MATCH_CARRY_BYTES: usize = 4 * 1024;
164
165struct ReadyMatcher {
172 pattern: Option<regex::Regex>,
174 hook: Option<HookMatcher>,
176 carried: String,
179}
180
181struct 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 fn matches(&mut self, clean: &str) -> bool {
197 let matched = match (&self.filter, &self.regex) {
198 (Some(substr), _) => clean.contains(substr.as_str()),
201 (None, Some(re)) => re.is_match(clean),
202 (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 fn is_watching(&self) -> bool {
232 self.pattern.is_some() || self.hook.is_some()
233 }
234
235 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 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 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 (ready_matched || fires_hook).then_some(Report {
283 text: candidate,
284 fires_hook,
285 })
286 }
287}
288
289async 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 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 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 if !line.is_empty() {
335 queue(&tx, &mut line, log_format, &mut matcher).await?;
336 }
337 Ok(())
338}
339
340async 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
359async 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
384fn 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#[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 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
414fn split_before_incomplete_utf8(bytes: &[u8]) -> usize {
421 let len = bytes.len();
422 for i in (len.saturating_sub(4)..len).rev() {
424 let byte = bytes[i];
425 if byte & 0b1100_0000 == 0b1000_0000 {
426 continue; }
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 _ => 1,
436 };
437 return if i + expected > len && i > 0 { i } else { len };
438 }
439 len
440}
441
442pub(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#[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#[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#[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 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#[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 if i > bytes.len() && bytes.len() > 1 {
550 bytes.len() - 1
551 } else {
552 bytes.len()
553 }
554}
555
556#[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
571async 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 let text = decode_line(line);
586 let text = text.trim_end_matches('\r');
587 let parsed = crate::log_parse::parse(text, log_format);
588 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
604async 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 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
639async fn report_line(id: &DaemonId, relay_token: u64, report: Report) {
645 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
661async 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 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 assert_eq!(m.consider("READY again", false), None);
712 }
713
714 #[test]
715 fn matches_a_pattern_split_across_the_line_cap() {
716 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 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 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 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 assert_eq!(m.consider("second", false), None);
773 }
774
775 #[test]
776 fn a_readiness_only_line_does_not_fire_the_hook() {
777 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 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 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 #[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 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 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 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 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 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 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 assert_eq!(super::split_before_incomplete_dbcs(b"a\x82\x82", 932), 3);
918 assert_eq!(super::split_before_incomplete_dbcs(b"\x82\x9f\x82", 932), 2);
920 assert_eq!(super::split_before_incomplete_dbcs(b"aa\x82", 1252), 3);
922 }
923
924 #[test]
925 fn never_fails_on_non_utf8() {
926 assert!(super::decode_line(CP932_CMD_ERROR).starts_with("'exec' "));
928 }
929}