#![warn(rust_2018_compatibility)]
#![warn(rust_2018_idioms)]
#![warn(rust_2021_compatibility)]
#![deny(warnings)]
#[macro_use]
extern crate log;
#[macro_use]
extern crate serde_derive;
use std::{
collections::BTreeMap,
convert::identity,
fs::File,
io::{self, BufWriter, Write},
path::Path,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
mpsc::RecvTimeoutError,
},
thread,
time::Duration,
};
use bus::{Bus, BusReader};
use clap::CommandFactory;
use color_eyre::Result;
use csv::Writer;
use env_logger::{Target, WriteStyle};
use log::LevelFilter;
use opts::{Commands, Generate};
use crate::record::Record;
mod file_parser;
mod flags;
mod opts;
mod record;
mod uniques;
mod version;
use mimalloc::MiMalloc;
#[global_allocator]
static GLOBAL: MiMalloc = MiMalloc;
fn main() -> Result<()> {
match opts::get_opts()?.command {
Commands::Dump(d) => dump(d),
Commands::Generate(g) => generate(g),
#[cfg(feature = "watch")]
Commands::Watch(w) => watch(w),
}
}
fn csv_write<I>(recv: BusReader<Arc<Record>>, mut writer: Writer<I>, _: bool, flush_all: bool)
where
I: Write,
{
for rec in recv {
if let Err(err) = writer.serialize(rec) {
error!("Couldn't serialize csv: {err}");
}
if flush_all && let Err(err) = writer.flush() {
error!("Couldn't flush csv: {err}");
}
}
}
fn json_write<I>(recv: BusReader<Arc<Record>>, mut writer: I, pretty: bool, flush_all: bool)
where
I: Write,
{
if pretty {
for rec in recv {
if let Err(err) = serde_json::to_writer_pretty(&mut writer, &rec) {
error!("Couldn't serialize json: {err}");
}
if let Err(err) = writeln!(writer) {
error!("Couldn't append json newline: {err}");
}
if flush_all && let Err(err) = writer.flush() {
error!("Couldn't flush json: {err}");
}
}
} else {
for rec in recv {
if let Err(err) = serde_json::to_writer(&mut writer, &rec) {
error!("Couldn't serialize json: {err}");
}
if let Err(err) = writeln!(writer) {
error!("Couldn't append json newline: {err}");
}
if flush_all && let Err(err) = writer.flush() {
error!("Couldn't flush json: {err}");
}
}
}
}
fn yaml_write<I>(recv: BusReader<Arc<Record>>, mut writer: I, _: bool, flush_all: bool)
where
I: Write,
{
for rec in recv {
if let Err(err) = writeln!(writer, "---") {
error!("Couldn't write yaml separator: {err}");
}
if let Err(err) = serde_yaml::to_writer(&mut writer, &rec) {
error!("Couldn't serialize yaml: {err}");
}
if let Err(err) = writeln!(writer) {
error!("Couldn't append yaml newline: {err}");
}
if flush_all && let Err(err) = writer.flush() {
error!("Couldn't flush yaml: {err}");
}
}
}
fn write_uniqs<I>(
recv: BusReader<Arc<Record>>,
mut writer: Writer<I>,
_: bool,
include_timestamps: bool,
) where
I: Write,
{
let mut u = BTreeMap::new();
for rec in recv {
u.entry(rec.path.clone())
.or_insert_with(uniques::UniqueCounts::default)
.update(rec.flag, rec.file_timestamp);
}
if include_timestamps {
for (path, v) in u {
if let Err(err) = writer.serialize(v.into_unique_out(path)) {
error!("Error writing the uniques: {err}");
}
}
} else {
#[cfg(feature = "alt_flags")]
let header = vec!["path", "counts", "flags", "alt_flags"];
#[cfg(not(feature = "alt_flags"))]
let header = vec!["path", "counts", "flags"];
if let Err(err) = writer.write_record(&header) {
error!("Error writing CSV header: {err}");
return;
}
for (path, v) in u {
let out = v.into_unique_out_no_timestamps(path);
let counts_str = out.counts.to_string();
#[cfg(feature = "alt_flags")]
let record = vec![
out.path.as_str(),
counts_str.as_str(),
out.flags,
out.alt_flags,
];
#[cfg(not(feature = "alt_flags"))]
let record = vec![out.path.as_str(), counts_str.as_str(), out.flags];
if let Err(err) = writer.write_record(&record) {
error!("Error writing unique record: {err}");
}
}
}
}
fn path_stdout(p: &Path) -> bool {
p.as_os_str() == "-"
}
#[inline]
fn icsv(rec: Arc<Record>, writer: &mut Writer<BufWriter<File>>) {
if let Err(err) = writer.serialize(&rec) {
error!("Error writing csv rec: {err}")
}
}
#[inline]
fn ijson(rec: Arc<Record>, writer: &mut BufWriter<File>) {
if let Err(err) = serde_json::to_writer(&mut *writer, &rec) {
error!("Error writing json rec: {err}")
}
if let Err(err) = writeln!(writer) {
error!("Error writing json newline: {err}")
}
}
#[inline]
fn iyaml(rec: Arc<Record>, writer: &mut BufWriter<File>) {
if let Err(err) = writeln!(writer, "---") {
error!("Error writing yaml separator: {err}")
}
if let Err(err) = serde_yaml::to_writer(&mut *writer, &rec) {
error!("Error writing yaml rec: {err}")
}
if let Err(err) = writeln!(writer) {
error!("Error writing yaml newline: {err}")
}
}
macro_rules! fdump {
( $bus: ident, $scope: ident, $ftype: expr, $path:ident, $proc_f:ident, $c_opt: ident, $creater:expr, ) => {
if let Some(p) = $path {
let recv = $bus.add_rx();
if path_stdout(&p) {
$scope.spawn(move |_| {
$proc_f(recv, $creater($c_opt.make_stdout()), false, false);
});
} else {
match File::create(&p) {
Err(err) => error!(
"Couldn't create {} output file {}: {err}",
$ftype,
p.display()
),
Ok(f) => {
$scope.spawn(move |_| {
if $c_opt.is_gz(&p) {
$proc_f(
recv,
$creater($c_opt.make_gzip(BufWriter::new(f))),
false,
false,
);
} else if $c_opt.is_zstd(&p) {
#[cfg(feature = "zstd")]
{
$proc_f(recv, $creater($c_opt.make_zstd(f)), false, false);
}
#[cfg(not(feature = "zstd"))]
unreachable!("zstd feature not enabled");
} else {
$proc_f(recv, $creater(BufWriter::new(f)), false, false);
};
});
}
}
}
};
};
}
macro_rules! idump {
( $want: ident, $bus: ident, $fscope: ident, $running: ident, $ftype: expr, $f: ident, $make_out: expr, $ifun: expr, ) => {
if $want {
let mut out_path = $f.clone();
out_path.as_mut_os_string().push(format!(".{}", $ftype));
match File::create(&out_path) {
Err(err) => error!(
"Couldn't open a {} writer at {}: {err}",
$ftype,
out_path.display()
),
Ok(w) => {
let mut recv = $bus.add_rx();
let running = $running.clone();
$fscope.spawn(move |_| {
let out = &mut $make_out(BufWriter::new(w));
'RUNNING: loop {
match recv.recv_timeout(Duration::from_millis(50)) {
Ok(r) => $ifun(r, out),
Err(e) => match e {
RecvTimeoutError::Timeout => {
if !running.load(Ordering::Acquire) {
break 'RUNNING;
}
thread::yield_now();
}
_ => return,
},
}
}
});
}
};
};
};
}
#[inline]
fn new_bus() -> Bus<Arc<Record>> {
Bus::new(4096)
}
fn dump(opts: opts::Dump) -> Result<()> {
let std_counts = opts.stdout_counts();
env_logger::Builder::new()
.filter(
None,
if std_counts == 1 {
LevelFilter::Error
} else {
LevelFilter::Info
},
)
.write_style(WriteStyle::Always)
.target(Target::Stderr)
.init();
color_eyre::install()?;
opts.validate(std_counts)?;
let file_paths = opts.real_files();
info!("Starting");
let opts::Dump {
csvs: individual_csvs,
jsons: individual_jsons,
yamls: individual_yamls,
csv: csv_path,
json: json_path,
yaml: yaml_path,
uniques: uniq_path,
unique_timestamps,
..
} = opts;
let rec_filter = opts.filter_opts.filter()?;
let copts = opts.compress_opts;
crossbeam::scope(|scope| {
let mut bus = new_bus();
fdump!(
bus,
scope,
"csv",
csv_path,
csv_write,
copts,
csv::Writer::from_writer,
);
if let Some(p) = uniq_path {
let recv = bus.add_rx();
if path_stdout(&p) {
scope.spawn(move |_| {
write_uniqs(
recv,
csv::Writer::from_writer(copts.make_stdout()),
false,
unique_timestamps,
);
});
} else {
match File::create(&p) {
Err(err) => error!(
"Couldn't create unique csv output file {}: {err}",
p.display()
),
Ok(f) => {
scope.spawn(move |_| {
if copts.is_gz(&p) {
write_uniqs(
recv,
csv::Writer::from_writer(copts.make_gzip(BufWriter::new(f))),
false,
unique_timestamps,
);
} else if copts.is_zstd(&p) {
#[cfg(feature = "zstd")]
{
write_uniqs(
recv,
csv::Writer::from_writer(copts.make_zstd(f)),
false,
unique_timestamps,
);
}
#[cfg(not(feature = "zstd"))]
unreachable!("zstd feature not enabled");
} else {
write_uniqs(
recv,
csv::Writer::from_writer(BufWriter::new(f)),
false,
unique_timestamps,
);
};
});
}
}
}
}
fdump!(bus, scope, "json", json_path, json_write, copts, identity,);
fdump!(bus, scope, "yaml", yaml_path, yaml_write, copts, identity,);
for f in file_paths {
let running = Arc::new(AtomicBool::new(true));
crossbeam::scope(|fscope| {
idump!(
individual_csvs,
bus,
fscope,
running,
"csv",
f,
Writer::from_writer,
icsv,
);
idump!(
individual_jsons,
bus,
fscope,
running,
"json",
f,
identity,
ijson,
);
idump!(
individual_yamls,
bus,
fscope,
running,
"yaml",
f,
identity,
iyaml,
);
match file_parser::parse_file(&f, &mut bus, &rec_filter) {
Ok(_) => info!("Finished parsing {}", f.display()),
Err(e) => error!("Couldn't parse '{}': {}", f.display(), e),
};
running.store(false, Ordering::Release);
})
.expect("Couldn't close all the threads");
}
})
.expect("Couldn't close all the threads");
Ok(())
}
fn generate(g: Generate) -> Result<()> {
let mut cmd = opts::Cli::command();
let name = cmd.get_name().to_string();
clap_complete::generate(g.shell, &mut cmd, name, &mut io::stdout().lock());
Ok(())
}
#[cfg(feature = "watch")]
fn watch(opts: opts::Watch) -> Result<()> {
use std::mem;
use notify_debouncer_full::{
DebounceEventResult, FileIdMap, new_debouncer_opt, notify::RecursiveMode,
};
use crate::file_parser::parse_file;
env_logger::Builder::new()
.filter(None, LevelFilter::Info)
.write_style(WriteStyle::Always)
.target(Target::Stderr)
.init();
color_eyre::install()?;
let rec_filter = opts.filter_opts.filter()?;
let (send, recv) = crossbeam_channel::bounded(128);
let debounce_time = Duration::from_secs(2);
if opts.poll {
let mut debouncer = new_debouncer_opt::<_, notify::PollWatcher, FileIdMap>(
debounce_time,
None,
move |result: DebounceEventResult| match result {
Ok(events) => events.iter().for_each(|event| {
if event.kind.is_create() {
for path in event.paths.iter() {
if path.exists()
&& let Err(err) =
send.send_timeout(path.clone(), Duration::from_secs(1))
{
error!("Error processing created file {}: {err}", path.display());
}
}
}
}),
Err(errors) => errors
.iter()
.for_each(|error| error!("Watch error: {error:?}")),
},
FileIdMap::new(),
notify::Config::default().with_poll_interval(Duration::from_secs(2)),
)?;
for path in opts.watch_dirs {
info!("Watching {}", path.display());
debouncer.watch(&path, RecursiveMode::Recursive)?;
}
mem::forget(debouncer);
} else {
let mut debouncer = new_debouncer_opt::<_, notify::RecommendedWatcher, FileIdMap>(
debounce_time,
None,
move |result: DebounceEventResult| match result {
Ok(events) => events.iter().for_each(|event| {
if event.kind.is_create() {
for path in event.paths.iter() {
if path.exists()
&& let Err(err) =
send.send_timeout(path.clone(), Duration::from_secs(1))
{
error!("Error processing created file {}: {err}", path.display());
}
}
}
}),
Err(errors) => errors
.iter()
.for_each(|error| error!("Watch error: {error:?}")),
},
FileIdMap::new(),
notify::Config::default(),
)?;
for path in opts.watch_dirs {
info!("Watching {}", path.display());
debouncer.watch(&path, RecursiveMode::Recursive)?;
}
mem::forget(debouncer);
};
let copts = opts.compress_opts;
crossbeam::scope(|fscope| {
let mut bus = new_bus();
let rec_recv = bus.add_rx();
fscope.spawn(move |_| {
let out = copts.make_stdout();
match opts.format {
opts::WatchFormat::Csv => {
csv_write(rec_recv, csv::Writer::from_writer(out), false, true)
}
opts::WatchFormat::Json => json_write(rec_recv, out, opts.pretty, true),
opts::WatchFormat::Yaml => yaml_write(rec_recv, out, false, true),
}
});
for path in recv {
if let Err(err) = parse_file(&path, &mut bus, &rec_filter) {
error!("Error parsing {}: {err}", path.display());
}
}
})
.unwrap();
Ok(())
}