1use std::collections::{BTreeMap, HashMap};
16use std::sync::Arc;
17
18use polars::prelude::*;
19
20use crate::formats::columns::{Builder, Cell, Kind};
21use crate::formats::fixed_records::{Bytes, ColumnLayout, Logical, Physical};
22use crate::formats::indexed::{IndexedRecords, Offsets};
23use crate::formats::model_files::MetaValue;
24use crate::formats::sqlite::Table;
25use crate::formats::text_formats::Detail;
26
27pub(crate) const READER: crate::formats::readers::Reader = crate::formats::readers::Reader {
29 scan: crate::formats::indexed::scan::<Index>,
30 signatures: &[crate::formats::readers::Signature {
31 says: |head, _| looks_like(head),
32 kind: crate::formats::readers::Kind::Magic,
33 trusted: crate::formats::readers::Trusted {
34 tables: true,
35 ..crate::formats::readers::EVERYWHERE
36 },
37 }],
38 tables: Some(crate::formats::indexed::listed::<Index>),
39 ..crate::formats::readers::BASE
40};
41
42pub const MAGIC: &[u8; 7] = b"ULog\x01\x12\x35";
44
45pub(crate) const SYNC: [u8; 8] = [0x2F, 0x73, 0x13, 0x20, 0x25, 0x0C, 0xBB, 0x12];
47
48pub const MAX_COLUMNS: usize = 4096;
50const MAX_DEPTH: usize = 16;
52const MAX_DEFINITIONS: usize = 65_536;
54const MAX_LOGGED: usize = 1_000_000;
56
57pub const LOGGED: &str = "logged_messages";
59pub const PARAMETERS: &str = "parameters";
60
61pub fn looks_like(head: &[u8]) -> bool {
63 head.starts_with(MAGIC)
64}
65
66#[derive(Debug, Clone, PartialEq)]
68struct FieldDef {
69 kind: String,
70 count: Option<usize>,
72 name: String,
73}
74
75#[derive(Debug)]
77pub struct Topic {
78 pub offsets: Arc<Offsets>,
80 pub columns: Vec<ColumnLayout>,
81}
82
83#[derive(Debug, Default)]
85pub struct Index {
86 pub version: u8,
87 pub topics: BTreeMap<u16, Topic>,
89 pub names: Vec<(String, u16)>,
91 pub info: Vec<(String, String)>,
92 pub parameters: Vec<(String, &'static str, f64, Option<u64>)>,
95 pub logged: Vec<(u64, &'static str, Option<u16>, String)>,
97 pub logged_left_out: usize,
98 pub dropouts: usize,
99 pub dropout_ms: u64,
100 pub skipped: usize,
102 pub damaged: usize,
103 pub short: usize,
105 pub unsubscribed: usize,
107 pub past_limit: usize,
109 pub unread: Vec<(String, String)>,
111 pub cut_short: bool,
113}
114
115fn primitive(kind: &str) -> Option<(Physical, usize)> {
117 Some(match kind {
118 "int8_t" => (Physical::Signed(1), 1),
119 "uint8_t" => (Physical::Unsigned(1), 1),
120 "int16_t" => (Physical::Signed(2), 2),
121 "uint16_t" => (Physical::Unsigned(2), 2),
122 "int32_t" => (Physical::Signed(4), 4),
123 "uint32_t" => (Physical::Unsigned(4), 4),
124 "int64_t" => (Physical::Signed(8), 8),
125 "uint64_t" => (Physical::Unsigned(8), 8),
126 "float" => (Physical::Float(4), 4),
127 "double" => (Physical::Float(8), 8),
128 "bool" => (Physical::Bool, 1),
129 "char" => (Physical::Text, 1),
130 _ => return None,
131 })
132}
133
134fn parse_format(text: &str) -> Option<(String, Vec<FieldDef>)> {
136 let (name, rest) = text.split_once(':')?;
137 let mut fields = Vec::new();
138 for part in rest.split(';').filter(|p| !p.trim().is_empty()) {
139 let (kind, field) = part.trim().split_once(' ')?;
140 let (kind, count) = match kind.split_once('[') {
141 Some((kind, n)) => (kind, Some(n.strip_suffix(']')?.parse().ok()?)),
142 None => (kind, None),
143 };
144 fields.push(FieldDef {
145 kind: kind.to_string(),
146 count,
147 name: field.trim().to_string(),
148 });
149 if fields.len() > MAX_COLUMNS {
150 return None;
151 }
152 }
153 Some((name.to_string(), fields))
154}
155
156fn size_of(
158 name: &str,
159 formats: &HashMap<String, Vec<FieldDef>>,
160 depth: usize,
161) -> Result<usize, String> {
162 if depth > MAX_DEPTH {
163 return Err("types nest too deep".into());
164 }
165 let fields = formats
166 .get(name)
167 .ok_or_else(|| format!("no format for {name}"))?;
168 let mut size = 0usize;
169 for f in fields {
170 let one = match primitive(&f.kind) {
171 Some((_, n)) => n,
172 None => size_of(&f.kind, formats, depth + 1)?,
173 };
174 size = one
175 .checked_mul(f.count.unwrap_or(1))
176 .and_then(|n| size.checked_add(n))
177 .ok_or("a type is too large")?;
178 }
179 Ok(size)
180}
181
182fn columns_of(
185 name: &str,
186 formats: &HashMap<String, Vec<FieldDef>>,
187 prefix: &str,
188 offset: usize,
189 out: &mut Vec<ColumnLayout>,
190 depth: usize,
191) -> Result<usize, String> {
192 if depth > MAX_DEPTH {
193 return Err("types nest too deep".into());
194 }
195 let fields = formats
196 .get(name)
197 .ok_or_else(|| format!("no format for {name}"))?;
198 let mut at = offset;
199 let mut needs = offset;
200 for f in fields {
201 let count = f.count.unwrap_or(1);
202 let full = format!("{prefix}{}", f.name);
203 match primitive(&f.kind) {
204 Some((physical, width)) => {
205 let bytes = width.checked_mul(count).ok_or("a field is too large")?;
206 if !f.name.starts_with("_padding") && bytes > 0 {
208 let mut column = match physical {
209 Physical::Text => ColumnLayout::new(&full, at, 0, physical, bytes),
211 _ => {
212 let mut c = ColumnLayout::new(&full, at, 0, physical, width);
213 c.count = count;
214 c
215 }
216 };
217 if depth == 0 && f.name == "timestamp" && f.kind == "uint64_t" && count == 1 {
219 column.logical = Logical::Duration {
220 unit: TimeUnit::Microseconds,
221 multiplier: 1,
222 };
223 }
224 out.push(column);
225 needs = at + bytes;
226 }
227 at = at.checked_add(bytes).ok_or("a type is too large")?;
228 }
229 None => {
230 let size = size_of(&f.kind, formats, depth + 1)?;
231 for i in 0..count {
232 let inner = match f.count {
233 Some(_) => format!("{full}[{i}]."),
234 None => format!("{full}."),
235 };
236 let end = columns_of(&f.kind, formats, &inner, at, out, depth + 1)?;
237 if end > at {
238 needs = needs.max(end);
239 }
240 at = at.checked_add(size).ok_or("a type is too large")?;
241 if out.len() > MAX_COLUMNS {
242 return Err(format!("its fields make more than {MAX_COLUMNS} columns"));
243 }
244 }
245 }
246 }
247 if out.len() > MAX_COLUMNS {
248 return Err(format!("its fields make more than {MAX_COLUMNS} columns"));
249 }
250 }
251 Ok(needs)
252}
253
254fn u16_at(b: &[u8], at: usize) -> Option<u16> {
255 Some(u16::from_le_bytes(b.get(at..at + 2)?.try_into().ok()?))
256}
257
258fn u64_at(b: &[u8], at: usize) -> Option<u64> {
259 Some(u64::from_le_bytes(b.get(at..at + 8)?.try_into().ok()?))
260}
261
262fn key_value(payload: &[u8]) -> Option<(String, String, String, Option<f64>)> {
264 let len = *payload.first()? as usize;
265 let key = std::str::from_utf8(payload.get(1..1 + len)?).ok()?;
266 let value = payload.get(1 + len..)?;
267 let (kind, name) = key.split_once(' ')?;
268 let (base, count) = match kind.split_once('[') {
269 Some((base, n)) => (base, n.trim_end_matches(']').parse::<usize>().ok()),
270 None => (kind, None),
271 };
272 let number = |v: Option<f64>| v.map(|n| (format_number(n), Some(n)));
273 let (text, number) = match (base, count) {
274 ("char", _) => (
275 String::from_utf8_lossy(value)
276 .trim_end_matches('\0')
277 .to_string(),
278 None,
279 ),
280 (_, Some(_)) => (crate::formats::fixed_records::hex(value), None),
281 _ => {
282 let (physical, width) = primitive(base)?;
283 let raw = value.get(..width)?;
284 let v = match physical {
285 Physical::Signed(_) => {
286 crate::formats::fixed_records::read_signed(raw, false) as f64
287 }
288 Physical::Unsigned(_) | Physical::Bool => {
289 crate::formats::fixed_records::read_unsigned(raw, false) as f64
290 }
291 Physical::Float(4) => f32::from_le_bytes(raw.try_into().ok()?) as f64,
292 Physical::Float(_) => f64::from_le_bytes(raw.try_into().ok()?),
293 _ => return None,
294 };
295 number(Some(v))?
296 }
297 };
298 Some((kind.to_string(), name.to_string(), text, number))
299}
300
301fn format_number(n: f64) -> String {
302 if n.fract() == 0.0 && n.abs() < 1e15 {
303 format!("{}", n as i64)
304 } else {
305 format!("{n}")
306 }
307}
308
309fn level(byte: u8) -> &'static str {
310 match byte {
311 b'0' => "emergency",
312 b'1' => "alert",
313 b'2' => "critical",
314 b'3' => "error",
315 b'4' => "warning",
316 b'5' => "notice",
317 b'6' => "info",
318 b'7' => "debug",
319 _ => "unknown",
320 }
321}
322
323pub fn index(data: &[u8]) -> Result<Index, String> {
325 if !looks_like(data) {
326 return Err("not a ULog file: no ULog magic at the start".into());
327 }
328 if data.len() < 16 {
329 return Err("the header is cut short".into());
330 }
331 let mut index = Index {
332 version: data[7],
333 ..Index::default()
334 };
335 let mut formats: HashMap<String, Vec<FieldDef>> = HashMap::new();
336 let mut subs: HashMap<u16, (String, u8, Offsets)> = HashMap::new();
338 let mut multi: Vec<(String, String)> = Vec::new();
339 let mut records = 0usize;
340 let limit = crate::limits::get().indexed_records;
341 let mut last_time: Option<u64> = None;
342 let mut appended: Vec<usize> = Vec::new();
344 let mut end = data.len();
345 let mut at = 16usize;
346 loop {
347 if at >= end {
348 match appended.first().copied() {
349 Some(next) if next >= at.min(end) && next < data.len() => {
350 appended.remove(0);
351 at = next;
352 end = appended.first().copied().unwrap_or(data.len()).max(next);
353 continue;
354 }
355 _ => break,
356 }
357 }
358 let Some(size) = u16_at(data, at) else {
359 index.cut_short = at < end;
360 break;
361 };
362 let size = size as usize;
363 let Some(&kind) = data.get(at + 2) else {
364 index.cut_short = true;
365 break;
366 };
367 let body = at + 3;
368 if body + size > end {
369 if let Some(next) = find_sync(data, at + 1, end) {
371 index.skipped += next - at;
372 index.damaged += 1;
373 at = next;
374 continue;
375 }
376 index.cut_short = true;
377 break;
378 }
379 let payload = &data[body..body + size];
380 match kind {
381 b'B' if size >= 40 => {
382 if payload[8] & 1 != 0 {
384 for i in 0..3 {
385 if let Some(o) = u64_at(payload, 16 + i * 8)
386 && o > 0
387 && let Ok(o) = usize::try_from(o)
388 && o > body + size
389 && o < data.len()
390 {
391 appended.push(o);
392 }
393 }
394 appended.sort_unstable();
395 if let Some(&first) = appended.first() {
396 end = first;
397 }
398 }
399 }
400 b'F' => {
401 if formats.len() < MAX_DEFINITIONS
402 && let Some((name, fields)) = parse_format(&String::from_utf8_lossy(payload))
403 {
404 formats.insert(name, fields);
405 }
406 }
407 b'A' if size >= 3 => {
408 let multi_id = payload[0];
409 let id = u16_at(payload, 1).unwrap_or_default();
410 let name = String::from_utf8_lossy(&payload[3..])
411 .trim_end_matches('\0')
412 .to_string();
413 if subs.len() < MAX_DEFINITIONS {
414 subs.insert(id, (name, multi_id, Offsets::for_file(data.len())));
415 }
416 }
417 b'D' if size >= 2 => {
418 let id = u16_at(payload, 0).unwrap_or_default();
419 match subs.get_mut(&id) {
420 Some((_, _, offsets)) if records < limit => {
421 offsets.push(body + 2);
422 records += 1;
423 }
424 Some(_) => index.past_limit += 1,
425 None => index.unsubscribed += 1,
426 }
427 }
428 b'I' => {
429 if let Some((_, name, text, _)) = key_value(payload)
430 && index.info.len() < MAX_DEFINITIONS
431 {
432 index.info.push((name, text));
433 }
434 }
435 b'M' if size >= 1 => {
436 if let Some((_, name, text, _)) = key_value(&payload[1..]) {
438 let continued = payload[0] != 0;
439 match multi.iter().position(|(n, _)| *n == name) {
440 Some(i) if continued => {
441 let value = &mut multi[i].1;
442 if value.len() < 1 << 20 {
443 value.push_str(&text);
444 }
445 }
446 _ if multi.len() < MAX_DEFINITIONS => multi.push((name, text)),
447 _ => {}
448 }
449 }
450 }
451 b'P' => {
452 if let Some((kind, name, _, Some(value))) = key_value(payload)
453 && index.parameters.len() < MAX_DEFINITIONS
454 {
455 let kind = if kind == "float" { "float" } else { "int32" };
456 index.parameters.push((name, kind, value, last_time));
457 }
458 }
459 b'L' if size >= 9 => {
460 let time = u64_at(payload, 1).unwrap_or_default();
461 last_time = Some(time);
462 push_logged(&mut index, time, payload[0], None, &payload[9..]);
463 }
464 b'C' if size >= 11 => {
465 let tag = u16_at(payload, 1);
466 let time = u64_at(payload, 3).unwrap_or_default();
467 last_time = Some(time);
468 push_logged(&mut index, time, payload[0], tag, &payload[11..]);
469 }
470 b'O' if size >= 2 => {
471 index.dropouts += 1;
472 index.dropout_ms += u64::from(u16_at(payload, 0).unwrap_or_default());
473 }
474 b'R' | b'S' | b'Q' | b'B' | b'A' | b'D' | b'L' | b'C' | b'O' | b'M' => {}
475 _ => {
476 match find_sync(data, at + 1, end) {
478 Some(next) => {
479 index.skipped += next - at;
480 index.damaged += 1;
481 at = next;
482 continue;
483 }
484 None => {
485 index.skipped += end - at;
486 index.damaged += 1;
487 at = end;
488 continue;
489 }
490 }
491 }
492 }
493 at = body + size;
494 }
495 for (name, value) in multi {
496 index.info.push((name, value));
497 }
498
499 let mut counts: HashMap<String, usize> = HashMap::new();
501 for (name, _, offsets) in subs.values() {
502 if !offsets.is_empty() {
503 *counts.entry(name.clone()).or_default() += 1;
504 }
505 }
506 let mut ids: Vec<u16> = subs.keys().copied().collect();
507 ids.sort_unstable_by_key(|id| (subs[id].0.clone(), subs[id].1));
508 for id in ids {
509 let (name, multi_id, mut offsets) = subs.remove(&id).expect("listed above");
510 if offsets.is_empty() {
511 continue;
512 }
513 let mut columns = Vec::new();
514 let needs = match columns_of(&name, &formats, "", 0, &mut columns, 0).and_then(|needs| {
515 let mut seen = std::collections::HashSet::new();
517 match columns.iter().find(|c| !seen.insert(c.name.clone())) {
518 Some(c) => Err(format!("it has two fields named {}", c.name)),
519 None => Ok(needs),
520 }
521 }) {
522 Ok(needs) => needs,
523 Err(why) => {
524 index.unread.push((name, why));
525 continue;
526 }
527 };
528 let before = offsets.len();
531 offsets = keep_long_enough(offsets, data, needs);
532 index.short += before - offsets.len();
533 if offsets.is_empty() {
534 continue;
535 }
536 offsets.shrink();
537 let table = if counts.get(&name).copied().unwrap_or(0) > 1 {
538 format!("{name}.{multi_id}")
539 } else {
540 name
541 };
542 index.names.push((table, id));
543 index.topics.insert(
544 id,
545 Topic {
546 offsets: Arc::new(offsets),
547 columns,
548 },
549 );
550 }
551 Ok(index)
552}
553
554fn keep_long_enough(offsets: Offsets, data: &[u8], needs: usize) -> Offsets {
557 let long_enough = |at: usize| {
558 u16_at(data, at - 5).is_some_and(|size| (size as usize).saturating_sub(2) >= needs)
559 };
560 let mut kept = Offsets::for_file(data.len());
561 for i in 0..offsets.len() {
562 let at = offsets.get(i);
563 if long_enough(at) {
564 kept.push(at);
565 }
566 }
567 kept
568}
569
570fn push_logged(index: &mut Index, time: u64, level_byte: u8, tag: Option<u16>, text: &[u8]) {
571 if index.logged.len() >= MAX_LOGGED {
572 index.logged_left_out += 1;
573 return;
574 }
575 let text = String::from_utf8_lossy(text)
576 .trim_end_matches('\0')
577 .to_string();
578 index.logged.push((time, level(level_byte), tag, text));
579}
580
581fn find_sync(data: &[u8], from: usize, end: usize) -> Option<usize> {
582 let hay = data.get(from..end)?;
583 memchr::memmem::find(hay, &SYNC).map(|i| from + i + SYNC.len())
584}
585
586impl crate::formats::indexed::Log for Index {
587 const EMPTY: &'static str = " The log has no data messages.";
588
589 fn index(data: &[u8]) -> Result<Self, String> {
590 index(data)
591 }
592
593 fn tables(&self) -> Vec<Table> {
594 let mut tables: Vec<Table> = self
595 .names
596 .iter()
597 .map(|(name, id)| {
598 let columns = self.topics[id].columns.iter().map(|c| c.name.as_str());
599 Table::plain(name, "topic", columns)
600 })
601 .collect();
602 let taken = |name: &str| self.names.iter().any(|(n, _)| n == name);
603 if !self.logged.is_empty() && !taken(LOGGED) {
604 let columns = ["timestamp", "level", "tag", "message"];
605 tables.push(Table::plain(LOGGED, "messages", columns));
606 }
607 if !self.parameters.is_empty() && !taken(PARAMETERS) {
608 let columns = ["name", "type", "value", "timestamp"];
609 tables.push(Table::plain(PARAMETERS, "parameters", columns));
610 }
611 tables
612 }
613
614 fn detail(&self) -> Detail {
615 let count = crate::formats::text_formats::count;
616 let mut lines = vec![
617 format!("Version: {}", self.version),
618 format!(
619 "Topics: {}",
620 count(self.names.len() as u64, "table", "tables")
621 ),
622 ];
623 if self.dropouts > 0 {
624 lines.push(format!(
625 "Dropouts: {}, {} ms in all",
626 self.dropouts, self.dropout_ms
627 ));
628 }
629 let mut list: Vec<(String, MetaValue)> = self
630 .info
631 .iter()
632 .map(|(k, v)| (k.clone(), MetaValue::Text(v.clone())))
633 .collect();
634 let mut seen = std::collections::HashSet::new();
635 for (name, _, value, at) in &self.parameters {
636 if at.is_none() && seen.insert(name.clone()) {
638 list.push((
639 format!("param {name}"),
640 MetaValue::Text(format_number(*value)),
641 ));
642 }
643 }
644 let total = list.len();
645 Detail {
646 tab: crate::formats::text_formats::tab(crate::FileFormat::Ulog),
647 lines,
648 list_title: "Info and parameters",
649 list: crate::formats::text_formats::capped_list(list.into_iter(), total),
650 first: false,
651 ..Default::default()
652 }
653 }
654
655 fn notes(&self) -> Vec<String> {
656 let group = |n: usize| crate::numfmt::group_chrome(n);
657 let mut notes = Vec::new();
658 if self.damaged > 0 {
659 notes.push(format!(
660 "{} damaged stretches skipped ({} bytes)",
661 group(self.damaged),
662 group(self.skipped)
663 ));
664 }
665 if self.cut_short {
666 notes.push("log cut short mid-message".to_string());
667 }
668 if self.short > 0 {
669 notes.push(format!(
670 "{} data messages left out: too short for their topic",
671 group(self.short)
672 ));
673 }
674 if self.unsubscribed > 0 {
675 notes.push(format!(
676 "{} data messages left out: no subscription",
677 group(self.unsubscribed)
678 ));
679 }
680 if self.past_limit > 0 {
681 notes.push(crate::limits::left_out(
682 &format!("{} messages", group(self.past_limit)),
683 crate::limits::get().indexed_records,
684 "indexed_records",
685 ));
686 }
687 if self.logged_left_out > 0 {
688 notes.push(format!(
689 "{} logged messages left out: past the first {}",
690 group(self.logged_left_out),
691 group(MAX_LOGGED)
692 ));
693 }
694 for (topic, why) in &self.unread {
695 notes.push(format!("topic {topic} not read: {why}"));
696 }
697 notes
698 }
699
700 fn table(
701 &self,
702 bytes: Arc<Bytes>,
703 name: &str,
704 opened: &mut crate::formats::members::Opened,
705 ) -> Result<LazyFrame, String> {
706 if let Some((_, id)) = self.names.iter().find(|(n, _)| *n == name) {
707 let topic = &self.topics[id];
708 let records = Arc::new(
709 IndexedRecords::new(bytes, topic.offsets.clone(), topic.columns.clone())
710 .map_err(|e| format!("table \"{name}\": {e}"))?,
711 );
712 opened.window = Some((records.clone(), records.rows()));
713 return Ok(records.lazy());
714 }
715 let mut table = if name == LOGGED {
716 let mut table = Builder::new(&[
717 ("timestamp", Kind::DurationUs),
718 ("level", Kind::Label),
719 ("tag", Kind::U16),
720 ("message", Kind::Str),
721 ]);
722 for (t, level, tag, text) in &self.logged {
723 table.push([
724 Cell::DurationUs(Some(*t as i64)),
725 Cell::Label(Some(level)),
726 Cell::U16(*tag),
727 Cell::Str(Some(text.clone())),
728 ]);
729 }
730 table
731 } else {
732 let mut table = Builder::new(&[
733 ("name", Kind::Str),
734 ("type", Kind::Label),
735 ("value", Kind::F64),
736 ("timestamp", Kind::DurationUs),
737 ]);
738 for (name, kind, value, t) in &self.parameters {
739 table.push([
740 Cell::Str(Some(name.clone())),
741 Cell::Label(Some(kind)),
742 Cell::F64(Some(*value)),
743 Cell::DurationUs(t.map(|t| t as i64)),
744 ]);
745 }
746 table
747 };
748 Ok(table.take().map_err(|e| e.to_string())?.lazy())
749 }
750}
751
752#[cfg(test)]
753mod tests {
754 use super::*;
755
756 #[test]
758 fn errors_name_the_file() {
759 crate::formats::readers::bad_input::each_names_its_file(
760 crate::FileFormat::Ulog,
761 &[
762 ("text.ulg", b"hello there, this is text", "Not a ULog file"),
763 ("cut.ulg", b"ULog\x01\x12\x35\x01", "cut short"),
764 ],
765 );
766 }
767
768 #[test]
769 fn topics_instances_and_text() {
770 let index = index(&crate::tests::fixtures::ulog()).unwrap();
771 let names: Vec<&str> = index.names.iter().map(|(n, _)| n.as_str()).collect();
772 assert_eq!(names, ["sensor.0", "sensor.1", "status"]);
773 let sensor = &index.topics[&1];
774 assert_eq!(sensor.offsets.len(), 2);
775 let columns: Vec<&str> = sensor.columns.iter().map(|c| c.name.as_str()).collect();
776 assert_eq!(
777 columns,
778 ["timestamp", "id", "v.x", "v.y", "v.z", "tag", "raw"]
779 );
780 assert_eq!(index.info, [("sys_name".to_string(), "PX4".to_string())]);
781 assert_eq!(index.parameters[0].0, "MPC_XY_VEL");
782 assert_eq!(index.logged[0].3, "low battery");
783 assert_eq!(index.damaged, 1);
784 assert!(index.cut_short);
785 }
786
787 #[test]
788 fn a_topic_decodes() {
789 let data = crate::tests::fixtures::ulog();
790 let index = index(&data).unwrap();
791 let topic = &index.topics[&1];
792 let records = Arc::new(
793 IndexedRecords::new(
794 Arc::new(Bytes::Owned(data.clone())),
795 topic.offsets.clone(),
796 topic.columns.clone(),
797 )
798 .unwrap(),
799 );
800 let df = records.lazy().collect().unwrap();
801 assert_eq!(
802 df.column("v.y").unwrap().f32().unwrap().to_vec(),
803 [Some(2.0), Some(4.0)]
804 );
805 assert_eq!(
806 df.column("timestamp").unwrap().dtype(),
807 &DataType::Duration(TimeUnit::Microseconds)
808 );
809 assert_eq!(df.column("tag").unwrap().str().unwrap().get(0), Some("ab"));
810 assert_eq!(
811 df.column("raw").unwrap().dtype(),
812 &DataType::Array(Box::new(DataType::Int16), 2)
813 );
814 }
815
816 #[test]
817 fn garbage_is_bounded() {
818 assert!(index(b"nope").is_err());
819 let mut data = MAGIC.to_vec();
820 data.extend([1, 0, 0, 0, 0, 0, 0, 0, 0]);
821 data.extend([0xFF; 64]);
822 let index = index(&data).unwrap();
823 assert!(index.names.is_empty());
824 }
825}