1use std::collections::{BTreeMap, HashMap};
16use std::path::{Path, PathBuf};
17use std::sync::{Arc, Mutex};
18
19use color_eyre::Result;
20use color_eyre::eyre::eyre;
21
22use crate::error_display::FileError;
23use polars::prelude::*;
24
25use crate::formats::columns::{Cell, Kind};
26use crate::formats::dbc::{Dbc, Message, Mux, Signal};
27use crate::formats::fixed_records::{Bytes, ColumnLayout, Logical, Physical};
28use crate::formats::indexed::Offsets;
29use crate::formats::model_files::MetaValue;
30use crate::formats::sqlite::Table;
31use crate::formats::text_formats::Detail;
32
33pub(crate) const READER: crate::formats::readers::Reader = crate::formats::readers::Reader {
35 scan,
36 signatures: &[crate::formats::readers::Signature {
37 says: |head, _| looks_like(head),
38 kind: crate::formats::readers::Kind::Text,
39 trusted: crate::formats::readers::Trusted {
40 tables: true,
41 ..crate::formats::readers::EVERYWHERE
42 },
43 }],
44 tables: Some(|path| {
45 listed(path).ok_or_else(|| color_eyre::eyre::eyre!("Open the log to list its tables."))
46 }),
47 ..crate::formats::readers::BASE
48};
49
50const MAX_LINE: usize = 4096;
52const MAX_INTERFACES: usize = 255;
54const MAX_LONG_PARTS: usize = 10_000;
56
57pub const FRAMES: &str = "frames";
59pub const SIGNALS: &str = "signals";
60
61#[derive(Debug, Clone, PartialEq)]
63pub struct Frame<'a> {
64 pub ts: Option<i64>,
66 pub iface: &'a str,
67 pub id: u32,
68 pub extended: bool,
69 pub fd: bool,
70 pub flags: Option<u8>,
72 pub remote: bool,
73 pub error: bool,
74 pub dlc: u8,
75 pub data: Vec<u8>,
76}
77
78fn hex_bytes(text: &str) -> Option<Vec<u8>> {
79 if !text.len().is_multiple_of(2) || text.len() > 128 {
80 return None;
81 }
82 (0..text.len())
83 .step_by(2)
84 .map(|i| u8::from_str_radix(text.get(i..i + 2)?, 16).ok())
85 .collect()
86}
87
88fn parse_id(text: &str) -> Option<(u32, bool)> {
90 if !text.bytes().all(|b| b.is_ascii_hexdigit()) {
91 return None;
92 }
93 match text.len() {
94 1..=3 => Some((u32::from_str_radix(text, 16).ok()?, false)),
95 8 => Some((u32::from_str_radix(text, 16).ok()?, true)),
96 _ => None,
97 }
98}
99
100fn parse_ts(text: &str) -> Option<i64> {
102 if let Some((s, frac)) = text.split_once('.')
103 && !s.is_empty()
104 && s.bytes().all(|b| b.is_ascii_digit())
105 && !frac.is_empty()
106 && frac.bytes().all(|b| b.is_ascii_digit())
107 {
108 let secs: i64 = s.parse().ok()?;
109 let digits = &frac[..frac.len().min(6)];
110 let micros: i64 = digits.parse::<i64>().ok()? * 10i64.pow(6 - digits.len() as u32);
111 return secs.checked_mul(1_000_000)?.checked_add(micros);
112 }
113 let parsed = chrono::NaiveDateTime::parse_from_str(text, "%Y-%m-%d %H:%M:%S%.f").ok()?;
115 Some(parsed.and_utc().timestamp_micros())
116}
117
118pub fn parse_line(line: &str) -> Option<Frame<'_>> {
120 let line = line.trim();
121 if line.is_empty() || line.len() > MAX_LINE {
122 return None;
123 }
124 let (ts, rest) = match line.strip_prefix('(') {
125 Some(after) => {
126 let (inside, rest) = after.split_once(')')?;
127 (Some(parse_ts(inside.trim())?), rest.trim_start())
128 }
129 None => (None, line),
130 };
131 let mut words = rest.split_whitespace();
132 let iface = words.next()?;
133 let next = words.next()?;
134 if let Some((id, frame)) = next.split_once('#') {
135 return parse_compact(ts, iface, id, frame);
136 }
137 let mut id_word = next;
139 let mut dlc_word = words.next()?;
140 let mut guard = 0;
141 while !(dlc_word.starts_with('[') && dlc_word.ends_with(']')) {
142 id_word = dlc_word;
143 dlc_word = words.next()?;
144 guard += 1;
145 if guard > 4 {
146 return None;
147 }
148 }
149 let (id, extended) = parse_id(id_word)?;
150 let dlc: u8 = dlc_word[1..dlc_word.len() - 1].parse().ok()?;
151 let rest: Vec<&str> = words.collect();
152 if rest.first() == Some(&"remote") {
153 return Some(Frame {
154 ts,
155 iface,
156 id,
157 extended,
158 fd: false,
159 flags: None,
160 remote: true,
161 error: false,
162 dlc,
163 data: Vec::new(),
164 });
165 }
166 let data: Vec<u8> = rest
167 .iter()
168 .take(dlc as usize)
169 .map(|b| {
170 (b.len() == 2)
171 .then(|| u8::from_str_radix(b, 16).ok())
172 .flatten()
173 })
174 .collect::<Option<_>>()?;
175 if data.len() != dlc as usize || dlc > 64 {
176 return None;
177 }
178 Some(Frame {
179 ts,
180 iface,
181 id,
182 extended,
183 fd: dlc > 8,
184 flags: None,
185 remote: false,
186 error: extended && id & 0x2000_0000 != 0,
187 dlc,
188 data,
189 })
190}
191
192fn parse_compact<'a>(ts: Option<i64>, iface: &'a str, id: &str, frame: &str) -> Option<Frame<'a>> {
193 let (id, extended) = parse_id(id)?;
194 let error = extended && id & 0x2000_0000 != 0;
195 let mut out = Frame {
196 ts,
197 iface,
198 id,
199 extended,
200 fd: false,
201 flags: None,
202 remote: false,
203 error,
204 dlc: 0,
205 data: Vec::new(),
206 };
207 if let Some(fd) = frame.strip_prefix('#') {
208 if fd.starts_with('#') {
210 return None;
211 }
212 let mut chars = fd.chars();
213 let flags = chars.next()?.to_digit(16)? as u8;
214 let data = hex_bytes(chars.as_str())?;
215 out.fd = true;
216 out.flags = Some(flags);
217 out.dlc = data.len() as u8;
218 out.data = data;
219 return Some(out);
220 }
221 if let Some(remote) = frame.strip_prefix('R').or_else(|| frame.strip_prefix('r')) {
222 out.remote = true;
223 out.dlc = if remote.is_empty() {
224 0
225 } else {
226 remote.parse().ok().filter(|d| *d <= 8)?
227 };
228 return Some(out);
229 }
230 let (data, code) = match frame.split_once('_') {
232 Some((data, code)) => (data, u8::from_str_radix(code, 16).ok()),
233 None => (frame, None),
234 };
235 let data = hex_bytes(data)?;
236 if data.len() > 8 {
237 return None;
238 }
239 out.dlc = code.unwrap_or(data.len() as u8);
240 out.data = data;
241 Some(out)
242}
243
244pub fn looks_like(head: &[u8]) -> bool {
246 let text = String::from_utf8_lossy(head);
247 let mut lines = text.lines().filter(|l| !l.trim().is_empty());
248 let Some(first) = lines.next() else {
249 return false;
250 };
251 let cut = !text.contains('\n');
254 parse_line(first).is_some()
255 || (cut
256 && first
257 .char_indices()
258 .last()
259 .is_some_and(|(at, _)| parse_line(&first[..at]).is_some()))
260}
261
262#[derive(Debug, Default)]
264pub struct Index {
265 pub offsets: Arc<Offsets>,
267 pub keys: Vec<u32>,
269 pub ifaces: Vec<u8>,
271 pub interfaces: Vec<String>,
272 pub absolute: bool,
274 pub skipped: usize,
276 pub past_limit: usize,
277}
278
279pub fn index(data: &[u8]) -> std::result::Result<Index, String> {
281 let limit = crate::limits::get().indexed_records;
282 let mut offsets = Offsets::for_file(data.len());
283 let mut index = Index::default();
284 let mut ifaces: HashMap<String, u8> = HashMap::new();
285 let mut first_ts = None;
286 let mut at = 0usize;
287 while at < data.len() {
288 let end = memchr::memchr(b'\n', &data[at..]).map_or(data.len(), |i| at + i);
289 let line = &data[at..end];
290 let parsed = (line.len() <= MAX_LINE)
291 .then(|| std::str::from_utf8(line).ok())
292 .flatten()
293 .and_then(parse_line);
294 match parsed {
295 Some(frame) if offsets.len() < limit => {
296 if first_ts.is_none() {
297 first_ts = frame.ts;
298 }
299 let iface = match ifaces.get(frame.iface) {
300 Some(&i) => i,
301 None if index.interfaces.len() < MAX_INTERFACES => {
302 let i = index.interfaces.len() as u8;
303 index.interfaces.push(frame.iface.to_string());
304 ifaces.insert(frame.iface.to_string(), i);
305 i
306 }
307 None => (MAX_INTERFACES - 1) as u8,
308 };
309 offsets.push(at);
310 index
311 .keys
312 .push(frame.id | (u32::from(frame.extended) << 31));
313 index.ifaces.push(iface);
314 }
315 Some(_) => index.past_limit += 1,
316 None if line.iter().all(|b| b.is_ascii_whitespace()) => {}
317 None => index.skipped += 1,
318 }
319 at = end + 1;
320 }
321 if offsets.is_empty() {
322 return Err("no line is a CAN frame as candump writes them".into());
323 }
324 offsets.shrink();
325 index.keys.shrink_to_fit();
326 index.ifaces.shrink_to_fit();
327 index.offsets = Arc::new(offsets);
328 index.absolute = first_ts.is_some_and(|ts| ts >= 978_307_200_000_000);
330 Ok(index)
331}
332
333fn line_of(bytes: &[u8], at: usize) -> &str {
335 let end = memchr::memchr(b'\n', &bytes[at..]).map_or(bytes.len(), |i| at + i);
336 std::str::from_utf8(&bytes[at..end.min(at + MAX_LINE)]).unwrap_or_default()
337}
338
339fn ts_kind(absolute: bool) -> Kind {
341 if absolute {
342 Kind::DatetimeUs
343 } else {
344 Kind::DurationUs
345 }
346}
347
348fn ts_cell(absolute: bool, ts: Option<i64>) -> Cell {
349 if absolute {
350 Cell::DatetimeUs(ts)
351 } else {
352 Cell::DurationUs(ts)
353 }
354}
355
356#[derive(Debug, Default)]
359struct Lines {
360 frames: Vec<Option<Parsed>>,
361 ifaces: Vec<String>,
363 data: Vec<u8>,
365}
366
367#[derive(Debug)]
369struct Parsed {
370 ts: Option<i64>,
371 iface: u32,
372 id: u32,
373 extended: bool,
374 fd: bool,
375 flags: Option<u8>,
376 remote: bool,
377 error: bool,
378 dlc: u8,
379 data: std::ops::Range<usize>,
380}
381
382impl Lines {
383 fn parse(bytes: &[u8], starts: impl ExactSizeIterator<Item = usize>) -> Self {
385 let mut lines = Lines {
386 frames: Vec::with_capacity(starts.len()),
387 ..Default::default()
388 };
389 for at in starts {
390 let parsed = parse_line(line_of(bytes, at)).map(|f| {
391 let iface = match lines.ifaces.iter().position(|i| i == f.iface) {
392 Some(i) => i,
393 None => {
394 lines.ifaces.push(f.iface.to_string());
395 lines.ifaces.len() - 1
396 }
397 } as u32;
398 let start = lines.data.len();
399 lines.data.extend_from_slice(&f.data);
400 Parsed {
401 ts: f.ts,
402 iface,
403 id: f.id,
404 extended: f.extended,
405 fd: f.fd,
406 flags: f.flags,
407 remote: f.remote,
408 error: f.error,
409 dlc: f.dlc,
410 data: start..lines.data.len(),
411 }
412 });
413 lines.frames.push(parsed);
414 }
415 lines
416 }
417
418 fn data(&self, frame: &Parsed) -> &[u8] {
419 &self.data[frame.data.clone()]
420 }
421}
422
423#[derive(Debug, Clone, Copy, PartialEq, Eq)]
426enum WindowKey {
427 Run {
428 first: IdxSize,
429 len: usize,
430 },
431 Rows {
432 first: IdxSize,
433 len: usize,
434 hash: u64,
435 },
436}
437
438impl WindowKey {
439 fn of(rows: &[IdxSize]) -> Self {
440 let first = rows.first().copied().unwrap_or(0);
441 let run = rows
442 .iter()
443 .enumerate()
444 .all(|(i, &r)| r as usize == first as usize + i);
445 if run {
446 return Self::Run {
447 first,
448 len: rows.len(),
449 };
450 }
451 use std::hash::{Hash, Hasher};
452 let mut hasher = std::collections::hash_map::DefaultHasher::new();
453 rows.hash(&mut hasher);
454 Self::Rows {
455 first,
456 len: rows.len(),
457 hash: hasher.finish(),
458 }
459 }
460}
461
462const KEPT_ROWS: usize = 1 << 20;
464
465struct Windows<T> {
471 width: usize,
473 slots: usize,
474 kept: Mutex<std::collections::VecDeque<Slot<T>>>,
475 #[cfg(test)]
476 parses: std::sync::atomic::AtomicUsize,
477}
478
479struct Slot<T> {
480 key: WindowKey,
481 listed: Option<Box<[IdxSize]>>,
484 rows: usize,
485 taken: usize,
486 parsed: Arc<std::sync::OnceLock<T>>,
487}
488
489impl<T> Windows<T> {
490 fn new(width: usize) -> Self {
491 let threads = std::thread::available_parallelism().map_or(4, |n| n.get());
492 Self {
493 width,
494 slots: (2 * threads).max(4),
495 kept: Default::default(),
496 #[cfg(test)]
497 parses: Default::default(),
498 }
499 }
500
501 fn with<R>(
503 &self,
504 rows: &[IdxSize],
505 parse: impl FnOnce(&[IdxSize]) -> T,
506 then: impl FnOnce(&T) -> R,
507 ) -> R {
508 let key = WindowKey::of(rows);
509 let parsed = {
510 let mut kept = self.kept.lock().unwrap_or_else(|e| e.into_inner());
511 let same = |s: &Slot<T>| {
512 s.key == key && s.listed.as_deref().is_none_or(|listed| listed == rows)
513 };
514 match kept.iter().position(same) {
515 Some(at) => {
516 let slot = &mut kept[at];
517 slot.taken += 1;
518 let parsed = slot.parsed.clone();
519 if slot.taken >= self.width {
520 kept.remove(at);
521 }
522 parsed
523 }
524 None => {
525 let parsed = Arc::new(std::sync::OnceLock::new());
526 if self.width > 1 {
527 kept.push_back(Slot {
528 key,
529 listed: matches!(key, WindowKey::Rows { .. }).then(|| rows.into()),
530 rows: rows.len(),
531 taken: 1,
532 parsed: parsed.clone(),
533 });
534 let mut total: usize = kept.iter().map(|s| s.rows).sum();
535 while kept.len() > 1 && (kept.len() > self.slots || total > KEPT_ROWS) {
536 total -= kept.pop_front().map_or(0, |s| s.rows);
537 }
538 }
539 parsed
540 }
541 }
542 };
543 then(parsed.get_or_init(|| {
544 #[cfg(test)]
545 self.parses
546 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
547 parse(rows)
548 }))
549 }
550}
551
552fn raw_columns(absolute: bool) -> [(&'static str, Kind); 9] {
554 [
555 ("ts", ts_kind(absolute)),
556 ("iface", Kind::Str),
557 ("id", Kind::Str),
558 ("ext", Kind::Bool),
559 ("dlc", Kind::U8),
560 ("data", Kind::Binary),
561 ("fd", Kind::Bool),
562 ("flags", Kind::U8),
563 ("kind", Kind::Label),
564 ]
565}
566
567pub struct RawFrames {
569 bytes: Arc<Bytes>,
570 offsets: Arc<Offsets>,
571 absolute: bool,
572 schema: SchemaRef,
573 windows: Windows<Lines>,
574}
575
576impl RawFrames {
577 pub fn new(bytes: Arc<Bytes>, index: &Index) -> Self {
578 let schema = raw_columns(index.absolute)
579 .into_iter()
580 .map(|(name, kind)| Field::new(name.into(), kind.dtype()))
581 .collect::<Schema>();
582 Self {
583 bytes,
584 offsets: index.offsets.clone(),
585 absolute: index.absolute,
586 windows: Windows::new(schema.len()),
587 schema: Arc::new(schema),
588 }
589 }
590
591 pub fn rows(&self) -> usize {
592 self.offsets.len().min(crate::formats::row_index::MAX_ROWS)
593 }
594
595 fn lines(&self, rows: impl ExactSizeIterator<Item = usize>) -> Lines {
596 Lines::parse(self.bytes.as_slice(), rows.map(|r| self.offsets.get(r)))
597 }
598
599 fn column(&self, column: usize, lines: &Lines) -> PolarsResult<Column> {
601 let (name, kind) = raw_columns(self.absolute)[column];
602 let cells = lines.frames.iter().map(|f| {
603 let f = f.as_ref();
604 match column {
605 0 => ts_cell(self.absolute, f.and_then(|f| f.ts)),
606 1 => Cell::Str(f.map(|f| lines.ifaces[f.iface as usize].clone())),
607 2 => Cell::Str(f.map(|f| {
608 if f.extended {
609 format!("{:08X}", f.id)
610 } else {
611 format!("{:03X}", f.id)
612 }
613 })),
614 3 => Cell::Bool(f.map(|f| f.extended)),
615 4 => Cell::U8(f.map(|f| f.dlc)),
616 5 => Cell::Binary(f.map(|f| lines.data(f).to_vec())),
617 6 => Cell::Bool(f.map(|f| f.fd)),
618 7 => Cell::U8(f.and_then(|f| f.flags)),
619 _ => Cell::Label(f.map(|f| {
620 if f.error {
621 "error"
622 } else if f.remote {
623 "remote"
624 } else {
625 "data"
626 }
627 })),
628 }
629 });
630 Ok(crate::formats::columns::series(name, kind, cells)?.into_column())
631 }
632
633 pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
634 self.bytes.still_whole()?;
635 let start = start.min(self.rows());
636 let len = len.min(self.rows() - start);
637 let lines = self.lines(start..start + len);
638 let columns = (0..self.schema.len())
639 .map(|c| self.column(c, &lines))
640 .collect::<PolarsResult<_>>()?;
641 DataFrame::new(len, columns)
642 }
643}
644
645impl crate::formats::row_index::RowSource for RawFrames {
646 fn height(&self) -> usize {
647 self.rows()
648 }
649
650 fn schema(&self) -> SchemaRef {
651 self.schema.clone()
652 }
653
654 fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
655 let rows = crate::formats::row_index::checked(index, self.rows())?;
656 self.bytes.still_whole()?;
657 self.windows.with(
658 &rows,
659 |rows| self.lines(rows.iter().map(|&r| r as usize)),
660 |lines| self.column(column, lines),
661 )
662 }
663}
664
665impl crate::formats::pushdown::Windowed for RawFrames {
666 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
667 Ok(self.collect_window(start, len)?.lazy())
668 }
669}
670
671pub struct Decoded {
673 bytes: Arc<Bytes>,
674 offsets: Arc<Offsets>,
675 rows: Arc<Vec<u32>>,
677 message: Arc<Message>,
678 multiplexer: Option<usize>,
680 absolute: bool,
681 named: bool,
683 schema: SchemaRef,
684 windows: Windows<Lines>,
685}
686
687fn signal_layout(signal: &Signal, named: bool) -> ColumnLayout {
690 let whole = signal.float == 0 && signal.factor == 1.0 && signal.offset.fract() == 0.0;
691 let physical = if signal.signed || signal.offset < 0.0 {
692 Physical::Signed(8)
693 } else {
694 Physical::Unsigned(8)
695 };
696 let mut layout = ColumnLayout::new(&signal.name, 0, 8, physical, 8);
697 layout.logical = if named && !signal.values.is_empty() && whole && signal.offset == 0.0 {
698 Logical::Enum(Arc::new(
699 signal.values.iter().map(|(k, v)| (*k, v.clone())).collect(),
700 ))
701 } else if whole && signal.offset == 0.0 {
702 Logical::Plain
703 } else {
704 Logical::Linear {
705 factor: signal.factor,
706 offset: signal.offset,
707 }
708 };
709 layout
710}
711
712impl Decoded {
713 pub fn new(
714 bytes: Arc<Bytes>,
715 index: &Index,
716 rows: Arc<Vec<u32>>,
717 message: Arc<Message>,
718 named: bool,
719 ) -> PolarsResult<Self> {
720 let mut fields = vec![Field::new("ts".into(), ts_kind(index.absolute).dtype())];
721 for s in &message.signals {
722 let dtype = if s.float != 0 {
723 DataType::Float64
724 } else {
725 signal_layout(s, named).dtype()
726 };
727 fields.push(Field::new(s.name.as_str().into(), dtype));
728 }
729 let schema: Schema = fields.into_iter().collect();
730 polars_ensure!(
731 schema.len() == message.signals.len() + 1,
732 Duplicate: "{} has two signals of one name, or one named ts", message.name
733 );
734 Ok(Self {
735 bytes,
736 offsets: index.offsets.clone(),
737 rows,
738 multiplexer: message
739 .signals
740 .iter()
741 .position(|s| s.mux == Mux::Multiplexer),
742 message,
743 absolute: index.absolute,
744 named,
745 windows: Windows::new(schema.len()),
746 schema: Arc::new(schema),
747 })
748 }
749
750 pub fn height(&self) -> usize {
751 self.rows.len().min(crate::formats::row_index::MAX_ROWS)
752 }
753
754 fn lines(&self, rows: impl ExactSizeIterator<Item = usize>) -> Lines {
755 Lines::parse(
756 self.bytes.as_slice(),
757 rows.map(|r| self.offsets.get(self.rows[r] as usize)),
758 )
759 }
760
761 fn column(&self, column: usize, lines: &Lines) -> PolarsResult<Column> {
763 let name = self.schema.get_at_index(column).map(|(n, _)| n.clone());
764 let name = name.unwrap_or_default();
765 let Some(signal) = column.checked_sub(1).map(|s| &self.message.signals[s]) else {
766 let ts = lines
767 .frames
768 .iter()
769 .map(|f| ts_cell(self.absolute, f.as_ref().and_then(|f| f.ts)));
770 let ts = crate::formats::columns::series(&name, ts_kind(self.absolute), ts)?;
771 return Ok(ts.into_column());
772 };
773 let multiplexer = self.multiplexer.map(|m| &self.message.signals[m]);
774 let data = |f: &Option<Parsed>| {
775 let data = lines.data(f.as_ref()?);
776 let mux = multiplexer.and_then(|m| crate::formats::dbc::raw(m, data));
777 crate::formats::dbc::present(signal, mux).then_some(data)
778 };
779 let series = if signal.float != 0 {
780 lines
781 .frames
782 .iter()
783 .map(|f| crate::formats::dbc::physical(signal, data(f)?))
784 .collect::<Float64Chunked>()
785 .into_series()
786 } else {
787 let ints: Vec<Option<i128>> = lines
788 .frames
789 .iter()
790 .map(|f| crate::formats::dbc::integer(signal, data(f)?))
791 .collect();
792 crate::formats::fixed_records::integers(&signal_layout(signal, self.named), ints)?
793 };
794 Ok(series.with_name(name).into_column())
795 }
796
797 pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
798 self.bytes.still_whole()?;
799 let start = start.min(self.height());
800 let len = len.min(self.height() - start);
801 let lines = self.lines(start..start + len);
802 let columns = (0..self.schema.len())
803 .map(|c| self.column(c, &lines))
804 .collect::<PolarsResult<_>>()?;
805 DataFrame::new(len, columns)
806 }
807}
808
809impl crate::formats::row_index::RowSource for Decoded {
810 fn height(&self) -> usize {
811 Decoded::height(self)
812 }
813
814 fn schema(&self) -> SchemaRef {
815 self.schema.clone()
816 }
817
818 fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
819 let rows = crate::formats::row_index::checked(index, self.height())?;
820 self.bytes.still_whole()?;
821 self.windows.with(
822 &rows,
823 |rows| self.lines(rows.iter().map(|&r| r as usize)),
824 |lines| self.column(column, lines),
825 )
826 }
827}
828
829impl crate::formats::pushdown::Windowed for Decoded {
830 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
831 Ok(self.collect_window(start, len)?.lazy())
832 }
833}
834
835#[derive(Debug, Clone, Default)]
839pub struct Layers {
840 pub dbcs: Vec<Arc<Dbc>>,
841}
842
843impl Layers {
844 pub fn new(registry: &crate::formats::Registry, dicts: &[PathBuf]) -> Result<Self> {
846 let mut dbcs: Vec<Arc<Dbc>> = registry.dbc.iter().map(|f| f.dbc.clone()).collect();
847 for path in dicts {
848 match crate::formats::dbc::load(path) {
849 Ok(Some(d)) => dbcs.push(Arc::new(d)),
850 Ok(None) => {
851 return Err(FileError::new(
852 path,
853 "not a DBC dictionary. --dict takes a .dbc file, or TOML with kind = \"dbc\".",
854 )
855 .into());
856 }
857 Err(e) => {
858 let at = e.path.as_deref().unwrap_or(path);
860 return Err(FileError::at(at, e.line, e.column, e.message).into());
861 }
862 }
863 }
864 Ok(Self { dbcs })
865 }
866
867 fn files(&self) -> Vec<PathBuf> {
868 self.dbcs.iter().filter_map(|d| d.path.clone()).collect()
869 }
870}
871
872#[derive(Debug, Default)]
875pub struct Listing {
876 pub layers: Layers,
877 pub messages: BTreeMap<String, (Arc<Message>, Arc<Vec<u32>>)>,
879 pub unknown: usize,
881}
882
883impl Listing {
884 pub fn resolve(index: &Index, layers: Layers) -> Self {
885 let mut known: HashMap<(u8, u32), Option<(usize, usize)>> = HashMap::new();
887 let lookup = |iface: u8, key: u32| -> Option<(usize, usize)> {
888 let name = index.interfaces.get(iface as usize)?;
889 let (id, extended) = (key & 0x7FFF_FFFF, key >> 31 == 1);
890 layers.dbcs.iter().enumerate().rev().find_map(|(l, dbc)| {
891 if !dbc.applies(name) {
892 return None;
893 }
894 dbc.messages
895 .iter()
896 .position(|m| m.id == id && m.extended == extended)
897 .map(|m| (l, m))
898 })
899 };
900 let mut rows: HashMap<(usize, usize), Vec<u32>> = HashMap::new();
901 let mut unknown = 0;
902 for (row, (&key, &iface)) in index.keys.iter().zip(&index.ifaces).enumerate() {
903 let found = *known
904 .entry((iface, key))
905 .or_insert_with(|| lookup(iface, key));
906 match found {
907 Some(at) => rows.entry(at).or_default().push(row as u32),
908 None => unknown += 1,
909 }
910 }
911 let mut messages = BTreeMap::new();
912 for ((l, m), rows) in rows {
913 let message = &layers.dbcs[l].messages[m];
914 let mut name = message.name.clone();
915 if name == FRAMES || name == SIGNALS || messages.contains_key(&name) {
916 name = format!("{name}.{}", layers.dbcs[l].name);
917 }
918 messages.insert(name, (Arc::new(message.clone()), Arc::new(rows)));
919 }
920 Self {
921 layers,
922 messages,
923 unknown,
924 }
925 }
926
927 pub fn tables(&self) -> Vec<Table> {
928 let frames = [
929 "ts", "iface", "id", "ext", "dlc", "data", "fd", "flags", "kind",
930 ];
931 let mut tables = vec![Table::plain(FRAMES, "frames", frames)];
932 if self.messages.is_empty() {
933 return tables;
934 }
935 let signals = ["ts", "message", "signal", "value", "unit"];
936 tables.push(Table::plain(SIGNALS, "signals", signals));
937 for (name, (message, _)) in &self.messages {
938 let columns =
939 std::iter::once("ts").chain(message.signals.iter().map(|s| s.name.as_str()));
940 tables.push(Table::plain(name, "message", columns));
941 }
942 tables
943 }
944}
945
946pub fn listed(path: &Path) -> Option<Vec<Table>> {
948 crate::formats::indexed::peek::<Listing>(path).map(|l| l.tables())
949}
950
951fn detail(index: &Index, listing: &Listing) -> Detail {
953 let group = crate::numfmt::group_chrome;
954 let mut lines = vec![
955 format!("Frames: {}", group(index.keys.len())),
956 format!("Interfaces: {}", index.interfaces.join(", ")),
957 format!(
958 "Timestamps: {}",
959 if index.absolute {
960 "wall clock"
961 } else {
962 "from the start of the capture"
963 }
964 ),
965 ];
966 if listing.layers.dbcs.is_empty() {
967 lines.push("DBC: none; --dict FILE or the format search path decodes signals".into());
968 } else {
969 for dbc in &listing.layers.dbcs {
970 lines.push(format!(
971 "DBC: {}{}, {} messages",
972 dbc.name,
973 dbc.interface
974 .as_ref()
975 .map(|i| format!(" on {i}"))
976 .unwrap_or_default(),
977 dbc.messages.len()
978 ));
979 }
980 lines.push(format!("Frames no DBC names: {}", group(listing.unknown)));
981 }
982 let list = listing.messages.iter().map(|(name, (message, rows))| {
983 (
984 name.clone(),
985 MetaValue::Text(format!(
986 "id {}, {} frames, {} signals{}",
987 if message.extended {
988 format!("{:08X}", message.id)
989 } else {
990 format!("{:03X}", message.id)
991 },
992 group(rows.len()),
993 message.signals.len(),
994 message
995 .comment
996 .as_ref()
997 .map(|c| format!(": {c}"))
998 .unwrap_or_default()
999 )),
1000 )
1001 });
1002 Detail {
1003 tab: crate::formats::text_formats::tab(crate::FileFormat::Candump),
1004 lines,
1005 list_title: "Messages",
1006 list: crate::formats::text_formats::capped_list(list, listing.messages.len()),
1007 first: false,
1008 ..Default::default()
1009 }
1010}
1011
1012fn long_table(
1014 bytes: &Arc<Bytes>,
1015 index: &Index,
1016 listing: &Listing,
1017 opened: &mut crate::formats::members::Opened,
1018) -> Result<LazyFrame> {
1019 let mut parts = Vec::new();
1020 let mut left_out = 0usize;
1021 for (name, (message, rows)) in &listing.messages {
1022 let decoded = Arc::new(
1023 Decoded::new(bytes.clone(), index, rows.clone(), message.clone(), false)
1024 .map_err(|e| eyre!("{e}"))?,
1025 );
1026 let lf = crate::formats::row_index::lazy(&decoded);
1027 for signal in &message.signals {
1028 if parts.len() >= MAX_LONG_PARTS {
1029 left_out += 1;
1030 continue;
1031 }
1032 parts.push(
1033 lf.clone()
1034 .select([
1035 col("ts"),
1036 lit(name.as_str()).alias("message"),
1037 lit(signal.name.as_str()).alias("signal"),
1038 col(signal.name.as_str())
1039 .cast(DataType::Float64)
1040 .alias("value"),
1041 lit(signal.unit.as_str()).alias("unit"),
1042 ])
1043 .filter(col("value").is_not_null()),
1044 );
1045 }
1046 }
1047 if left_out > 0 {
1048 opened.notes.push(crate::formats::text_formats::note(
1049 format!(
1050 "{} signals left out: past the first {MAX_LONG_PARTS} {} each message's table has them",
1051 crate::numfmt::group_chrome(left_out),
1052 crate::glyphs::get().middot
1053 ),
1054 "the dictionaries".to_string(),
1055 ));
1056 }
1057 if parts.is_empty() {
1058 return Err(eyre!(
1059 "no signals to decode: no frame of the log is a message of the dictionaries"
1060 ));
1061 }
1062 let all = concat(parts, UnionArgs::default())?;
1063 Ok(all.sort(
1064 ["ts"],
1065 SortMultipleOptions::default()
1066 .with_maintain_order(true)
1067 .with_nulls_last(true),
1068 ))
1069}
1070
1071fn scan(input: crate::formats::readers::ScanIn<'_>) -> Result<crate::loading::scan::Scan> {
1075 let path = &input.path().to_path_buf();
1076 let wanted = input.options.table.as_deref();
1077 let layers = Layers::new(input.formats, &input.options.dicts)?;
1078 let (bytes, index) = crate::formats::indexed::indexed(path, index)?;
1079 let layers = match crate::formats::indexed::peek::<Listing>(path) {
1082 Some(last) if layers.files().is_empty() || last.layers.files() == layers.files() => {
1083 last.layers.clone()
1084 }
1085 _ => layers,
1086 };
1087 let listing =
1088 crate::formats::indexed::cached::<Listing, std::convert::Infallible>(path, || {
1089 Ok(Listing::resolve(&index, layers.clone()))
1090 })
1091 .unwrap_or_else(|never| match never {});
1092 let listing = if listing.layers.files() == layers.files() {
1093 listing
1094 } else {
1095 crate::formats::indexed::forget::<Listing>(path);
1097 crate::formats::indexed::cached::<Listing, std::convert::Infallible>(path, || {
1098 Ok(Listing::resolve(&index, layers))
1099 })
1100 .unwrap_or_else(|never| match never {})
1101 };
1102 let tables = listing.tables();
1103 if let Some(wanted) = wanted
1106 && listing.layers.dbcs.is_empty()
1107 && !tables.iter().any(|t| t.name.eq_ignore_ascii_case(wanted))
1108 {
1109 return Err(FileError::new(
1110 path,
1111 format!(
1112 "no table \"{wanted}\": with no dictionary only its {FRAMES} are read. --dict names one that decodes its messages."
1113 ),
1114 )
1115 .into());
1116 }
1117 let picked = match crate::formats::members::pick(tables.clone(), wanted, path, "")? {
1118 crate::formats::sqlite::Pick::One(table) => table.name,
1119 crate::formats::sqlite::Pick::Several(tables) => {
1120 return Ok(crate::formats::members::several(&input, tables));
1121 }
1122 };
1123 let mut notes: Vec<String> = Vec::new();
1124 if index.skipped > 0 {
1125 notes.push(format!(
1126 "{} non-frame lines skipped",
1127 crate::numfmt::group_chrome(index.skipped)
1128 ));
1129 }
1130 if index.past_limit > 0 {
1131 notes.push(crate::limits::left_out(
1132 &format!("{} frames", crate::numfmt::group_chrome(index.past_limit)),
1133 crate::limits::get().indexed_records,
1134 "indexed_records",
1135 ));
1136 }
1137 for dbc in &listing.layers.dbcs {
1138 notes.extend(dbc.notes.iter().cloned());
1139 }
1140 let mut opened = crate::formats::members::Opened::for_table(
1141 detail(&index, &listing),
1142 &tables,
1143 &picked,
1144 notes,
1145 "the log",
1146 );
1147 let lf = if picked == FRAMES {
1148 let raw = Arc::new(RawFrames::new(bytes, &index));
1149 opened.window = Some((raw.clone(), raw.rows()));
1150 crate::formats::row_index::lazy(&raw)
1151 } else if picked == SIGNALS {
1152 long_table(&bytes, &index, &listing, &mut opened)?
1153 } else {
1154 let (message, rows) = listing
1155 .messages
1156 .get(&picked)
1157 .ok_or_else(|| FileError::new(path, format!("no message \"{picked}\"")))?;
1158 let decoded = Arc::new(
1159 Decoded::new(bytes, &index, rows.clone(), message.clone(), true)
1160 .map_err(|e| eyre!("{e}"))?,
1161 );
1162 opened.units = message
1163 .signals
1164 .iter()
1165 .filter(|s| !s.unit.is_empty())
1166 .map(|s| (s.name.clone(), s.unit.clone()))
1167 .collect();
1168 opened.window = Some((decoded.clone(), decoded.height()));
1169 crate::formats::row_index::lazy(&decoded)
1170 };
1171 Ok(opened.scan(input, lf))
1172}
1173
1174#[cfg(test)]
1175pub(crate) mod tests {
1176 use super::*;
1177
1178 #[test]
1181 fn a_window_is_parsed_once_for_all_its_columns() {
1182 let windows = Windows::new(2);
1183 let parse = |rows: &[IdxSize]| rows.to_vec();
1184 let parses = || windows.parses.load(std::sync::atomic::Ordering::Relaxed);
1185 let first = windows.with(&[1, 2], parse, |rows| rows.clone());
1186 let second = windows.with(&[1, 2], parse, |rows| rows.clone());
1187 assert_eq!((first, second, parses()), (vec![1, 2], vec![1, 2], 1));
1188 assert!(windows.kept.lock().unwrap().is_empty());
1189 windows.with(&[1, 3], parse, |_| ());
1190 windows.with(&[3], parse, |_| ());
1191 windows.with(&[1, 3], parse, |_| ());
1192 assert_eq!(parses(), 3);
1193 }
1194
1195 #[test]
1197 fn rows_that_hash_alike_are_parsed_apart() {
1198 let windows = Windows::new(2);
1199 let parse = |rows: &[IdxSize]| rows.to_vec();
1200 windows.with(&[1, 3], parse, |_| ());
1201 windows.kept.lock().unwrap()[0].key = WindowKey::of(&[2, 4]);
1203 let other = windows.with(&[2, 4], parse, |rows| rows.clone());
1204 assert_eq!(other, vec![2, 4]);
1205 assert_eq!(windows.parses.load(std::sync::atomic::Ordering::Relaxed), 2);
1206 }
1207
1208 #[test]
1210 fn windows_parse_in_parallel() {
1211 let windows = &Windows::new(9);
1212 let (parsed_b, b_done) = std::sync::mpsc::channel();
1213 let (began_a, a_running) = std::sync::mpsc::channel();
1214 std::thread::scope(|s| {
1215 s.spawn(move || {
1216 windows.with(
1217 &[0, 1],
1218 |_| {
1219 began_a.send(()).unwrap();
1220 b_done
1221 .recv_timeout(std::time::Duration::from_secs(10))
1222 .expect("the second window waited on the first");
1223 },
1224 |_| (),
1225 );
1226 });
1227 s.spawn(move || {
1228 a_running.recv().unwrap();
1230 windows.with(&[5, 6], |_| (), |_| ());
1231 parsed_b.send(()).unwrap();
1232 });
1233 });
1234 }
1235
1236 #[test]
1239 fn a_decoded_window_parses_each_line_once() {
1240 use crate::formats::row_index::RowSource;
1241 let log = b"(1.000001) can0 123#0102\n(1.000002) can1 1FFFFFFF#R\nnot a frame\n(1.000003) can0 456#03\n";
1242 let index = index(log).unwrap();
1243 let raw = RawFrames::new(Arc::new(Bytes::Owned(log.to_vec())), &index);
1244 let rows = IdxCa::from_vec("".into(), vec![0, 2]);
1245 let columns: Vec<Column> = (0..9).map(|c| raw.decode(c, &rows).unwrap()).collect();
1246 assert_eq!(
1247 raw.windows
1248 .parses
1249 .load(std::sync::atomic::Ordering::Relaxed),
1250 1
1251 );
1252 let df = DataFrame::new(2, columns).unwrap();
1253 let window = raw.collect_window(0, 3).unwrap();
1254 let picked = window.take(&rows).unwrap();
1255 assert!(df.equals_missing(&picked), "{df}\n{picked}");
1256 assert_eq!(
1257 df.column("iface").unwrap().str().unwrap().get(1),
1258 Some("can0")
1259 );
1260 }
1261
1262 #[test]
1265 fn errors_name_the_file() {
1266 use crate::formats::readers::bad_input::{assert_shape, each_names_its_file, opening};
1267 each_names_its_file(
1268 crate::FileFormat::Candump,
1269 &[("text.log", b"hello there\n", "No line is a CAN frame")],
1270 );
1271 let dir = tempfile::tempdir().unwrap();
1272 let log = b"(1436509052.249713) can0 123#DEADBEEF\n";
1273 let dbc = dir.path().join("plain.toml");
1274 std::fs::write(&dbc, "a = 1\n").unwrap();
1275 for (options, named, says) in [
1276 (
1277 crate::OpenOptions {
1278 table: Some("Engine".into()),
1279 ..Default::default()
1280 },
1281 dir.path().join("a.log"),
1282 "--dict names one",
1283 ),
1284 (
1285 crate::OpenOptions {
1286 dicts: vec![dbc.clone()],
1287 ..Default::default()
1288 },
1289 dbc.clone(),
1290 "--dict takes",
1291 ),
1292 ] {
1293 let message = opening(
1294 dir.path(),
1295 "a.log",
1296 log,
1297 crate::FileFormat::Candump,
1298 &options,
1299 )
1300 .expect("refused");
1301 eprintln!("{message}");
1302 assert_shape(&message, &named);
1303 assert!(message.contains(says), "{message}");
1304 }
1305 }
1306
1307 pub(crate) const LOG: &str = "(1700000000.000100) can0 123#401F7602\n\
1308(1700000000.000200) can1 18FEF1FE#1234FFF000000000\n\
1309# a comment\n\
1310(1700000000.000300) can0 200#01102700\n\
1311(1700000000.000400) can0 200#02F6FF0000000000\n\
1312(1700000000.000500) can0 7DF#R\n\
1313(1700000000.000600) can0 123##3112233445566778899AABBCC\n\
1314(1700000000.000700) can0 456#DEAD\n";
1315
1316 #[test]
1317 fn each_line_form() {
1318 let f = parse_line("(1436509052.249713) vcan0 044#2A366C2BBA").unwrap();
1319 assert_eq!(f.ts, Some(1_436_509_052_249_713));
1320 assert_eq!((f.id, f.extended, f.dlc), (0x44, false, 5));
1321 let f = parse_line("(0.5) can0 12345678#R").unwrap();
1322 assert!(f.remote && f.extended);
1323 let f = parse_line("(1.0) can0 123##1DEADBEEF").unwrap();
1324 assert!(f.fd);
1325 assert_eq!(f.flags, Some(1));
1326 assert_eq!(f.data, [0xDE, 0xAD, 0xBE, 0xEF]);
1327 let f = parse_line(" can0 123 [4] DE AD BE EF").unwrap();
1328 assert_eq!((f.ts, f.dlc), (None, 4));
1329 let f = parse_line(" (1436509052.249713) can0 1F334455 [2] 01 02").unwrap();
1330 assert!(f.extended);
1331 let f = parse_line("(2024-01-31 08:15:00.123456) can0 123 [1] FF").unwrap();
1332 assert!(f.ts.is_some());
1333 let f = parse_line(" can0 321 [8] remote request").unwrap();
1334 assert!(f.remote);
1335 assert_eq!(parse_line("can0 123#ABC"), None);
1336 assert_eq!(parse_line("hello world"), None);
1337 assert!(looks_like(LOG.as_bytes()));
1338 assert!(!looks_like(b"a,b,c\n1,2,3\n"));
1339 }
1340
1341 #[test]
1342 fn frames_and_decoded_messages() {
1343 let index = index(LOG.as_bytes()).unwrap();
1344 assert_eq!(index.keys.len(), 7);
1345 assert_eq!(index.skipped, 1);
1346 assert!(index.absolute);
1347 let bytes = Arc::new(Bytes::Owned(LOG.as_bytes().to_vec()));
1348 let raw = Arc::new(RawFrames::new(bytes.clone(), &index));
1349 let df = crate::formats::row_index::lazy(&raw).collect().unwrap();
1350 assert_eq!(df.height(), 7);
1351 assert_eq!(
1352 df.column("id").unwrap().str().unwrap().get(1),
1353 Some("18FEF1FE")
1354 );
1355 assert_eq!(
1356 df.column("kind").unwrap().str().unwrap().get(4),
1357 Some("remote")
1358 );
1359
1360 let dbc = crate::formats::dbc::parse(crate::tests::fixtures::DBC, "car", None).unwrap();
1361 let layers = Layers {
1362 dbcs: vec![Arc::new(dbc)],
1363 };
1364 let listing = Listing::resolve(&index, layers);
1365 let names: Vec<&String> = listing.messages.keys().collect();
1366 assert_eq!(names, ["BODY", "ENGINE", "MUXED"]);
1367 assert_eq!(listing.unknown, 2);
1368 let (message, rows) = &listing.messages["ENGINE"];
1369 let decoded = Arc::new(
1370 Decoded::new(bytes.clone(), &index, rows.clone(), message.clone(), true).unwrap(),
1371 );
1372 let df = crate::formats::row_index::lazy(&decoded).collect().unwrap();
1373 assert_eq!(df.height(), 2);
1375 assert_eq!(
1376 df.column("Speed").unwrap().f64().unwrap().get(0),
1377 Some(1000.0)
1378 );
1379 assert_eq!(df.column("Temp").unwrap().f64().unwrap().get(0), Some(78.0));
1380 assert_eq!(
1381 df.column("Gear").unwrap().str().unwrap().get(0),
1382 Some("Second")
1383 );
1384 let (message, rows) = &listing.messages["MUXED"];
1385 let decoded = Arc::new(
1386 Decoded::new(bytes.clone(), &index, rows.clone(), message.clone(), true).unwrap(),
1387 );
1388 let df = crate::formats::row_index::lazy(&decoded).collect().unwrap();
1389 let volts = df.column("Volts").unwrap().f64().unwrap();
1390 let amps = df.column("Amps").unwrap().f64().unwrap();
1391 assert_eq!(volts.get(0), Some(100.0));
1392 assert_eq!(amps.get(0), None);
1393 assert_eq!(volts.get(1), None);
1394 assert!((amps.get(1).unwrap() - -1.0).abs() < 1e-9);
1395 }
1396}