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, clap::Args)]
#[clap(hide = true, verbatim_doc_comment)]
pub struct LogSink {
#[clap(long)]
daemon_id: String,
#[clap(long, default_value = "text")]
log_format: String,
}
impl LogSink {
pub async fn run(&self) -> Result<()> {
let id = DaemonId::parse(&self.daemon_id)?;
let (tx, rx) = tokio::sync::mpsc::channel::<ParsedLog>(QUEUE_DEPTH);
let writer = tokio::spawn(write_batches(id.clone(), rx));
let read_result = read_lines(tx, &self.log_format).await;
let _ = writer.await;
read_result.map_err(|e| {
miette::miette!("log sink for {id} could not read the daemon's output: {e}")
})
}
}
async fn read_lines(
tx: tokio::sync::mpsc::Sender<ParsedLog>,
log_format: &str,
) -> 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).await?;
} else {
line.push(byte);
split_at_cap = false;
if line.len() >= MAX_LINE_BYTES {
queue_capped(&tx, &mut line, log_format).await?;
split_at_cap = true;
}
}
}
}
if !line.is_empty() {
queue(&tx, &mut line, log_format).await?;
}
Ok(())
}
async fn queue_capped(
tx: &tokio::sync::mpsc::Sender<ParsedLog>,
line: &mut Vec<u8>,
log_format: &str,
) -> std::io::Result<()> {
let split = split_before_incomplete_char(line);
let tail = line.split_off(split);
let result = queue(tx, line, log_format).await;
*line = tail;
result
}
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<ParsedLog>,
line: &mut Vec<u8>,
log_format: &str,
) -> std::io::Result<()> {
let text = String::from_utf8_lossy(line);
let parsed = crate::log_parse::parse(text.trim_end_matches('\r'), log_format);
line.clear();
tx.send(parsed)
.await
.map_err(|_| std::io::Error::other("log writer stopped"))
}
async fn write_batches(id: DaemonId, mut rx: tokio::sync::mpsc::Receiver<ParsedLog>) {
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 batch, BATCH_SIZE) => received == 0,
_ = flush_interval.tick() => false,
};
flush(&id, &mut batch).await;
if closed {
break;
}
}
}
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}");
}
}