1use std::collections::{BTreeMap, HashMap};
16use std::path::Path;
17use std::sync::Arc;
18
19use color_eyre::Result;
20use color_eyre::eyre::eyre;
21
22use crate::error_display::{FileError, in_file};
23use polars::prelude::*;
24
25use crate::fixed_records::{Bytes, ColumnLayout, Logical, Physical};
26use crate::indexed::{IndexedRecords, Offsets};
27use crate::model_files::MetaValue;
28use crate::sqlite::Table;
29use crate::text_formats::Detail;
30
31pub(crate) const READER: crate::readers::Reader = crate::readers::Reader {
33 scan,
34 signatures: &[crate::readers::Signature {
35 says: |head, _| looks_like(head),
36 kind: crate::readers::Kind::Magic,
37 trusted: crate::readers::Trusted {
38 tables: true,
39 ..crate::readers::EVERYWHERE
40 },
41 }],
42 tables: Some(listed),
43 ..crate::readers::BASE
44};
45
46pub const MAGIC: &[u8; 7] = b"ULog\x01\x12\x35";
48
49const SYNC: [u8; 8] = [0x2F, 0x73, 0x13, 0x20, 0x25, 0x0C, 0xBB, 0x12];
51
52pub const MAX_COLUMNS: usize = 4096;
54const MAX_DEPTH: usize = 16;
56const MAX_DEFINITIONS: usize = 65_536;
58const MAX_LOGGED: usize = 1_000_000;
60
61pub const LOGGED: &str = "logged_messages";
63pub const PARAMETERS: &str = "parameters";
64
65pub fn looks_like(head: &[u8]) -> bool {
67 head.starts_with(MAGIC)
68}
69
70#[derive(Debug, Clone, PartialEq)]
72struct FieldDef {
73 kind: String,
74 count: Option<usize>,
76 name: String,
77}
78
79#[derive(Debug)]
81pub struct Topic {
82 pub name: String,
83 pub multi_id: u8,
84 pub offsets: Arc<Offsets>,
86 pub needs: usize,
88 pub columns: Vec<ColumnLayout>,
89}
90
91#[derive(Debug, Default)]
93pub struct Index {
94 pub version: u8,
95 pub start_us: u64,
96 pub topics: BTreeMap<u16, Topic>,
98 pub names: Vec<(String, u16)>,
100 pub info: Vec<(String, String)>,
101 pub parameters: Vec<(String, &'static str, f64, Option<u64>)>,
104 pub logged: Vec<(u64, &'static str, Option<u16>, String)>,
106 pub logged_left_out: usize,
107 pub dropouts: usize,
108 pub dropout_ms: u64,
109 pub skipped: usize,
111 pub damaged: usize,
112 pub short: usize,
114 pub unsubscribed: usize,
116 pub past_limit: usize,
118 pub unread: Vec<(String, String)>,
120 pub cut_short: bool,
122}
123
124fn primitive(kind: &str) -> Option<(Physical, usize)> {
126 Some(match kind {
127 "int8_t" => (Physical::Signed(1), 1),
128 "uint8_t" => (Physical::Unsigned(1), 1),
129 "int16_t" => (Physical::Signed(2), 2),
130 "uint16_t" => (Physical::Unsigned(2), 2),
131 "int32_t" => (Physical::Signed(4), 4),
132 "uint32_t" => (Physical::Unsigned(4), 4),
133 "int64_t" => (Physical::Signed(8), 8),
134 "uint64_t" => (Physical::Unsigned(8), 8),
135 "float" => (Physical::Float(4), 4),
136 "double" => (Physical::Float(8), 8),
137 "bool" => (Physical::Bool, 1),
138 "char" => (Physical::Text, 1),
139 _ => return None,
140 })
141}
142
143fn parse_format(text: &str) -> Option<(String, Vec<FieldDef>)> {
145 let (name, rest) = text.split_once(':')?;
146 let mut fields = Vec::new();
147 for part in rest.split(';').filter(|p| !p.trim().is_empty()) {
148 let (kind, field) = part.trim().split_once(' ')?;
149 let (kind, count) = match kind.split_once('[') {
150 Some((kind, n)) => (kind, Some(n.strip_suffix(']')?.parse().ok()?)),
151 None => (kind, None),
152 };
153 fields.push(FieldDef {
154 kind: kind.to_string(),
155 count,
156 name: field.trim().to_string(),
157 });
158 if fields.len() > MAX_COLUMNS {
159 return None;
160 }
161 }
162 Some((name.to_string(), fields))
163}
164
165fn size_of(
167 name: &str,
168 formats: &HashMap<String, Vec<FieldDef>>,
169 depth: usize,
170) -> std::result::Result<usize, String> {
171 if depth > MAX_DEPTH {
172 return Err("types nest too deep".into());
173 }
174 let fields = formats
175 .get(name)
176 .ok_or_else(|| format!("no format for {name}"))?;
177 let mut size = 0usize;
178 for f in fields {
179 let one = match primitive(&f.kind) {
180 Some((_, n)) => n,
181 None => size_of(&f.kind, formats, depth + 1)?,
182 };
183 size = one
184 .checked_mul(f.count.unwrap_or(1))
185 .and_then(|n| size.checked_add(n))
186 .ok_or("a type is too large")?;
187 }
188 Ok(size)
189}
190
191fn columns_of(
194 name: &str,
195 formats: &HashMap<String, Vec<FieldDef>>,
196 prefix: &str,
197 offset: usize,
198 out: &mut Vec<ColumnLayout>,
199 depth: usize,
200) -> std::result::Result<usize, String> {
201 if depth > MAX_DEPTH {
202 return Err("types nest too deep".into());
203 }
204 let fields = formats
205 .get(name)
206 .ok_or_else(|| format!("no format for {name}"))?;
207 let mut at = offset;
208 let mut needs = offset;
209 for f in fields {
210 let count = f.count.unwrap_or(1);
211 let full = format!("{prefix}{}", f.name);
212 match primitive(&f.kind) {
213 Some((physical, width)) => {
214 let bytes = width.checked_mul(count).ok_or("a field is too large")?;
215 if !f.name.starts_with("_padding") && bytes > 0 {
217 let mut column = match physical {
218 Physical::Text => ColumnLayout::new(&full, at, 0, physical, bytes),
220 _ => {
221 let mut c = ColumnLayout::new(&full, at, 0, physical, width);
222 c.count = count;
223 c
224 }
225 };
226 if depth == 0 && f.name == "timestamp" && f.kind == "uint64_t" && count == 1 {
228 column.logical = Logical::Duration {
229 unit: TimeUnit::Microseconds,
230 multiplier: 1,
231 };
232 }
233 out.push(column);
234 needs = at + bytes;
235 }
236 at = at.checked_add(bytes).ok_or("a type is too large")?;
237 }
238 None => {
239 let size = size_of(&f.kind, formats, depth + 1)?;
240 for i in 0..count {
241 let inner = match f.count {
242 Some(_) => format!("{full}[{i}]."),
243 None => format!("{full}."),
244 };
245 let end = columns_of(&f.kind, formats, &inner, at, out, depth + 1)?;
246 if end > at {
247 needs = needs.max(end);
248 }
249 at = at.checked_add(size).ok_or("a type is too large")?;
250 if out.len() > MAX_COLUMNS {
251 return Err(format!("its fields make more than {MAX_COLUMNS} columns"));
252 }
253 }
254 }
255 }
256 if out.len() > MAX_COLUMNS {
257 return Err(format!("its fields make more than {MAX_COLUMNS} columns"));
258 }
259 }
260 Ok(needs)
261}
262
263fn u16_at(b: &[u8], at: usize) -> Option<u16> {
264 Some(u16::from_le_bytes(b.get(at..at + 2)?.try_into().ok()?))
265}
266
267fn u64_at(b: &[u8], at: usize) -> Option<u64> {
268 Some(u64::from_le_bytes(b.get(at..at + 8)?.try_into().ok()?))
269}
270
271fn key_value(payload: &[u8]) -> Option<(String, String, String, Option<f64>)> {
273 let len = *payload.first()? as usize;
274 let key = std::str::from_utf8(payload.get(1..1 + len)?).ok()?;
275 let value = payload.get(1 + len..)?;
276 let (kind, name) = key.split_once(' ')?;
277 let (base, count) = match kind.split_once('[') {
278 Some((base, n)) => (base, n.trim_end_matches(']').parse::<usize>().ok()),
279 None => (kind, None),
280 };
281 let number = |v: Option<f64>| v.map(|n| (format_number(n), Some(n)));
282 let (text, number) = match (base, count) {
283 ("char", _) => (
284 String::from_utf8_lossy(value)
285 .trim_end_matches('\0')
286 .to_string(),
287 None,
288 ),
289 (_, Some(_)) => (crate::fixed_records::hex(value), None),
290 _ => {
291 let (physical, width) = primitive(base)?;
292 let raw = value.get(..width)?;
293 let v = match physical {
294 Physical::Signed(_) => crate::fixed_records::read_signed(raw, false) as f64,
295 Physical::Unsigned(_) | Physical::Bool => {
296 crate::fixed_records::read_unsigned(raw, false) as f64
297 }
298 Physical::Float(4) => f32::from_le_bytes(raw.try_into().ok()?) as f64,
299 Physical::Float(_) => f64::from_le_bytes(raw.try_into().ok()?),
300 _ => return None,
301 };
302 number(Some(v))?
303 }
304 };
305 Some((kind.to_string(), name.to_string(), text, number))
306}
307
308fn format_number(n: f64) -> String {
309 if n.fract() == 0.0 && n.abs() < 1e15 {
310 format!("{}", n as i64)
311 } else {
312 format!("{n}")
313 }
314}
315
316fn level(byte: u8) -> &'static str {
317 match byte {
318 b'0' => "emergency",
319 b'1' => "alert",
320 b'2' => "critical",
321 b'3' => "error",
322 b'4' => "warning",
323 b'5' => "notice",
324 b'6' => "info",
325 b'7' => "debug",
326 _ => "unknown",
327 }
328}
329
330pub fn index(data: &[u8]) -> std::result::Result<Index, String> {
332 if !looks_like(data) {
333 return Err("not a ULog file: no ULog magic at the start".into());
334 }
335 if data.len() < 16 {
336 return Err("the header is cut short".into());
337 }
338 let mut index = Index {
339 version: data[7],
340 start_us: u64_at(data, 8).unwrap_or(0),
341 ..Index::default()
342 };
343 let mut formats: HashMap<String, Vec<FieldDef>> = HashMap::new();
344 let mut subs: HashMap<u16, (String, u8, Offsets)> = HashMap::new();
346 let mut multi: Vec<(String, String)> = Vec::new();
347 let mut records = 0usize;
348 let mut last_time: Option<u64> = None;
349 let mut appended: Vec<usize> = Vec::new();
351 let mut end = data.len();
352 let mut at = 16usize;
353 loop {
354 if at >= end {
355 match appended.first().copied() {
356 Some(next) if next >= at.min(end) && next < data.len() => {
357 appended.remove(0);
358 at = next;
359 end = appended.first().copied().unwrap_or(data.len()).max(next);
360 continue;
361 }
362 _ => break,
363 }
364 }
365 let Some(size) = u16_at(data, at) else {
366 index.cut_short = at < end;
367 break;
368 };
369 let size = size as usize;
370 let Some(&kind) = data.get(at + 2) else {
371 index.cut_short = true;
372 break;
373 };
374 let body = at + 3;
375 if body + size > end {
376 if let Some(next) = find_sync(data, at + 1, end) {
378 index.skipped += next - at;
379 index.damaged += 1;
380 at = next;
381 continue;
382 }
383 index.cut_short = true;
384 break;
385 }
386 let payload = &data[body..body + size];
387 match kind {
388 b'B' if size >= 40 => {
389 if payload[8] & 1 != 0 {
391 for i in 0..3 {
392 if let Some(o) = u64_at(payload, 16 + i * 8)
393 && o > 0
394 && let Ok(o) = usize::try_from(o)
395 && o > body + size
396 && o < data.len()
397 {
398 appended.push(o);
399 }
400 }
401 appended.sort_unstable();
402 if let Some(&first) = appended.first() {
403 end = first;
404 }
405 }
406 }
407 b'F' => {
408 if formats.len() < MAX_DEFINITIONS
409 && let Some((name, fields)) = parse_format(&String::from_utf8_lossy(payload))
410 {
411 formats.insert(name, fields);
412 }
413 }
414 b'A' if size >= 3 => {
415 let multi_id = payload[0];
416 let id = u16_at(payload, 1).unwrap_or_default();
417 let name = String::from_utf8_lossy(&payload[3..])
418 .trim_end_matches('\0')
419 .to_string();
420 if subs.len() < MAX_DEFINITIONS {
421 subs.insert(id, (name, multi_id, Offsets::for_file(data.len())));
422 }
423 }
424 b'D' if size >= 2 => {
425 let id = u16_at(payload, 0).unwrap_or_default();
426 match subs.get_mut(&id) {
427 Some((_, _, offsets)) if records < crate::indexed::MAX_RECORDS => {
428 offsets.push(body + 2);
429 records += 1;
430 }
431 Some(_) => index.past_limit += 1,
432 None => index.unsubscribed += 1,
433 }
434 }
435 b'I' => {
436 if let Some((_, name, text, _)) = key_value(payload)
437 && index.info.len() < MAX_DEFINITIONS
438 {
439 index.info.push((name, text));
440 }
441 }
442 b'M' if size >= 1 => {
443 if let Some((_, name, text, _)) = key_value(&payload[1..]) {
445 let continued = payload[0] != 0;
446 match multi.iter().position(|(n, _)| *n == name) {
447 Some(i) if continued => {
448 let value = &mut multi[i].1;
449 if value.len() < 1 << 20 {
450 value.push_str(&text);
451 }
452 }
453 _ if multi.len() < MAX_DEFINITIONS => multi.push((name, text)),
454 _ => {}
455 }
456 }
457 }
458 b'P' => {
459 if let Some((kind, name, _, Some(value))) = key_value(payload)
460 && index.parameters.len() < MAX_DEFINITIONS
461 {
462 let kind = if kind == "float" { "float" } else { "int32" };
463 index.parameters.push((name, kind, value, last_time));
464 }
465 }
466 b'L' if size >= 9 => {
467 let time = u64_at(payload, 1).unwrap_or_default();
468 last_time = Some(time);
469 push_logged(&mut index, time, payload[0], None, &payload[9..]);
470 }
471 b'C' if size >= 11 => {
472 let tag = u16_at(payload, 1);
473 let time = u64_at(payload, 3).unwrap_or_default();
474 last_time = Some(time);
475 push_logged(&mut index, time, payload[0], tag, &payload[11..]);
476 }
477 b'O' if size >= 2 => {
478 index.dropouts += 1;
479 index.dropout_ms += u64::from(u16_at(payload, 0).unwrap_or_default());
480 }
481 b'R' | b'S' | b'Q' | b'B' | b'A' | b'D' | b'L' | b'C' | b'O' | b'M' => {}
482 _ => {
483 match find_sync(data, at + 1, end) {
485 Some(next) => {
486 index.skipped += next - at;
487 index.damaged += 1;
488 at = next;
489 continue;
490 }
491 None => {
492 index.skipped += end - at;
493 index.damaged += 1;
494 at = end;
495 continue;
496 }
497 }
498 }
499 }
500 at = body + size;
501 }
502 for (name, value) in multi {
503 index.info.push((name, value));
504 }
505
506 let mut counts: HashMap<String, usize> = HashMap::new();
508 for (name, _, offsets) in subs.values() {
509 if !offsets.is_empty() {
510 *counts.entry(name.clone()).or_default() += 1;
511 }
512 }
513 let mut ids: Vec<u16> = subs.keys().copied().collect();
514 ids.sort_unstable_by_key(|id| (subs[id].0.clone(), subs[id].1));
515 for id in ids {
516 let (name, multi_id, mut offsets) = subs.remove(&id).expect("listed above");
517 if offsets.is_empty() {
518 continue;
519 }
520 let mut columns = Vec::new();
521 let needs = match columns_of(&name, &formats, "", 0, &mut columns, 0).and_then(|needs| {
522 let mut seen = std::collections::HashSet::new();
524 match columns.iter().find(|c| !seen.insert(c.name.clone())) {
525 Some(c) => Err(format!("it has two fields named {}", c.name)),
526 None => Ok(needs),
527 }
528 }) {
529 Ok(needs) => needs,
530 Err(why) => {
531 index.unread.push((name, why));
532 continue;
533 }
534 };
535 let before = offsets.len();
538 offsets = keep_long_enough(offsets, data, needs);
539 index.short += before - offsets.len();
540 if offsets.is_empty() {
541 continue;
542 }
543 offsets.shrink();
544 let table = if counts.get(&name).copied().unwrap_or(0) > 1 {
545 format!("{name}.{multi_id}")
546 } else {
547 name.clone()
548 };
549 index.names.push((table, id));
550 index.topics.insert(
551 id,
552 Topic {
553 name,
554 multi_id,
555 offsets: Arc::new(offsets),
556 needs,
557 columns,
558 },
559 );
560 }
561 Ok(index)
562}
563
564fn keep_long_enough(offsets: Offsets, data: &[u8], needs: usize) -> Offsets {
567 let long_enough = |at: usize| {
568 u16_at(data, at - 5).is_some_and(|size| (size as usize).saturating_sub(2) >= needs)
569 };
570 let mut kept = Offsets::for_file(data.len());
571 for i in 0..offsets.len() {
572 let at = offsets.get(i);
573 if long_enough(at) {
574 kept.push(at);
575 }
576 }
577 kept
578}
579
580fn push_logged(index: &mut Index, time: u64, level_byte: u8, tag: Option<u16>, text: &[u8]) {
581 if index.logged.len() >= MAX_LOGGED {
582 index.logged_left_out += 1;
583 return;
584 }
585 let text = String::from_utf8_lossy(text)
586 .trim_end_matches('\0')
587 .to_string();
588 index.logged.push((time, level(level_byte), tag, text));
589}
590
591fn find_sync(data: &[u8], from: usize, end: usize) -> Option<usize> {
592 let hay = data.get(from..end)?;
593 memchr::memmem::find(hay, &SYNC).map(|i| from + i + SYNC.len())
594}
595
596pub fn listed(file: &Path) -> Result<Vec<Table>> {
599 crate::indexed::peek::<Index>(file)
600 .map(|index| tables(&index))
601 .ok_or_else(|| eyre!("Open the log to list its tables."))
602}
603
604pub fn tables(index: &Index) -> Vec<Table> {
606 let mut tables: Vec<Table> = index
607 .names
608 .iter()
609 .map(|(name, id)| Table {
610 name: name.clone(),
611 kind: "topic".to_string(),
612 internal: false,
613 columns: index.topics[id]
614 .columns
615 .iter()
616 .map(|c| (c.name.to_string(), String::new()))
617 .collect(),
618 })
619 .collect();
620 let taken = |name: &str| index.names.iter().any(|(n, _)| n == name);
621 if !index.logged.is_empty() && !taken(LOGGED) {
622 tables.push(Table {
623 name: LOGGED.to_string(),
624 kind: "messages".to_string(),
625 internal: false,
626 columns: ["timestamp", "level", "tag", "message"]
627 .iter()
628 .map(|c| (c.to_string(), String::new()))
629 .collect(),
630 });
631 }
632 if !index.parameters.is_empty() && !taken(PARAMETERS) {
633 tables.push(Table {
634 name: PARAMETERS.to_string(),
635 kind: "parameters".to_string(),
636 internal: false,
637 columns: ["name", "type", "value", "timestamp"]
638 .iter()
639 .map(|c| (c.to_string(), String::new()))
640 .collect(),
641 });
642 }
643 tables
644}
645
646pub fn detail(index: &Index) -> Detail {
648 let count = crate::text_formats::count;
649 let mut lines = vec![
650 format!("Version: {}", index.version),
651 format!(
652 "Topics: {}",
653 count(index.names.len() as u64, "table", "tables")
654 ),
655 ];
656 if index.dropouts > 0 {
657 lines.push(format!(
658 "Dropouts: {}, {} ms in all",
659 index.dropouts, index.dropout_ms
660 ));
661 }
662 let mut list: Vec<(String, MetaValue)> = index
663 .info
664 .iter()
665 .map(|(k, v)| (k.clone(), MetaValue::Text(v.clone())))
666 .collect();
667 let mut seen = std::collections::HashSet::new();
668 for (name, _, value, at) in &index.parameters {
669 if at.is_none() && seen.insert(name.clone()) {
671 list.push((
672 format!("param {name}"),
673 MetaValue::Text(format_number(*value)),
674 ));
675 }
676 }
677 let total = list.len();
678 Detail {
679 tab: crate::text_formats::tab(crate::FileFormat::Ulog),
680 lines,
681 list_title: "Info and parameters",
682 list: crate::text_formats::capped_list(list.into_iter(), total),
683 first: false,
684 ..Default::default()
685 }
686}
687
688pub fn notes(index: &Index) -> Vec<String> {
690 let group = |n: usize| crate::numfmt::group_chrome(n);
691 let mut notes = Vec::new();
692 if index.damaged > 0 {
693 notes.push(format!(
694 "{} damaged stretches skipped ({} bytes)",
695 group(index.damaged),
696 group(index.skipped)
697 ));
698 }
699 if index.cut_short {
700 notes.push("log cut short mid-message".to_string());
701 }
702 if index.short > 0 {
703 notes.push(format!(
704 "{} data messages left out: too short for their topic",
705 group(index.short)
706 ));
707 }
708 if index.unsubscribed > 0 {
709 notes.push(format!(
710 "{} data messages left out: no subscription",
711 group(index.unsubscribed)
712 ));
713 }
714 if index.past_limit > 0 {
715 notes.push(format!(
716 "{} messages left out: past the first {}",
717 group(index.past_limit),
718 group(crate::indexed::MAX_RECORDS)
719 ));
720 }
721 if index.logged_left_out > 0 {
722 notes.push(format!(
723 "{} logged messages left out: past the first {}",
724 group(index.logged_left_out),
725 group(MAX_LOGGED)
726 ));
727 }
728 for (topic, why) in &index.unread {
729 notes.push(format!("topic {topic} not read: {why}"));
730 }
731 notes
732}
733
734pub fn indexed(path: &Path) -> Result<(Arc<Bytes>, Arc<Index>)> {
736 let bytes = Arc::new(Bytes::map(path).map_err(|e| in_file(path, e.into()))?);
737 let index = crate::indexed::cached(path, || index(bytes.as_slice()))
738 .map_err(|e| FileError::new(path, e))?;
739 Ok((bytes, index))
740}
741
742pub enum Open {
744 Table {
745 lf: Box<LazyFrame>,
746 opened: Box<crate::members::Opened>,
747 },
748 Several(Vec<String>),
749}
750
751pub fn open(path: &Path, wanted: Option<&str>) -> Result<Open> {
753 let (bytes, index) = indexed(path)?;
754 let tables = tables(&index);
755 let picked = match crate::members::pick(
756 tables.clone(),
757 wanted,
758 path,
759 " The log has no data messages.",
760 )? {
761 crate::sqlite::Pick::One(table) => table.name,
762 crate::sqlite::Pick::Several(tables) => {
763 return Ok(Open::Several(tables.into_iter().map(|t| t.name).collect()));
764 }
765 };
766 let mut opened = crate::members::Opened {
767 detail: Some(Arc::new(detail(&index))),
768 other_tables: crate::members::others(&tables, &picked),
769 notes: notes(&index)
770 .into_iter()
771 .map(|n| crate::text_formats::note(n, "the log".to_string()))
772 .collect(),
773 ..Default::default()
774 };
775 let lf = if let Some((_, id)) = index.names.iter().find(|(n, _)| *n == picked) {
776 let topic = &index.topics[id];
777 let records = Arc::new(
778 IndexedRecords::new(bytes, topic.offsets.clone(), topic.columns.clone())
779 .map_err(|e| FileError::new(path, format!("table \"{picked}\": {e}")))?,
780 );
781 opened.window = Some((records.clone(), records.rows()));
782 records.lazy()
783 } else if picked == LOGGED {
784 let (mut time, mut lvl, mut tag, mut text) =
785 (Vec::new(), Vec::new(), Vec::new(), Vec::new());
786 for (t, l, g, m) in &index.logged {
787 time.push(*t as i64);
788 lvl.push(*l);
789 tag.push(*g);
790 text.push(m.as_str());
791 }
792 df!(
793 "timestamp" => Int64Chunked::from_vec("timestamp".into(), time).into_duration(TimeUnit::Microseconds).into_series(),
794 "level" => lvl,
795 "tag" => tag,
796 "message" => text,
797 )?
798 .lazy()
799 } else {
800 let (mut name, mut kind, mut value, mut time) =
801 (Vec::new(), Vec::new(), Vec::new(), Vec::new());
802 for (n, k, v, t) in &index.parameters {
803 name.push(n.as_str());
804 kind.push(*k);
805 value.push(*v);
806 time.push(t.map(|t| t as i64));
807 }
808 df!(
809 "name" => name,
810 "type" => kind,
811 "value" => value,
812 "timestamp" => time.into_iter().collect::<Int64Chunked>().into_duration(TimeUnit::Microseconds).into_series(),
813 )?
814 .lazy()
815 };
816 Ok(Open::Table {
817 lf: Box::new(lf),
818 opened: Box::new(opened),
819 })
820}
821
822fn scan(input: crate::readers::ScanIn<'_>) -> Result<crate::scan::Scan> {
826 let file = input.path();
827 Ok(match open(file, input.options.table.as_deref())? {
828 Open::Table { lf, opened } => {
829 input.report.opened = Some(Arc::new(*opened));
830 (*lf).into()
831 }
832 Open::Several(tables) => crate::scan::Scan::Tables {
833 file: file.to_path_buf(),
834 tables,
835 format: input.format,
836 },
837 })
838}
839
840#[cfg(test)]
841pub(crate) mod tests {
842 use super::*;
843
844 #[test]
846 fn errors_name_the_file() {
847 crate::readers::bad_input::each_names_its_file(
848 crate::FileFormat::Ulog,
849 &[
850 ("text.ulg", b"hello there, this is text", "Not a ULog file"),
851 ("cut.ulg", b"ULog\x01\x12\x35\x01", "cut short"),
852 ],
853 );
854 }
855
856 fn message(kind: u8, payload: &[u8]) -> Vec<u8> {
857 let mut out = (payload.len() as u16).to_le_bytes().to_vec();
858 out.push(kind);
859 out.extend(payload);
860 out
861 }
862
863 pub(crate) fn tiny() -> Vec<u8> {
867 let mut log = MAGIC.to_vec();
868 log.push(1);
869 log.extend(1_000u64.to_le_bytes());
870 log.extend(message(b'F', b"vec3:float x;float y;float z;"));
871 log.extend(message(
872 b'F',
873 b"sensor:uint64_t timestamp;uint8_t id;vec3 v;char[4] tag;int16_t[2] raw;uint8_t[3] _padding0;",
874 ));
875 log.extend(message(b'F', b"status:uint64_t timestamp;bool armed;"));
876 let mut info = vec![b"char[6] sys_name".len() as u8];
877 info.extend(b"char[6] sys_name");
878 info.extend(b"PX4\0\0\0");
879 log.extend(message(b'I', &info));
880 let mut param = vec![b"float MPC_XY_VEL".len() as u8];
881 param.extend(b"float MPC_XY_VEL");
882 param.extend(2.5f32.to_le_bytes());
883 log.extend(message(b'P', ¶m));
884 for (multi, id, name) in [(0u8, 1u16, "sensor"), (1, 2, "sensor"), (0, 3, "status")] {
885 let mut a = vec![multi];
886 a.extend(id.to_le_bytes());
887 a.extend(name.as_bytes());
888 log.extend(message(b'A', &a));
889 }
890 let sensor = |id: u16, t: u64, x: f32, pad: bool| {
891 let mut d = id.to_le_bytes().to_vec();
892 d.extend(t.to_le_bytes());
893 d.push(id as u8);
894 for v in [x, x * 2.0, x * 3.0] {
895 d.extend(v.to_le_bytes());
896 }
897 d.extend(b"ab\0\0");
898 d.extend(7i16.to_le_bytes());
899 d.extend((-7i16).to_le_bytes());
900 if pad {
901 d.extend([0u8; 3]);
902 }
903 message(b'D', &d)
904 };
905 log.extend(sensor(1, 100, 1.0, true));
906 log.extend(sensor(2, 110, 5.0, false));
907 let mut l = vec![b'4'];
908 l.extend(120u64.to_le_bytes());
909 l.extend(b"low battery");
910 log.extend(message(b'L', &l));
911 log.extend([0xEE; 7]);
913 log.extend(message(b'S', &SYNC));
914 log.extend(sensor(1, 200, 2.0, false));
915 let mut s = 3u16.to_le_bytes().to_vec();
916 s.extend(150u64.to_le_bytes());
917 s.push(1);
918 log.extend(message(b'D', &s));
919 let cut = sensor(1, 300, 3.0, false);
921 log.extend(&cut[..cut.len() - 4]);
922 log
923 }
924
925 #[test]
926 fn topics_instances_and_text() {
927 let index = index(&tiny()).unwrap();
928 let names: Vec<&str> = index.names.iter().map(|(n, _)| n.as_str()).collect();
929 assert_eq!(names, ["sensor.0", "sensor.1", "status"]);
930 let sensor = &index.topics[&1];
931 assert_eq!(sensor.offsets.len(), 2);
932 let columns: Vec<&str> = sensor.columns.iter().map(|c| c.name.as_str()).collect();
933 assert_eq!(
934 columns,
935 ["timestamp", "id", "v.x", "v.y", "v.z", "tag", "raw"]
936 );
937 assert_eq!(index.info, [("sys_name".to_string(), "PX4".to_string())]);
938 assert_eq!(index.parameters[0].0, "MPC_XY_VEL");
939 assert_eq!(index.logged[0].3, "low battery");
940 assert_eq!(index.damaged, 1);
941 assert!(index.cut_short);
942 }
943
944 #[test]
945 fn a_topic_decodes() {
946 let data = tiny();
947 let index = index(&data).unwrap();
948 let topic = &index.topics[&1];
949 let records = Arc::new(
950 IndexedRecords::new(
951 Arc::new(Bytes::Owned(data.clone())),
952 topic.offsets.clone(),
953 topic.columns.clone(),
954 )
955 .unwrap(),
956 );
957 let df = records.lazy().collect().unwrap();
958 assert_eq!(
959 df.column("v.y").unwrap().f32().unwrap().to_vec(),
960 [Some(2.0), Some(4.0)]
961 );
962 assert_eq!(
963 df.column("timestamp").unwrap().dtype(),
964 &DataType::Duration(TimeUnit::Microseconds)
965 );
966 assert_eq!(df.column("tag").unwrap().str().unwrap().get(0), Some("ab"));
967 assert_eq!(
968 df.column("raw").unwrap().dtype(),
969 &DataType::Array(Box::new(DataType::Int16), 2)
970 );
971 }
972
973 #[test]
974 fn garbage_is_bounded() {
975 assert!(index(b"nope").is_err());
976 let mut data = MAGIC.to_vec();
977 data.extend([1, 0, 0, 0, 0, 0, 0, 0, 0]);
978 data.extend([0xFF; 64]);
979 let index = index(&data).unwrap();
980 assert!(index.names.is_empty());
981 }
982}