Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ ctor = "0.6"
eyre = "0.6"
gungraun = "0.17"
io-uring = "0.7"
jiff = { version = "0.2", default-features = false, features = ["std"] }
libc = "0.2"
neli = "0.7"
pretty_assertions = "1"
Expand Down
2 changes: 2 additions & 0 deletions candumpr/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ color-eyre.workspace = true
crossbeam-channel.workspace = true
eyre.workspace = true
io-uring.workspace = true
jiff.workspace = true
libc.workspace = true
neli.workspace = true
tracing.workspace = true
Expand All @@ -26,6 +27,7 @@ ctor.workspace = true
gungraun.workspace = true
pretty_assertions.workspace = true
tabled.workspace = true
tempfile.workspace = true
vcan-fixture = { path = "../vcan-fixture" }

[[example]]
Expand Down
36 changes: 0 additions & 36 deletions candumpr/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,42 +8,6 @@ pub mod recv;
pub mod sink;
pub mod writer;

#[cfg(test)]
pub(crate) mod test_util {
pub(crate) struct TestBufWriter {
pub(crate) bytes: Vec<u8>,
}

impl TestBufWriter {
pub(crate) fn new() -> Self {
Self { bytes: Vec::new() }
}
}

impl crate::writer::Writer for TestBufWriter {
fn write(&mut self, b: &[u8]) -> eyre::Result<()> {
self.bytes.extend_from_slice(b);
Ok(())
}

fn flush(&mut self) -> eyre::Result<()> {
Ok(())
}

fn sync(&mut self) -> eyre::Result<()> {
Ok(())
}

fn finish(&mut self) -> eyre::Result<()> {
Ok(())
}

fn as_any_mut(&mut self) -> &mut dyn std::any::Any {
self
}
}
}

#[cfg(test)]
#[ctor::ctor]
fn test_setup() {
Expand Down
61 changes: 50 additions & 11 deletions candumpr/src/main.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
use std::os::unix::io::AsFd;
use std::path::PathBuf;
use std::process::ExitCode;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
Expand All @@ -11,8 +12,7 @@ use candumpr::frame::CanFrame;
use candumpr::pipeline::Pipeline;
use candumpr::recv::netlink::{self, LinkEvent};
use candumpr::recv::receiver::{BATCH_CAPACITY, Receiver};
use candumpr::sink::Sink;
use candumpr::writer::StdoutWriter;
use candumpr::sink::{Output, Sink, SinkConfig};
use clap::Parser;
use crossbeam_channel::select;

Expand Down Expand Up @@ -95,6 +95,16 @@ enum Format {
CandumpConsole,
}

impl Format {
/// Log file extension for each output format
fn ext(&self) -> &'static str {
match self {
Format::CandumpFile => "log",
Format::CandumpConsole => "txt",
}
}
}

/// Log CAN traffic from multiple networks.
#[derive(Parser)]
#[command(version)]
Expand All @@ -103,6 +113,16 @@ struct Cli {
#[arg(required = true)]
interfaces: Vec<String>,

/// Log to a file in the current directory, instead of stdout.
///
/// TODO: Accepts exactly one interface (for now).
#[arg(long, short = 'l', conflicts_with = "output")]
log: bool,

/// Log to this file path. Truncated if it already exists.
#[arg(long, short = 'o', value_name = "FILE")]
output: Option<PathBuf>,

/// Output format for received frames.
#[arg(long, value_enum, default_value = "candump-file")]
format: Format,
Expand All @@ -128,6 +148,11 @@ fn main() -> ExitCode {
.with_max_level(cli.log_level)
.init();

// TODO: Add multiple interface support for --log
if cli.log && cli.interfaces.len() > 1 {
unimplemented!("--log does not yet support multiple interfaces");
}

// The sockets vector defines the canonical interface ordering. The orderings of:
//
// 1. cli.interfaces
Expand Down Expand Up @@ -198,19 +223,29 @@ fn main() -> ExitCode {
Box::new(CanutilsConsoleFormatter::new(cli.interfaces, cli.timestamp))
}
};
let header = formatter.header().map(|h| h.to_vec());
let sink = Sink::new(
StdoutWriter::new(),
header,
64 * 1024,
Some(Duration::from_secs(5)),
Some(Duration::from_secs(5 * 60)),
);
let output = if cli.log {
Output::Template {
dir: ".".into(),
interface: names[0].clone(),
ext: cli.format.ext().to_string(),
next_index: 0,
}
} else if let Some(path) = cli.output {
Output::Path(path)
} else {
Output::Stdout
};
let mut sink_config = SinkConfig::new(output);
sink_config.header = formatter.header().map(|h| h.to_vec());
let sink = Sink::new(sink_config);
let mut pipeline = Pipeline::new(formatter, vec![sink]);

// Write-path errors are logged and recorded rather than propagated: returning early would skip
// draining the remaining batches and closing the pipeline, both of which can lose buffered data.
// Every error sets `failed` so the process still exits nonzero.
//
// A write_batch error means the Sink gave up on the output, so the loop breaks instead of
// continuing to write into something broken.
let mut failed = false;

// Debounce link state and bus state events.
Expand All @@ -230,6 +265,7 @@ fn main() -> ExitCode {
}
tracing::error!(error = ?e, "failed to write batch");
failed = true;
break;
}
batch.clear();
let _ = empty_tx.try_send(batch);
Expand Down Expand Up @@ -290,14 +326,17 @@ fn main() -> ExitCode {
}
}

