diff --git a/.github/workflows/lint.yml b/.github/workflows/lint.yml index 38da336..262d82a 100644 --- a/.github/workflows/lint.yml +++ b/.github/workflows/lint.yml @@ -85,6 +85,7 @@ jobs: # The vcan-fixture crate needs user namespaces to create isolated vcan # interfaces without root. sudo sysctl -w kernel.apparmor_restrict_unprivileged_userns=0 2>/dev/null || true + sudo apt-get update sudo apt-get install -y linux-modules-extra-"$(uname -r)" && \ sudo modprobe vcan && \ echo "available=true" >> "$GITHUB_OUTPUT" || true @@ -150,5 +151,6 @@ jobs: - name: Setup vcan run: | sudo sysctl -w kernel.apparmor_restrict_unprivileged_userns=0 2>/dev/null || true + sudo apt-get update sudo apt-get install -y linux-modules-extra-"$(uname -r)" sudo modprobe vcan diff --git a/CHANGELOG.md b/CHANGELOG.md index 8d5ffa0..8ec04c6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,21 @@ focus on the user impact** rather than the actual changes made. ## Fixed +# candemonium - 0.1.1 - (2026-08-10) + +Productionizability polish on the 0.1.0 MVP release. + +* Add a `[interfaces.].name` config field for human-meaningful interface names. Shows up in + logs and the filename if present. +* Rx/Tx direction is captured and added to `candump-file` format +* `[receiver].batch_size` can now be set +* Various DEBUG logging fixes. Error frames are now logged at TRACE level, with future logging + enhancements planned as a part of and + . + +The addition of the new config fields _should_ make this a SemVer 0.2.0 release, but candemonium has +no consumers yet ;) + # candemonium - 0.1.0 - (2026-08-07) First release of candemonium with an MVP release of the `candumpr` logging tool. diff --git a/Cargo.lock b/Cargo.lock index 054a02e..366d7a8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -150,7 +150,7 @@ checksum = "7354288c522e7e980fafd2075d63d1285794c3a6a16cdd492f189ea406e5f18b" [[package]] name = "candumpr" -version = "0.1.0" +version = "0.1.1" dependencies = [ "bytesize", "clap", @@ -177,7 +177,7 @@ dependencies = [ [[package]] name = "cangenr" -version = "0.1.0" +version = "0.1.1" [[package]] name = "cc" @@ -1260,7 +1260,7 @@ checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" [[package]] name = "vcan-fixture" -version = "0.1.0" +version = "0.1.1" dependencies = [ "assert_cmd", "ctor", diff --git a/Cargo.toml b/Cargo.toml index 9f21a12..8b34b7f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,7 +7,7 @@ members = [ ] [workspace.package] -version = "0.1.0" +version = "0.1.1" edition = "2024" license = "MIT" rust-version = "1.89" diff --git a/candumpr/src/config.rs b/candumpr/src/config.rs index 26278a3..735ac3c 100644 --- a/candumpr/src/config.rs +++ b/candumpr/src/config.rs @@ -7,6 +7,7 @@ use eyre::WrapErr; use crate::format::TimestampMode; use crate::quantity::Quantity; +use crate::recv::receiver::BATCH_CAPACITY; use crate::sink::{DEFAULT_FLUSH_INTERVAL, DEFAULT_SYNC_INTERVAL, Output}; /// The first interface name that appears more than once. @@ -90,6 +91,19 @@ pub struct Cli { #[arg(long, conflicts_with = "daemon")] pub no_request_address_claims: bool, + /// Number of received frames to wait for before each receiver wakeup. + /// + /// Larger batches reduce CPU overhead at high frame rates at the cost of recv latency, up to + /// 100ms. + #[arg( + long, + value_name = "N", + default_value = "4", + value_parser = clap::value_parser!(u16).range(1..=BATCH_CAPACITY as i64), + conflicts_with = "daemon" + )] + pub batch_size: u16, + /// Log level for tracing output on stderr. #[arg(long, default_value = "INFO")] pub log_level: tracing::Level, @@ -99,6 +113,8 @@ pub struct Cli { #[derive(Debug, PartialEq, Eq)] pub struct InterfaceConfig { pub name: String, + /// Optional user-given name, as written in the config file, to include in log messages + pub display_name: Option, pub request_address_claims: bool, } @@ -115,6 +131,8 @@ pub struct Config { pub streams: Vec, /// Whether recoverable [Sink](crate::sink::Sink) activation failures are retried or are fatal. pub retry_activation_failures: bool, + /// Number of received frames to wait for before each receiver wakeup. + pub batch_size: usize, } /// Configuration for one output stream @@ -221,6 +239,8 @@ struct RawStreamConfig { format: Format, compress: bool, timestamp: TimestampMode, + // Only allowed in [interface.] sections + name: Option, // This is the one TOML setting that's required; the rest have default values taken from // [RawStreamConfig::default]. directory: Option, @@ -241,6 +261,7 @@ impl Default for RawStreamConfig { format: Format::CandumpFile, compress: true, timestamp: TimestampMode::Absolute, + name: None, directory: None, flush_every: Interval::Every(DEFAULT_FLUSH_INTERVAL), sync_every: Interval::Every(DEFAULT_SYNC_INTERVAL), @@ -251,6 +272,18 @@ impl Default for RawStreamConfig { } } +#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Deserialize)] +#[serde(default)] +struct RawReceiverConfig { + batch_size: usize, +} + +impl Default for RawReceiverConfig { + fn default() -> Self { + RawReceiverConfig { batch_size: 32 } + } +} + /// The TOML config file settings /// /// This struct exists to facilitate the TOML parsing. But the [Config] struct is what we end up @@ -264,6 +297,25 @@ struct Raw { defaults: RawStreamConfig, #[serde(default)] interface: HashMap, + #[serde(default)] + receiver: RawReceiverConfig, +} + +/// Lowercase `name` and collapse every run of characters outside [a-z0-9] to a single hyphen. +fn sluggify(name: &str) -> String { + let mut slug = String::with_capacity(name.len()); + for ch in name.chars() { + let ch = ch.to_ascii_lowercase(); + if ch.is_ascii_alphanumeric() { + slug.push(ch); + } else if !slug.is_empty() && !slug.ends_with('-') { + slug.push('-'); + } + } + if slug.ends_with('-') { + slug.pop(); + } + slug } /// Layer `over` onto `base`: every key in `over` replaces the same key in `base`. @@ -325,6 +377,7 @@ impl Config { .map(|interface| Output::Template { dir: ".".into(), interface: interface.clone(), + name: None, ext: cli.format.ext(cli.compress).to_string(), }) .collect() @@ -355,6 +408,7 @@ impl Config { .iter() .map(|name| InterfaceConfig { name: name.clone(), + display_name: None, request_address_claims: !cli.no_request_address_claims, }) .collect(); @@ -363,6 +417,7 @@ impl Config { interfaces, streams, retry_activation_failures: false, + batch_size: cli.batch_size as usize, }) } @@ -385,11 +440,12 @@ impl Config { let mut names: Vec = raw.interface.keys().cloned().collect(); names.sort_unstable(); - drop(raw); // just parsed for type-checking and valid TOML syntax. Overlaying is all done on the toml::Table level below. + let batch_size = raw.receiver.batch_size; eyre::ensure!( - !names.is_empty(), - "no [interface.] sections; at least one is required" + (1..=BATCH_CAPACITY).contains(&batch_size), + "receiver batch_size must be between 1 and {BATCH_CAPACITY}, got {batch_size}" ); + drop(raw); // just parsed for type-checking and valid TOML syntax. Overlaying is all done on the toml::Table level below. // Second pass: Overlay each interface's table over the [defaults] table and deserialize the // result, so absent keys fall through to Settings::default(). @@ -399,6 +455,10 @@ impl Config { .and_then(toml::Value::as_table) .cloned() .unwrap_or_default(); + eyre::ensure!( + !defaults.contains_key("name"), + "`name` is not allowed in [defaults]; set it in each [interface.] section" + ); let sections = table .get("interface") .and_then(toml::Value::as_table) @@ -423,6 +483,17 @@ impl Config { "interface {interface}: missing `directory` setting in [interface.{interface}] or [defaults]" ); }; + let name = match &settings.name { + Some(raw) => { + let slug = sluggify(raw); + eyre::ensure!( + !slug.is_empty(), + "interface {interface}: name {raw:?} has no usable characters after sluggification" + ); + Some(slug) + } + None => None, + }; // A retention limit at or below the rotation limit can never be met. Only same-kind // limits are comparable and therefore validatable; mixed kinds are best-effort. match (settings.retain, settings.rotate_every) { @@ -440,12 +511,14 @@ impl Config { } interfaces.push(InterfaceConfig { name: interface.clone(), + display_name: settings.name.clone(), request_address_claims: settings.request_address_claims, }); streams.push(StreamConfig { output: Output::Template { dir: directory.join(interface), interface: interface.clone(), + name, ext: settings.format.ext(settings.compress).to_string(), }, format: settings.format, @@ -462,6 +535,7 @@ impl Config { interfaces, streams, retry_activation_failures: true, + batch_size, }) } } @@ -490,6 +564,7 @@ mod tests { flush_every = "250ms" rotate_every = "off" retain = "10 files" + name = "Engine Bus (J1939)" "#; let config = Config::from_toml(src).unwrap(); @@ -498,15 +573,18 @@ mod tests { [ InterfaceConfig { name: "can0".to_string(), + display_name: Some("Engine Bus (J1939)".to_string()), request_address_claims: true, }, InterfaceConfig { name: "can1".to_string(), + display_name: None, request_address_claims: true, }, ] ); assert!(config.retry_activation_failures); + assert_eq!(config.batch_size, 32); assert_eq!( config.streams, [ @@ -514,6 +592,7 @@ mod tests { output: Output::Template { dir: "/var/log/can/can0".into(), interface: "can0".to_string(), + name: Some("engine-bus-j1939".to_string()), ext: "txt.zst".to_string(), }, format: Format::CandumpConsole, @@ -529,6 +608,7 @@ mod tests { output: Output::Template { dir: "/var/log/can/can1".into(), interface: "can1".to_string(), + name: None, ext: "log".to_string(), }, format: Format::CandumpFile, @@ -646,15 +726,6 @@ mod tests { "got: {err}" ); - let err = format!( - "{:#}", - Config::from_toml("[defaults]\ndirectory = \"/x\"\n").unwrap_err() - ); - assert!(err.contains("at least one is required"), "got: {err}"); - - let err = format!("{:#}", Config::from_toml("[interface]\n").unwrap_err()); - assert!(err.contains("at least one is required"), "got: {err}"); - let err = format!( "{:#}", Config::from_toml( @@ -672,5 +743,33 @@ mod tests { "[defaults]\ndirectory = \"/x\"\nrotate_every = \"200KB\"\nretain = \"1 day\"\n[interface.can0]\n" ) .unwrap(); + + let err = format!( + "{:#}", + Config::from_toml( + "[defaults]\ndirectory = \"/x\"\nname = \"shared\"\n[interface.can0]\n" + ) + .unwrap_err() + ); + assert!(err.contains("not allowed in [defaults]"), "got: {err}"); + + let err = format!( + "{:#}", + Config::from_toml("[defaults]\ndirectory = \"/x\"\n[interface.can0]\nname = \"!!!\"\n") + .unwrap_err() + ); + assert!(err.contains("no usable characters"), "got: {err}"); + + let err = format!( + "{:#}", + Config::from_toml( + "[defaults]\ndirectory = \"/x\"\n[interface.can0]\n[receiver]\nbatch_size = 0\n" + ) + .unwrap_err() + ); + assert!( + err.contains("batch_size must be between 1 and 256"), + "got: {err}" + ); } } diff --git a/candumpr/src/format.rs b/candumpr/src/format.rs index d423dbf..7ac9f54 100644 --- a/candumpr/src/format.rs +++ b/candumpr/src/format.rs @@ -2,7 +2,7 @@ use std::io::Write; use std::time::Duration; use crate::errframe::ErrorFrame; -use crate::frame::CanFrame; +use crate::frame::{CanFrame, Direction}; use crate::recv::Timestamp; /// How candumpr renders the timestamp prefix on candump-format output. @@ -120,7 +120,7 @@ fn canid_and_width(can_id: u32) -> (u32, usize) { /// Formats frames in the can-utils candump file (log) format. /// -/// Output format: `(TIMESTAMP) IFACE CANID#DATA\n`, e.g. `(1616161616.123456) can0 18FECA00#AABB0011`. +/// Output format: `(TIMESTAMP) IFACE CANID#DATA DIR\n`, e.g. `(1616161616.123456) can0 18FECA00#AABB0011 R`. /// The timestamp is rendered according to the configured [TimestampMode]. pub struct CanutilsFileFormatter { iface_names: Vec, @@ -147,7 +147,10 @@ impl Formatter for CanutilsFileFormatter { for i in 0..frame.raw.len as usize { write!(buf, "{:02X}", frame.raw.data[i]).unwrap(); } - buf.push(b'\n'); + buf.extend_from_slice(match frame.direction { + Direction::Tx => b" T\n", + Direction::Rx => b" R\n", + }); } } @@ -230,7 +233,7 @@ mod tests { ); assert_eq!( String::from_utf8(buf).unwrap(), - "(1616161616.000123) can0 18FECA00#AABB0011\n" + "(1616161616.000123) can0 18FECA00#AABB0011 R\n" ); } @@ -247,7 +250,7 @@ mod tests { ); assert_eq!( String::from_utf8(buf).unwrap(), - "(0.000000) vcan1 00000123#\n" + "(0.000000) vcan1 00000123# R\n" ); } @@ -261,7 +264,7 @@ mod tests { ); assert_eq!( String::from_utf8(buf).unwrap(), - "(1.000000) can0 20000040#0000000000000000\n" + "(1.000000) can0 20000040#0000000000000000 R\n" ); } @@ -269,8 +272,13 @@ mod tests { fn canutils_format_standard_frame_uses_three_digits() { let mut fmt = CanutilsFileFormatter::new(vec!["can0".to_string()], TimestampMode::Absolute); let mut buf = Vec::new(); - fmt.format(&frame(0, ts(0, 0), 0x123, &[0xAB]), &mut buf); - assert_eq!(String::from_utf8(buf).unwrap(), "(0.000000) can0 123#AB\n"); + let mut msg = frame(0, ts(0, 0), 0x123, &[0xAB]); + msg.direction = Direction::Tx; + fmt.format(&msg, &mut buf); + assert_eq!( + String::from_utf8(buf).unwrap(), + "(0.000000) can0 123#AB T\n" + ); } #[test] @@ -281,7 +289,7 @@ mod tests { fmt.format(&frame(0, ts(100, 250_000_000), 0x123, &[0xAB]), &mut buf); assert_eq!( String::from_utf8(buf).unwrap(), - "(000.000000) can0 123#AB\n(000.250000) can0 123#AB\n" + "(000.000000) can0 123#AB R\n(000.250000) can0 123#AB R\n" ); } @@ -293,7 +301,7 @@ mod tests { fmt.format(&frame(0, ts(105, 0), 0x123, &[0xAB]), &mut buf); assert_eq!( String::from_utf8(buf).unwrap(), - "(000.000000) can0 123#AB\n(005.000000) can0 123#AB\n" + "(000.000000) can0 123#AB R\n(005.000000) can0 123#AB R\n" ); } diff --git a/candumpr/src/main.rs b/candumpr/src/main.rs index 2dcc4f3..1d673c6 100644 --- a/candumpr/src/main.rs +++ b/candumpr/src/main.rs @@ -41,7 +41,12 @@ fn is_broken_pipe(err: &eyre::Report) -> bool { /// Log a link-state edge to stderr, ignoring repeats of the last observed state. /// /// Returns true if the interface transitioned to down. -fn handle_link_event(event: LinkEvent, link_up: &mut [Option], names: &[String]) -> bool { +fn handle_link_event( + event: LinkEvent, + link_up: &mut [Option], + names: &[String], + display_names: &[Option], +) -> bool { let (sock_id, up) = match event { LinkEvent::LinkUp { sock_id } => (sock_id, true), LinkEvent::LinkDown { sock_id } => (sock_id, false), @@ -51,10 +56,11 @@ fn handle_link_event(event: LinkEvent, link_up: &mut [Option], names: &[St } link_up[sock_id] = Some(up); let interface = &names[sock_id]; + let name = display_names[sock_id].as_deref(); if up { - tracing::info!(interface = %interface, "interface link up"); + tracing::info!(interface = %interface, name, "interface link up"); } else { - tracing::warn!(interface = %interface, "interface link down"); + tracing::warn!(interface = %interface, name, "interface link down"); } !up } @@ -62,14 +68,20 @@ fn handle_link_event(event: LinkEvent, link_up: &mut [Option], names: &[St /// Log each error frame in `batch` at debug level, and log bus-state transitions (edges only). /// /// Returns true if any interface transitioned into bus-off. -fn log_error_frames(batch: &[CanFrame], bus_state: &mut [BusState], names: &[String]) -> bool { +fn log_error_frames( + batch: &[CanFrame], + bus_state: &mut [BusState], + names: &[String], + display_names: &[Option], +) -> bool { let mut bus_off = false; for frame in batch { let Some(err) = ErrorFrame::parse(&frame.raw) else { continue; }; let interface = &names[frame.sock_id]; - tracing::debug!(interface = %interface, error = %err, "CAN error frame"); + let name = display_names[frame.sock_id].as_deref(); + tracing::trace!(interface = %interface, name, error = %err, "CAN error frame"); let Some(new) = err.bus_state() else { continue; @@ -81,13 +93,13 @@ fn log_error_frames(batch: &[CanFrame], bus_state: &mut [BusState], names: &[Str bus_state[frame.sock_id] = new; match new { BusState::ErrorActive => { - tracing::info!(interface = %interface, "bus state {old} -> {new}") + tracing::info!(interface = %interface, name, "bus state {old} -> {new}") } BusState::ErrorWarning | BusState::ErrorPassive => { - tracing::warn!(interface = %interface, "bus state {old} -> {new}") + tracing::warn!(interface = %interface, name, "bus state {old} -> {new}") } BusState::BusOff => { - tracing::error!(interface = %interface, "bus state {old} -> {new}"); + tracing::error!(interface = %interface, name, "bus state {old} -> {new}"); bus_off = true; } } @@ -108,6 +120,7 @@ fn main() -> ExitCode { .with_target("neli", tracing::Level::WARN); tracing_subscriber::fmt() .with_writer(std::io::stderr) + .with_max_level(cli.log_level) .finish() .with(filter) .init(); @@ -120,6 +133,11 @@ fn main() -> ExitCode { } }; + if config.interfaces.is_empty() { + tracing::info!("no interfaces configured; nothing to do"); + return ExitCode::SUCCESS; + } + // The sockets vector defines the canonical interface ordering. The orderings of: // // 1. config.interfaces @@ -130,6 +148,11 @@ fn main() -> ExitCode { // // all follow the same ordering, and are all indexed by CanFrame::sock_id let names: Vec = config.interfaces.iter().map(|i| i.name.clone()).collect(); + let display_names: Vec> = config + .interfaces + .iter() + .map(|i| i.display_name.clone()) + .collect(); let sockets: Vec<_> = match config .interfaces .iter() @@ -143,12 +166,16 @@ fn main() -> ExitCode { } }; - for (name, sock) in names.iter().zip(&sockets) { + for (i, sock) in sockets.iter().enumerate() { + let interface = &names[i]; + let name = display_names[i].as_deref(); match can::get_recv_buffer(sock.as_fd()) { Ok(bytes) => { - tracing::info!(interface = %name, rcvbuf_bytes = bytes, "opened CAN socket") + tracing::info!(interface = %interface, name, rcvbuf_bytes = bytes, "opened CAN socket") + } + Err(e) => { + tracing::warn!(interface = %interface, name, error = ?e, "failed to query rcvbuf size") } - Err(e) => tracing::warn!(interface = %name, error = ?e, "failed to query rcvbuf size"), } } @@ -176,16 +203,31 @@ fn main() -> ExitCode { .collect(); let recv_names = names.clone(); + let recv_display_names = display_names.clone(); let recv_flags = claim_flags.clone(); - let recv_handle = std::thread::spawn(move || -> eyre::Result { - let mut recv = Receiver::new(sockets)?; - let total = recv.run(&SIGNAL_STOP, &recv_names, &recv_flags, &full_tx, &empty_rx)?; - Ok(total) - }); + let batch_size = config.batch_size; + let recv_handle = std::thread::Builder::new() + .name("candumpr-recv".into()) + .spawn(move || -> eyre::Result { + let mut recv = Receiver::new(sockets, batch_size)?; + let total = recv.run( + &SIGNAL_STOP, + &recv_names, + &recv_display_names, + &recv_flags, + &full_tx, + &empty_rx, + )?; + Ok(total) + }) + .expect("failed to spawn recv thread"); let (event_tx, event_rx) = crossbeam_channel::unbounded::(); let nl_names = names.clone(); - let nl_handle = std::thread::spawn(move || netlink::run(&SIGNAL_STOP, &nl_names, &event_tx)); + let nl_handle = std::thread::Builder::new() + .name("candumpr-net".into()) + .spawn(move || netlink::run(&SIGNAL_STOP, &nl_names, &event_tx)) + .expect("failed to spawn netlink thread"); // Last observed link state per sock_id, so we log only edges. let mut link_up: Vec> = vec![None; names.len()]; @@ -247,7 +289,7 @@ fn main() -> ExitCode { select! { recv(full_rx) -> msg => match msg { Ok(mut batch) => { - if log_error_frames(&batch, &mut bus_state, &names) { + if log_error_frames(&batch, &mut bus_state, &names, &display_names) { state_debounce.trigger(Instant::now()); } if let Err(e) = pipeline.write_batch(&batch) { @@ -266,7 +308,7 @@ fn main() -> ExitCode { }, recv(event_rx) -> msg => match msg { Ok(event) => { - if handle_link_event(event, &mut link_up, &names) { + if handle_link_event(event, &mut link_up, &names, &display_names) { state_debounce.trigger(Instant::now()); } } @@ -343,7 +385,7 @@ fn main() -> ExitCode { // Log any link transitions the netlink thread queued before it exited. while let Ok(event) = event_rx.try_recv() { - handle_link_event(event, &mut link_up, &names); + handle_link_event(event, &mut link_up, &names, &display_names); } // close() always runs, even after write errors: for file and zstd writers it is what writes the diff --git a/candumpr/src/recv/receiver.rs b/candumpr/src/recv/receiver.rs index 03d8c9a..0207a7a 100644 --- a/candumpr/src/recv/receiver.rs +++ b/candumpr/src/recv/receiver.rs @@ -20,9 +20,6 @@ pub const BATCH_CAPACITY: usize = FRAMEBUF_COUNT as usize; const BGID: u16 = 0; -/// Number of CQEs to wait for before waking. -const BATCH_SIZE: usize = 4; - /// Size of the `io_uring_recvmsg_out` header the kernel writes at the start of each provided buffer const RECVMSG_OUT_HDR: usize = 16; @@ -39,6 +36,8 @@ const BUF_ENTRY_SIZE: usize = RECVMSG_OUT_HDR + CMSG_BUF_SIZE + FRAME_SIZE; pub struct Receiver { ring: IoUring, sockets: Vec, + /// Number of CQEs to wait for before waking. + batch_size: usize, framebuf_ring_ptr: *mut u8, framebuf_ring_layout: Layout, framebuf_data: Box<[u8]>, @@ -52,7 +51,10 @@ impl Receiver { /// [open_can_raw](crate::can::open_can_raw)). /// /// Enables hardware timestamping, drop count reporting, and own-message echo on each socket. - pub fn new(sockets: Vec) -> std::io::Result { + /// + /// `batch_size` is the number of CQEs to wait for before waking. Larger batches reduce wakeup + /// frequency at the cost of added recv latency, bounded by the 100ms timeout in [Self::run]. + pub fn new(sockets: Vec, batch_size: usize) -> std::io::Result { for sock in &sockets { can::enable_timestamps(sock.as_fd())?; can::enable_drop_count(sock.as_fd())?; @@ -110,6 +112,7 @@ impl Receiver { Ok(Self { ring, sockets, + batch_size, framebuf_ring_ptr, framebuf_ring_layout, framebuf_data, @@ -126,12 +129,13 @@ impl Receiver { &mut self, stop: &AtomicBool, names: &[String], + display_names: &[Option], claim_flags: &[Arc], full_tx: &crossbeam_channel::Sender>, empty_rx: &crossbeam_channel::Receiver>, ) -> std::io::Result { let mut total = 0u64; - // We wait for BATCH_SIZE completions before waking up, but we still want to be able to + // We wait for batch_size completions before waking up, but we still want to be able to // react to the stop signal in a timely manner, so we set a 100ms timeout. let timeout = types::Timespec::new().nsec(100_000_000); let args = types::SubmitArgs::new().timespec(&timeout); @@ -152,7 +156,11 @@ impl Receiver { } while !stop.load(Ordering::Relaxed) { - match self.ring.submitter().submit_with_args(BATCH_SIZE, &args) { + match self + .ring + .submitter() + .submit_with_args(self.batch_size, &args) + { Ok(_) => {} Err(e) if e.raw_os_error() == Some(libc::ETIME) => {} Err(e) if e.raw_os_error() == Some(libc::EINTR) => continue, @@ -201,10 +209,16 @@ impl Receiver { let meta = parse_control_data(out.control_data()); let timestamp = meta.timestamp.unwrap_or(Timestamp { sec: 0, nsec: 0 }); + let direction = if out.flags() & libc::MSG_DONTROUTE as u32 != 0 { + Direction::Tx + } else { + Direction::Rx + }; + batch.push(CanFrame { sock_id: idx, timestamp, - direction: Direction::Rx, + direction, raw, }); total += 1; @@ -236,11 +250,14 @@ impl Receiver { if flag.swap(false, Ordering::Relaxed) { let frame = can::address_claim_pgn_request(); match can::send_frame(self.sockets[idx].as_fd(), &frame) { - Ok(()) => { - tracing::debug!(interface = %names[idx], "sent address claim PGN request") - } + Ok(()) => tracing::debug!( + interface = %names[idx], + name = display_names[idx].as_deref(), + "sent address claim PGN request" + ), Err(e) => tracing::error!( interface = %names[idx], + name = display_names[idx].as_deref(), error = %e, "failed to send address claim request" ), @@ -382,8 +399,8 @@ mod tests { // Receiver must be created on the same thread that calls run() due to SINGLE_ISSUER. let handle = std::thread::spawn(move || { - let mut recv = Receiver::new(rx_sockets).unwrap(); - recv.run(&STOP, &[], &[], &full_tx, &empty_rx) + let mut recv = Receiver::new(rx_sockets, 4).unwrap(); + recv.run(&STOP, &[], &[], &[], &full_tx, &empty_rx) }); let mut count = 0u64; @@ -393,7 +410,8 @@ mod tests { .unwrap(); for frame in &batch { assert!(frame.sock_id < IFACE_COUNT); - assert_eq!(frame.direction, Direction::Rx); + // Sent by another socket on this host, so it's a local TX as far as the kernel sees + assert_eq!(frame.direction, Direction::Tx); assert!(frame.raw.len <= 8); count += 1; } diff --git a/candumpr/src/sink.rs b/candumpr/src/sink.rs index b948556..cb1c1d4 100644 --- a/candumpr/src/sink.rs +++ b/candumpr/src/sink.rs @@ -20,6 +20,8 @@ pub enum Output { Template { dir: PathBuf, interface: String, + /// Optional user-given stream name, already sluggified, to include in filenames + name: Option, /// File extension to use ext: String, }, @@ -73,10 +75,12 @@ impl SinkConfig { Output::Template { dir, interface, + name, ext, } => dir.join(template::render( template::next_index_in(dir, interface), interface, + name.as_deref(), timestamp.sec, ext, )), @@ -444,6 +448,7 @@ mod tests { let mut config = SinkConfig::new(Output::Template { dir: blocker.join("logs"), interface: "can0".to_string(), + name: None, ext: "log".to_string(), }); config.flush_every = Interval::Off; @@ -458,6 +463,7 @@ mod tests { let mut config = SinkConfig::new(Output::Template { dir: dir.path().to_path_buf(), interface: "can0".to_string(), + name: None, ext: "log".to_string(), }); config.header = header; @@ -648,6 +654,7 @@ mod tests { let mut config = SinkConfig::new(Output::Template { dir: dir.path().to_path_buf(), interface: "can0".to_string(), + name: None, ext: "log".to_string(), }); // Every byte lands on disk immediately, so the tests can read files mid-stream. diff --git a/candumpr/src/template.rs b/candumpr/src/template.rs index 2877d6c..acd4f9e 100644 --- a/candumpr/src/template.rs +++ b/candumpr/src/template.rs @@ -15,8 +15,12 @@ fn iso_utc(sec: i64) -> String { } /// Render a filename given various parameters. -pub fn render(index: u64, interface: &str, sec: i64, ext: &str) -> String { - format!("i{index:04}_{interface}_{}.{ext}", iso_utc(sec)) +pub fn render(index: u64, interface: &str, name: Option<&str>, sec: i64, ext: &str) -> String { + let ts = iso_utc(sec); + match name { + Some(name) => format!("i{index:04}_{interface}_{name}_{ts}.{ext}"), + None => format!("i{index:04}_{interface}_{ts}.{ext}"), + } } /// Scan the given directory for files matching candumpr's filename template generated by [render] @@ -108,21 +112,28 @@ mod tests { #[test] fn render_filename() { assert_eq!( - render(0, "can0", 1732117385, "log"), + render(0, "can0", None, 1732117385, "log"), "i0000_can0_2024-11-20T15-43-05Z.log" ); + // A name slots in between the interface and the timestamp. + assert_eq!( + render(0, "can0", Some("engine"), 1732117385, "log"), + "i0000_can0_engine_2024-11-20T15-43-05Z.log" + ); // The index pads to width 4 and keeps counting past 9999 as plain decimal. assert_eq!( - render(10000, "vcan1", 0, "pcap"), + render(10000, "vcan1", None, 0, "pcap"), "i10000_vcan1_1970-01-01T00-00-00Z.pcap" ); } #[test] fn parse_index_reads_back_what_render_writes() { - for index in [0, 7, 10000, u64::MAX] { - let name = render(index, "can0", 1732117385, "log"); - assert_eq!(parse_index(&name, "can0"), Some(index), "name={name}"); + for stream_name in [None, Some("engine")] { + for index in [0, 7, 10000, u64::MAX] { + let name = render(index, "can0", stream_name, 1732117385, "log"); + assert_eq!(parse_index(&name, "can0"), Some(index), "name={name}"); + } } } diff --git a/candumpr/tests/basic.rs b/candumpr/tests/basic.rs index af347b8..c19e0b4 100644 --- a/candumpr/tests/basic.rs +++ b/candumpr/tests/basic.rs @@ -57,7 +57,7 @@ fn canutils_stdout_output() { let stdout = String::from_utf8(output.stdout).unwrap(); // Strip the dynamic timestamp prefix from each line. The format is: - // (SECONDS.MICROSECONDS) IFACE CANID#DATA + // (SECONDS.MICROSECONDS) IFACE CANID#DATA DIR // Everything after ") " is deterministic. let lines: Vec<&str> = stdout .lines() @@ -65,6 +65,6 @@ fn canutils_stdout_output() { .collect(); assert_eq!(lines.len(), 2); - assert_eq!(lines[0], format!("{iface0} 18FECA00#AABB")); - assert_eq!(lines[1], format!("{iface1} 00000123#010203")); + assert_eq!(lines[0], format!("{iface0} 18FECA00#AABB T")); + assert_eq!(lines[1], format!("{iface1} 00000123#010203 T")); } diff --git a/candumpr/tests/compress.rs b/candumpr/tests/compress.rs index 938f2c4..4d4ec0c 100644 --- a/candumpr/tests/compress.rs +++ b/candumpr/tests/compress.rs @@ -92,7 +92,7 @@ fn logs_compressed_to_a_zst_file() { assert!(status.success(), "zstd -d exited {status}"); assert_eq!( String::from_utf8(stdout).unwrap(), - format!("(000.000000) {iface} 123#AB\n") + format!("(000.000000) {iface} 123#AB T\n") ); } @@ -172,7 +172,7 @@ fn sigkill_leaves_a_decodable_prefix() { &text[text.len().saturating_sub(60)..] ); for (i, line) in text.lines().enumerate() { - let counter = line.rsplit('#').next().unwrap(); + let counter = line.rsplit('#').next().unwrap().strip_suffix(" T").unwrap(); assert_eq!( u32::from_str_radix(counter, 16).unwrap(), i as u32, diff --git a/candumpr/tests/daemon.rs b/candumpr/tests/daemon.rs index 71000ca..57def83 100644 --- a/candumpr/tests/daemon.rs +++ b/candumpr/tests/daemon.rs @@ -66,6 +66,7 @@ fn logs_each_interface_to_its_own_subdirectory() { [interface.{iface1}] compress = true + name = "Engine Bus (J1939)" [interface.{iface2}] "#, @@ -98,7 +99,7 @@ fn logs_each_interface_to_its_own_subdirectory() { ); let name = paths[0].file_name().unwrap().to_str().unwrap(); - let prefix = format!("i0000_{iface1}_"); + let prefix = format!("i0000_{iface1}_engine-bus-j1939_"); assert_eq!(paths[0].parent().unwrap(), log_dir.path().join(iface1)); assert!( name.starts_with(&prefix) && name.ends_with(".log.zst"), @@ -115,7 +116,7 @@ fn logs_each_interface_to_its_own_subdirectory() { assert_eq!( std::fs::read_to_string(&paths[1]).unwrap(), - format!("(000.000000) {iface2} 456#CD\n") + format!("(000.000000) {iface2} 456#CD T\n") ); } diff --git a/candumpr/tests/file_output.rs b/candumpr/tests/file_output.rs index 6936af8..f24b61a 100644 --- a/candumpr/tests/file_output.rs +++ b/candumpr/tests/file_output.rs @@ -85,7 +85,7 @@ fn logs_one_interface_to_a_scheme_named_file() { // Deterministic: --timestamp zero renders the first frame at 0.0, and only one frame is sent. assert_eq!( std::fs::read_to_string(&paths[0]).unwrap(), - format!("(000.000000) {iface} 123#AB\n") + format!("(000.000000) {iface} 123#AB T\n") ); assert_eq!(String::from_utf8_lossy(&output.stdout), ""); let stderr = String::from_utf8_lossy(&output.stderr); @@ -138,8 +138,8 @@ fn interleaves_multiple_interfaces_to_one_file() { let name = paths[0].file_name().unwrap().to_str().unwrap(); assert_eq!(name, "foo.bar"); let log = std::fs::read_to_string(&paths[0]).unwrap(); - assert!(log.contains(&format!("{iface1} 123#AB\n"))); - assert!(log.contains(&format!("{iface2} 456#CD\n"))); + assert!(log.contains(&format!("{iface1} 123#AB T\n"))); + assert!(log.contains(&format!("{iface2} 456#CD T\n"))); } #[test] @@ -193,7 +193,7 @@ fn logs_each_interface_to_its_own_file() { .unwrap_or_else(|| panic!("expected a {prefix}*.log file, got {paths:?}")); assert_eq!( std::fs::read_to_string(path).unwrap(), - format!("(000.000000) {iface} {frame}\n") + format!("(000.000000) {iface} {frame} T\n") ); } } diff --git a/docs/user/candumpr-configuration.md b/docs/user/candumpr-configuration.md index e46db3d..1bec7ee 100644 --- a/docs/user/candumpr-configuration.md +++ b/docs/user/candumpr-configuration.md @@ -50,6 +50,7 @@ Each of the `[defaults]` and `[interface.]` tables support the following k | `rotate_every` | `"30min"` | Rotate the log once it exceeds a size or an age | | `retain` | `"1 GB"` | Delete the oldest log files once the interface directory exceeds a total size, an age, or a file count | | `request_address_claims` | `true` | Broadcast a J1939 Address Claim PGN request on the interface whenever one of its log files is opened | +| `name` | none | Human-meaningful stream name added to log filenames. Not allowed in `[defaults]` | Durations are parsed using [jiff's friendly format](https://docs.rs/jiff/latest/jiff/fmt/friendly/index.html). Days and weeks @@ -60,6 +61,22 @@ The `rotate_every`, `flush_every`, and `sync_every` parameters accept durations The `retain` parameter additionally accepts file counts, like `"10 files"`. +## Receiver table + +Receiver thread settings can be set in the `[receiver]` table: + +```toml +[receiver] +batch_size = 32 +``` + +| Key | Default | Description | +| ------------ | ------- | ------------------------------------------------------------------------------------ | +| `batch_size` | `32` | Number of received frames to wait for before each receiver wakeup, between 1 and 256 | + +Larger batches reduce CPU overhead at high frame rates at the cost of recv latency, up to a 100ms +timeout. In CLI mode this is set with `--batch-size`. + ## Output files Each interface logs to its own subdirectory using files with the following format: