use crate::Result;
use crate::daemon_id::DaemonId;
use crate::log_parse::ParsedLog;
use crate::log_store::LogStore;
use crate::log_store::sqlite::LOG_STORE;
use tokio::io::AsyncReadExt;
const BATCH_SIZE: usize = 100;
const FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
const QUEUE_DEPTH: usize = 8192;
const READ_CHUNK: usize = 8192;
const MAX_LINE_BYTES: usize = 64 * 1024;
#[derive(Debug, usage_rs::Args)]
#[usage(verbatim_doc_comment)]
pub struct LogSink {
#[usage(long)]
daemon_id: String,
#[usage(long, default = "text")]
log_format: String,
#[usage(long)]
ready_pattern: Option<String>,
#[usage(long, default_value_t = 0, default = "0")]
relay_token: u64,
#[usage(long)]
report_output: bool,
#[usage(long)]
output_filter: Option<String>,
#[usage(long)]
output_regex: Option<String>,
#[usage(long, default_value_t = 1000, default = "1000")]
output_debounce_ms: u64,
}
impl LogSink {
pub async fn run(&self) -> Result<()> {
let id = DaemonId::parse(&self.daemon_id)?;
let compile = |what: &str, pattern: &str| {
regex::Regex::new(pattern)
.map_err(|e| error!("log sink for {id} ignoring unparsable {what}: {e}"))
.ok()
};
let ready_pattern = self
.ready_pattern
.as_deref()
.and_then(|p| compile("ready pattern", p));
let hook = self.report_output.then(|| HookMatcher {
filter: self.output_filter.clone(),
regex: self
.output_regex
.as_deref()
.and_then(|p| compile("output pattern", p)),
debounce: std::time::Duration::from_millis(self.output_debounce_ms),
last_reported: None,
});
let (tx, rx) = tokio::sync::mpsc::channel::<SinkEvent>(QUEUE_DEPTH);
let writer = tokio::spawn(write_batches(id.clone(), self.relay_token, rx));
let read_result =
read_lines(tx, &self.log_format, ReadyMatcher::new(ready_pattern, hook)).await;
let _ = writer.await;
read_result.map_err(|e| {
miette::miette!("log sink for {id} could not read the daemon's output: {e}")
})
}
}
enum SinkEvent {
Line(ParsedLog),
Report(Report),
}
#[derive(Debug, PartialEq, Eq)]
struct Report {
text: String,
fires_hook: bool,
}
const MATCH_CARRY_BYTES: usize = 4 * 1024;
struct ReadyMatcher {
pattern: Option<regex::Regex>,
hook: Option<HookMatcher>,
carried: String,
}
struct HookMatcher {
filter: Option<String>,
regex: Option<regex::Regex>,
debounce: std::time::Duration,
last_reported: Option<std::time::Instant>,
}
impl HookMatcher {
fn matches(&mut self, clean: &str) -> bool {
let matched = match (&self.filter, &self.regex) {
(Some(substr), _) => clean.contains(substr.as_str()),
(None, Some(re)) => re.is_match(clean),
(None, None) => true,
};
if !matched {
return false;
}
let now = std::time::Instant::now();
if self
.last_reported
.is_some_and(|last| now.duration_since(last) < self.debounce)
{
return false;
}
self.last_reported = Some(now);
true
}
}
impl ReadyMatcher {
fn new(pattern: Option<regex::Regex>, hook: Option<HookMatcher>) -> Self {
Self {
pattern,
hook,
carried: String::new(),
}
}
fn is_watching(&self) -> bool {
self.pattern.is_some() || self.hook.is_some()
}
fn consider(&mut self, text: &str, split_at_cap: bool) -> Option<Report> {
if !self.is_watching() {
self.carried.clear();
return None;
}
let clean = console::strip_ansi_codes(text);
let candidate = if self.carried.is_empty() {
clean.into_owned()
} else {
format!("{}{clean}", self.carried)
};
self.carried = if split_at_cap {
let start = candidate.len().saturating_sub(MATCH_CARRY_BYTES);
let start = (start..candidate.len())
.find(|i| candidate.is_char_boundary(*i))
.unwrap_or(candidate.len());
candidate[start..].to_string()
} else {
String::new()
};
let mut ready_matched = false;
if self
.pattern
.as_ref()
.is_some_and(|re| re.is_match(&candidate))
{
self.pattern = None;
ready_matched = true;
}
let fires_hook = self
.hook
.as_mut()
.is_some_and(|hook| hook.matches(&candidate));
(ready_matched || fires_hook).then_some(Report {
text: candidate,
fires_hook,
})
}
}
async fn read_lines(
tx: tokio::sync::mpsc::Sender<SinkEvent>,
log_format: &str,
mut matcher: ReadyMatcher,
) -> std::io::Result<()> {
let mut stdin = tokio::io::stdin();
let mut chunk = vec![0u8; READ_CHUNK];
let mut line: Vec<u8> = Vec::with_capacity(256);
let mut split_at_cap = false;
loop {
let read = stdin.read(&mut chunk).await?;
if read == 0 {
break;
}
for &byte in &chunk[..read] {
if byte == b'\n' {
if split_at_cap && line.is_empty() {
split_at_cap = false;
continue;
}
split_at_cap = false;
queue(&tx, &mut line, log_format, &mut matcher).await?;
} else {
line.push(byte);
split_at_cap = false;
if line.len() >= MAX_LINE_BYTES {
queue_capped(&tx, &mut line, log_format, &mut matcher).await?;
split_at_cap = true;
}
}
}
}
if !line.is_empty() {
queue(&tx, &mut line, log_format, &mut matcher).await?;
}
Ok(())
}
async fn queue_capped(
tx: &tokio::sync::mpsc::Sender<SinkEvent>,
line: &mut Vec<u8>,
log_format: &str,
matcher: &mut ReadyMatcher,
) -> std::io::Result<()> {
let split = split_before_incomplete_char(line);
let tail = line.split_off(split);
let result = queue_piece(tx, line, log_format, matcher).await;
*line = tail;
result
}
async fn queue_piece(
tx: &tokio::sync::mpsc::Sender<SinkEvent>,
line: &mut Vec<u8>,
log_format: &str,
matcher: &mut ReadyMatcher,
) -> std::io::Result<()> {
let text = String::from_utf8_lossy(line);
let text = text.trim_end_matches('\r');
let parsed = crate::log_parse::parse(text, log_format);
let report = matcher.consider(text, true);
line.clear();
tx.send(SinkEvent::Line(parsed))
.await
.map_err(|_| std::io::Error::other("log writer stopped"))?;
if let Some(report) = report {
tx.send(SinkEvent::Report(report))
.await
.map_err(|_| std::io::Error::other("log writer stopped"))?;
}
Ok(())
}
fn split_before_incomplete_char(bytes: &[u8]) -> usize {
let len = bytes.len();
for i in (len.saturating_sub(4)..len).rev() {
let byte = bytes[i];
if byte & 0b1100_0000 == 0b1000_0000 {
continue; }
let expected = match byte {
0x00..=0x7f => 1,
b if b >> 5 == 0b110 => 2,
b if b >> 4 == 0b1110 => 3,
b if b >> 3 == 0b11110 => 4,
_ => 1,
};
return if i + expected > len && i > 0 { i } else { len };
}
len
}
async fn queue(
tx: &tokio::sync::mpsc::Sender<SinkEvent>,
line: &mut Vec<u8>,
log_format: &str,
matcher: &mut ReadyMatcher,
) -> std::io::Result<()> {
let text = String::from_utf8_lossy(line);
let text = text.trim_end_matches('\r');
let parsed = crate::log_parse::parse(text, log_format);
let report = matcher.consider(text, false);
line.clear();
tx.send(SinkEvent::Line(parsed))
.await
.map_err(|_| std::io::Error::other("log writer stopped"))?;
if let Some(report) = report {
tx.send(SinkEvent::Report(report))
.await
.map_err(|_| std::io::Error::other("log writer stopped"))?;
}
Ok(())
}
async fn write_batches(
id: DaemonId,
relay_token: u64,
mut rx: tokio::sync::mpsc::Receiver<SinkEvent>,
) {
let mut events: Vec<SinkEvent> = Vec::with_capacity(BATCH_SIZE);
let mut batch: Vec<ParsedLog> = Vec::with_capacity(BATCH_SIZE);
let mut flush_interval = tokio::time::interval(FLUSH_INTERVAL);
flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
let closed = tokio::select! {
received = rx.recv_many(&mut events, BATCH_SIZE) => received == 0,
_ = flush_interval.tick() => false,
};
for event in events.drain(..) {
match event {
SinkEvent::Line(parsed) => batch.push(parsed),
SinkEvent::Report(report) => {
flush(&id, &mut batch).await;
report_line(&id, relay_token, report).await;
}
}
}
flush(&id, &mut batch).await;
if closed {
break;
}
}
}
async fn report_line(id: &DaemonId, relay_token: u64, report: Report) {
match crate::ipc::client::IpcClient::connect(false).await {
Ok(client) => {
if let Err(e) = client
.sink_output_line(id.clone(), relay_token, report.fires_hook, report.text)
.await
{
warn!("log sink for {id} could not report a line of output: {e}");
}
}
Err(e) => {
warn!("log sink for {id} could not reach the supervisor to report output: {e}");
}
}
}
async fn flush(id: &DaemonId, batch: &mut Vec<ParsedLog>) {
if batch.is_empty() {
return;
}
let daemon_id = id.clone();
let entries = std::mem::take(batch);
let written = tokio::task::spawn_blocking(move || {
LOG_STORE.append_structured_batch(&daemon_id, &entries)
})
.await;
if let Ok(Err(e)) = written {
error!("log sink failed to write batch for {id}: {e}");
}
}
#[cfg(test)]
mod tests {
use super::{HookMatcher, MATCH_CARRY_BYTES, ReadyMatcher};
fn matcher(pattern: &str) -> ReadyMatcher {
ReadyMatcher::new(Some(regex::Regex::new(pattern).unwrap()), None)
}
fn hook_matcher(filter: Option<&str>, debounce_ms: u64) -> ReadyMatcher {
ReadyMatcher::new(
None,
Some(HookMatcher {
filter: filter.map(str::to_string),
regex: None,
debounce: std::time::Duration::from_millis(debounce_ms),
last_reported: None,
}),
)
}
#[test]
fn reports_the_first_match_and_then_stops_looking() {
let mut m = matcher("READY");
assert_eq!(m.consider("starting up", false), None);
assert_eq!(
m.consider("READY to serve", false)
.map(|r| r.text)
.as_deref(),
Some("READY to serve")
);
assert_eq!(m.consider("READY again", false), None);
}
#[test]
fn matches_a_pattern_split_across_the_line_cap() {
let mut m = matcher("SERVER READY");
assert_eq!(m.consider("....SERVER ", true), None);
let report = m
.consider("READY....", false)
.expect("should match across the split");
assert!(
report.text.contains("SERVER READY"),
"reported {:?}",
report.text
);
}
#[test]
fn does_not_match_across_a_completed_line() {
let mut m = matcher("SERVER READY");
assert_eq!(m.consider("SERVER ", false), None);
assert_eq!(m.consider("READY", false), None);
}
#[test]
fn carries_a_bounded_amount_of_a_capped_line() {
let mut m = matcher("nothing-matches-this");
m.consider(&"x".repeat(MATCH_CARRY_BYTES * 3), true);
assert!(m.carried.len() <= MATCH_CARRY_BYTES);
}
#[test]
fn carrying_never_splits_a_character() {
let mut m = matcher("nothing-matches-this");
m.consider(&"é".repeat(MATCH_CARRY_BYTES), true);
assert!(m.carried.chars().all(|c| c == 'é'));
}
#[test]
fn hook_reports_matching_lines_within_the_debounce_window() {
let mut m = hook_matcher(Some("ALERT"), 0);
assert_eq!(m.consider("nothing here", false), None);
let report = m.consider("ALERT disk full", false).expect("should report");
assert_eq!(report.text, "ALERT disk full");
assert!(report.fires_hook);
assert!(m.consider("ALERT again", false).is_some());
}
#[test]
fn hook_debounce_suppresses_a_second_line_in_the_same_window() {
let mut m = hook_matcher(None, 60_000);
assert!(m.consider("first", false).is_some());
assert_eq!(m.consider("second", false), None);
}
#[test]
fn a_readiness_only_line_does_not_fire_the_hook() {
let mut m = ReadyMatcher::new(
Some(regex::Regex::new("READY").unwrap()),
Some(HookMatcher {
filter: Some("ALERT".to_string()),
regex: None,
debounce: std::time::Duration::from_millis(0),
last_reported: None,
}),
);
let report = m.consider("READY to serve", false).expect("should report");
assert!(!report.fires_hook, "readiness line must not fire the hook");
let mut m = ReadyMatcher::new(
Some(regex::Regex::new("READY").unwrap()),
Some(HookMatcher {
filter: Some("READY".to_string()),
regex: None,
debounce: std::time::Duration::from_millis(0),
last_reported: None,
}),
);
assert!(
m.consider("READY", false)
.expect("should report")
.fires_hook
);
}
#[test]
fn strips_ansi_before_matching() {
let mut m = matcher("^READY$");
assert!(m.consider("\x1b[32mREADY\x1b[0m", false).is_some());
}
}