// Drain everything the receiver queued before it exited.
// Drain everything the receiver queued before it exited. A failure here is the same
// unrecoverable class as above, so stop rather than retry the same broken output once per
// queued batch.
while let Ok(mut batch) = full_rx.try_recv() {
if let Err(e) = pipeline.write_batch(&batch) {
if is_broken_pipe(&e) {
break;
}
tracing::error!(error = ?e, "failed to write batch during drain");
failed = true;
break;
}
batch.clear();
let _ = empty_tx.try_send(batch);
Expand Down
44 changes: 21 additions & 23 deletions candumpr/src/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -106,13 +106,13 @@ impl Pipeline {
#[cfg(test)]
mod tests {
use pretty_assertions::assert_eq;
use tempfile::TempDir;

use super::*;
use crate::can::LinuxCanFrame;
use crate::format::{CanutilsFileFormatter, TimestampMode};
use crate::frame::Direction;
use crate::sink::Sink;
use crate::test_util::TestBufWriter;
use crate::sink::{Output, Sink, SinkConfig};

fn frame(sock_id: usize, id: u32, data: &[u8]) -> CanFrame {
CanFrame {
Expand All @@ -126,17 +126,17 @@ mod tests {
}
}

fn sink() -> Sink {
Sink::new(TestBufWriter::new(), None, 64 * 1024, None, None)
/// A Path-output Sink writing to `sink<i>.log` in `dir`, with time-based flush/sync disabled.
fn sink(dir: &TempDir, i: usize) -> Sink {
let mut config = SinkConfig::new(Output::Path(dir.path().join(format!("sink{i}.log"))));
config.flush_interval = None;
config.sync_interval = None;
Sink::new(config)
}

fn bytes_in(sink: &mut Sink) -> Vec<u8> {
sink.writer
.as_any_mut()
.downcast_mut::<TestBufWriter>()
.unwrap()
.bytes
.clone()
/// Contents of sink `i`'s file. The pipeline must be flushed first.
fn bytes_in(dir: &TempDir, i: usize) -> Vec<u8> {
std::fs::read(dir.path().join(format!("sink{i}.log"))).unwrap()
}

fn formatted(names: &[String], frames: &[&CanFrame]) -> Vec<u8> {
Expand All @@ -156,25 +156,28 @@ mod tests {
"can2".to_string(),
"can3".to_string(),
];
let dir = TempDir::new().unwrap();
let mut pipeline = Pipeline::new(
Box::new(CanutilsFileFormatter::new(
names.clone(),
TimestampMode::Absolute,
)),
vec![sink()],
vec![sink(&dir, 0)],
);

let frames = vec![frame(0, 0x100, &[0x01]), frame(3, 0x200, &[0x02])];
pipeline.write_batch(&frames).unwrap();
pipeline.flush().unwrap();

let expected = formatted(&names, &[&frames[0], &frames[1]]);
assert_eq!(bytes_in(&mut pipeline.sinks[0]), expected);
assert_eq!(bytes_in(&dir, 0), expected);
}

#[test]
fn per_interface_dispatches_by_sock_id() {
let names = vec!["can0".to_string(), "can1".to_string(), "can2".to_string()];
let sinks = vec![sink(), sink(), sink()];
let dir = TempDir::new().unwrap();
let sinks = vec![sink(&dir, 0), sink(&dir, 1), sink(&dir, 2)];
let mut pipeline = Pipeline::new(
Box::new(CanutilsFileFormatter::new(
names.clone(),
Expand All @@ -190,18 +193,13 @@ mod tests {
frame(1, 0x400, &[0x0D]),
];
pipeline.write_batch(&frames).unwrap();
pipeline.flush().unwrap();

assert_eq!(
bytes_in(&mut pipeline.sinks[0]),
bytes_in(&dir, 0),
formatted(&names, &[&frames[0], &frames[2]])
);
assert_eq!(
bytes_in(&mut pipeline.sinks[1]),
formatted(&names, &[&frames[3]])
);
assert_eq!(
bytes_in(&mut pipeline.sinks[2]),
formatted(&names, &[&frames[1]])
);
assert_eq!(bytes_in(&dir, 1), formatted(&names, &[&frames[3]]));
assert_eq!(bytes_in(&dir, 2), formatted(&names, &[&frames[1]]));
}
}
Loading