mod render;
mod stats_cmd;
use std::fs::File;
use std::io::{self, BufRead, BufReader, IsTerminal, Write};
use std::path::PathBuf;
use std::process::ExitCode;
use std::time::{Duration, Instant};
use clap::{Parser, Subcommand};
use tracing_subscriber::EnvFilter;
use logdive_core::{
BATCH_SIZE, Indexer, InsertStats, LogEntry, LogdiveError, QueryParseError, Result, db_path,
execute, parse_line, parse_query,
};
use crate::render::{OutputFormat, render};
use crate::stats_cmd::{StatsArgs, run_stats};
#[derive(Parser, Debug)]
#[command(
name = "logdive",
version,
about = "Fast, self-hosted query engine for structured JSON logs",
long_about = None,
)]
struct Cli {
#[arg(long, global = true, value_name = "PATH")]
db: Option<PathBuf>,
#[command(subcommand)]
command: Command,
}
#[derive(Subcommand, Debug)]
enum Command {
Ingest(IngestArgs),
Query(QueryArgs),
Stats(StatsArgs),
}
#[derive(clap::Args, Debug)]
struct IngestArgs {
#[arg(long, short = 'f', value_name = "PATH")]
file: Option<PathBuf>,
#[arg(long, short = 't', value_name = "TAG")]
tag: Option<String>,
}
#[derive(clap::Args, Debug)]
struct QueryArgs {
#[arg(value_name = "QUERY")]
query: String,
#[arg(long, value_enum, default_value_t = OutputFormat::Pretty)]
format: OutputFormat,
#[arg(long, default_value_t = 1000)]
limit: usize,
}
fn main() -> ExitCode {
init_tracing();
let cli = Cli::parse();
match run(cli) {
Ok(()) => ExitCode::SUCCESS,
Err(e) => {
report_error(&e);
ExitCode::FAILURE
}
}
}
fn report_error(e: &LogdiveError) {
if let LogdiveError::QueryParse(qpe) = e {
let qpe: &QueryParseError = qpe;
eprintln!("logdive: query error: {qpe}");
} else {
eprintln!("logdive: {e}");
}
}
fn init_tracing() {
let filter = EnvFilter::try_from_env("LOGDIVE_LOG").unwrap_or_else(|_| EnvFilter::new("warn"));
tracing_subscriber::fmt()
.with_env_filter(filter)
.with_writer(io::stderr)
.init();
}
fn run(cli: Cli) -> Result<()> {
let db = db_path(cli.db.as_deref());
match cli.command {
Command::Ingest(args) => {
let mut indexer = Indexer::open(&db)?;
run_ingest(&mut indexer, args)
}
Command::Query(args) => {
let indexer = Indexer::open(&db)?;
run_query(&indexer, args)
}
Command::Stats(args) => run_stats(&db, args),
}
}
#[derive(Default, Debug)]
struct IngestReport {
inserted: usize,
deduplicated: usize,
skipped_no_timestamp: usize,
malformed: usize,
io_errors: usize,
elapsed: Duration,
}
impl IngestReport {
fn fold_insert_stats(&mut self, s: InsertStats) {
self.inserted += s.inserted;
self.deduplicated += s.deduplicated;
self.skipped_no_timestamp += s.skipped_no_timestamp;
}
fn total_seen(&self) -> usize {
self.inserted
+ self.deduplicated
+ self.skipped_no_timestamp
+ self.malformed
+ self.io_errors
}
}
fn run_ingest(indexer: &mut Indexer, args: IngestArgs) -> Result<()> {
let tag = args.tag;
let reader: Box<dyn BufRead> = match args.file {
Some(ref p) => {
let f = File::open(p).map_err(|e| LogdiveError::io_at(p.clone(), e))?;
Box::new(BufReader::new(f))
}
None => Box::new(BufReader::new(io::stdin().lock())),
};
let tty = io::stderr().is_terminal();
let mut progress = Progress::new(tty);
let mut report = IngestReport::default();
let started = Instant::now();
let mut batch: Vec<LogEntry> = Vec::with_capacity(BATCH_SIZE);
for line_result in reader.lines() {
let line = match line_result {
Ok(l) => l,
Err(_) => {
report.io_errors += 1;
continue;
}
};
match parse_line(&line) {
Some(entry) => batch.push(entry.with_tag(tag.clone())),
None if line.trim().is_empty() => {
}
None => report.malformed += 1,
}
if batch.len() >= BATCH_SIZE {
let stats = indexer.insert_batch(&batch)?;
report.fold_insert_stats(stats);
batch.clear();
progress.tick(&report, started.elapsed());
}
}
if !batch.is_empty() {
let stats = indexer.insert_batch(&batch)?;
report.fold_insert_stats(stats);
}
report.elapsed = started.elapsed();
progress.finish(&report);
print_summary(&report);
Ok(())
}
fn run_query(indexer: &Indexer, args: QueryArgs) -> Result<()> {
let ast = parse_query(&args.query)?;
let limit = if args.limit == 0 {
None
} else {
Some(args.limit)
};
tracing::debug!(?limit, "executing query");
let rows = execute(&ast, indexer.connection(), limit)?;
tracing::debug!(result_count = rows.len(), "query returned results");
render(&rows, args.format)
}
struct Progress {
tty: bool,
last_tick: Instant,
tick_interval: Duration,
}
impl Progress {
fn new(tty: bool) -> Self {
Self {
tty,
last_tick: Instant::now() - Duration::from_secs(3600),
tick_interval: Duration::from_millis(250),
}
}
fn tick(&mut self, report: &IngestReport, elapsed: Duration) {
if self.last_tick.elapsed() < self.tick_interval {
return;
}
self.last_tick = Instant::now();
self.render(report, elapsed);
}
fn finish(&mut self, report: &IngestReport) {
self.render(report, report.elapsed);
if self.tty {
eprintln!();
}
}
fn render(&self, report: &IngestReport, elapsed: Duration) {
let rate = lines_per_sec(report.total_seen(), elapsed);
let payload = format!(
"ingesting: {total:>7} seen | {ins:>7} new | {dedup:>5} dup | {bad:>5} skip | {rate:>7.0} lines/s",
total = report.total_seen(),
ins = report.inserted,
dedup = report.deduplicated,
bad = report.malformed + report.skipped_no_timestamp + report.io_errors,
rate = rate,
);
if self.tty {
let mut err = io::stderr().lock();
let _ = write!(err, "\r{payload}");
let _ = err.flush();
} else {
eprintln!("{payload}");
}
}
}
fn lines_per_sec(n: usize, elapsed: Duration) -> f64 {
let secs = elapsed.as_secs_f64();
if secs <= 0.0 { 0.0 } else { n as f64 / secs }
}
fn print_summary(report: &IngestReport) {
let rate = lines_per_sec(report.total_seen(), report.elapsed);
eprintln!(
"ingest complete in {:.2}s ({:.0} lines/s)",
report.elapsed.as_secs_f64(),
rate
);
eprintln!(" inserted: {}", report.inserted);
eprintln!(" deduplicated: {}", report.deduplicated);
eprintln!(" no timestamp: {}", report.skipped_no_timestamp);
eprintln!(" malformed: {}", report.malformed);
if report.io_errors > 0 {
eprintln!(" io errors: {}", report.io_errors);
}
}