use std::collections::{BTreeMap, HashMap};
use std::path::Path;
use std::sync::Arc;
use color_eyre::Result;
use color_eyre::eyre::eyre;
use crate::error_display::{FileError, in_file};
use polars::prelude::*;
use crate::fixed_records::{Bytes, ColumnLayout, Logical, Physical};
use crate::indexed::{IndexedRecords, Offsets};
use crate::model_files::MetaValue;
use crate::sqlite::Table;
use crate::text_formats::Detail;
pub(crate) const READER: crate::readers::Reader = crate::readers::Reader {
scan,
signatures: &[crate::readers::Signature {
says: |head, _| looks_like(head),
kind: crate::readers::Kind::Magic,
trusted: crate::readers::Trusted {
tables: true,
..crate::readers::EVERYWHERE
},
}],
tables: Some(listed),
..crate::readers::BASE
};
pub const MAGIC: &[u8; 7] = b"ULog\x01\x12\x35";
const SYNC: [u8; 8] = [0x2F, 0x73, 0x13, 0x20, 0x25, 0x0C, 0xBB, 0x12];
pub const MAX_COLUMNS: usize = 4096;
const MAX_DEPTH: usize = 16;
const MAX_DEFINITIONS: usize = 65_536;
const MAX_LOGGED: usize = 1_000_000;
pub const LOGGED: &str = "logged_messages";
pub const PARAMETERS: &str = "parameters";
pub fn looks_like(head: &[u8]) -> bool {
head.starts_with(MAGIC)
}
#[derive(Debug, Clone, PartialEq)]
struct FieldDef {
kind: String,
count: Option<usize>,
name: String,
}
#[derive(Debug)]
pub struct Topic {
pub name: String,
pub multi_id: u8,
pub offsets: Arc<Offsets>,
pub needs: usize,
pub columns: Vec<ColumnLayout>,
}
#[derive(Debug, Default)]
pub struct Index {
pub version: u8,
pub start_us: u64,
pub topics: BTreeMap<u16, Topic>,
pub names: Vec<(String, u16)>,
pub info: Vec<(String, String)>,
pub parameters: Vec<(String, &'static str, f64, Option<u64>)>,
pub logged: Vec<(u64, &'static str, Option<u16>, String)>,
pub logged_left_out: usize,
pub dropouts: usize,
pub dropout_ms: u64,
pub skipped: usize,
pub damaged: usize,
pub short: usize,
pub unsubscribed: usize,
pub past_limit: usize,
pub unread: Vec<(String, String)>,
pub cut_short: bool,
}
fn primitive(kind: &str) -> Option<(Physical, usize)> {
Some(match kind {
"int8_t" => (Physical::Signed(1), 1),
"uint8_t" => (Physical::Unsigned(1), 1),
"int16_t" => (Physical::Signed(2), 2),
"uint16_t" => (Physical::Unsigned(2), 2),
"int32_t" => (Physical::Signed(4), 4),
"uint32_t" => (Physical::Unsigned(4), 4),
"int64_t" => (Physical::Signed(8), 8),
"uint64_t" => (Physical::Unsigned(8), 8),
"float" => (Physical::Float(4), 4),
"double" => (Physical::Float(8), 8),
"bool" => (Physical::Bool, 1),
"char" => (Physical::Text, 1),
_ => return None,
})
}
fn parse_format(text: &str) -> Option<(String, Vec<FieldDef>)> {
let (name, rest) = text.split_once(':')?;
let mut fields = Vec::new();
for part in rest.split(';').filter(|p| !p.trim().is_empty()) {
let (kind, field) = part.trim().split_once(' ')?;
let (kind, count) = match kind.split_once('[') {
Some((kind, n)) => (kind, Some(n.strip_suffix(']')?.parse().ok()?)),
None => (kind, None),
};
fields.push(FieldDef {
kind: kind.to_string(),
count,
name: field.trim().to_string(),
});
if fields.len() > MAX_COLUMNS {
return None;
}
}
Some((name.to_string(), fields))
}
fn size_of(
name: &str,
formats: &HashMap<String, Vec<FieldDef>>,
depth: usize,
) -> std::result::Result<usize, String> {
if depth > MAX_DEPTH {
return Err("types nest too deep".into());
}
let fields = formats
.get(name)
.ok_or_else(|| format!("no format for {name}"))?;
let mut size = 0usize;
for f in fields {
let one = match primitive(&f.kind) {
Some((_, n)) => n,
None => size_of(&f.kind, formats, depth + 1)?,
};
size = one
.checked_mul(f.count.unwrap_or(1))
.and_then(|n| size.checked_add(n))
.ok_or("a type is too large")?;
}
Ok(size)
}
fn columns_of(
name: &str,
formats: &HashMap<String, Vec<FieldDef>>,
prefix: &str,
offset: usize,
out: &mut Vec<ColumnLayout>,
depth: usize,
) -> std::result::Result<usize, String> {
if depth > MAX_DEPTH {
return Err("types nest too deep".into());
}
let fields = formats
.get(name)
.ok_or_else(|| format!("no format for {name}"))?;
let mut at = offset;
let mut needs = offset;
for f in fields {
let count = f.count.unwrap_or(1);
let full = format!("{prefix}{}", f.name);
match primitive(&f.kind) {
Some((physical, width)) => {
let bytes = width.checked_mul(count).ok_or("a field is too large")?;
if !f.name.starts_with("_padding") && bytes > 0 {
let mut column = match physical {
Physical::Text => ColumnLayout::new(&full, at, 0, physical, bytes),
_ => {
let mut c = ColumnLayout::new(&full, at, 0, physical, width);
c.count = count;
c
}
};
if depth == 0 && f.name == "timestamp" && f.kind == "uint64_t" && count == 1 {
column.logical = Logical::Duration {
unit: TimeUnit::Microseconds,
multiplier: 1,
};
}
out.push(column);
needs = at + bytes;
}
at = at.checked_add(bytes).ok_or("a type is too large")?;
}
None => {
let size = size_of(&f.kind, formats, depth + 1)?;
for i in 0..count {
let inner = match f.count {
Some(_) => format!("{full}[{i}]."),
None => format!("{full}."),
};
let end = columns_of(&f.kind, formats, &inner, at, out, depth + 1)?;
if end > at {
needs = needs.max(end);
}
at = at.checked_add(size).ok_or("a type is too large")?;
if out.len() > MAX_COLUMNS {
return Err(format!("its fields make more than {MAX_COLUMNS} columns"));
}
}
}
}
if out.len() > MAX_COLUMNS {
return Err(format!("its fields make more than {MAX_COLUMNS} columns"));
}
}
Ok(needs)
}
fn u16_at(b: &[u8], at: usize) -> Option<u16> {
Some(u16::from_le_bytes(b.get(at..at + 2)?.try_into().ok()?))
}
fn u64_at(b: &[u8], at: usize) -> Option<u64> {
Some(u64::from_le_bytes(b.get(at..at + 8)?.try_into().ok()?))
}
fn key_value(payload: &[u8]) -> Option<(String, String, String, Option<f64>)> {
let len = *payload.first()? as usize;
let key = std::str::from_utf8(payload.get(1..1 + len)?).ok()?;
let value = payload.get(1 + len..)?;
let (kind, name) = key.split_once(' ')?;
let (base, count) = match kind.split_once('[') {
Some((base, n)) => (base, n.trim_end_matches(']').parse::<usize>().ok()),
None => (kind, None),
};
let number = |v: Option<f64>| v.map(|n| (format_number(n), Some(n)));
let (text, number) = match (base, count) {
("char", _) => (
String::from_utf8_lossy(value)
.trim_end_matches('\0')
.to_string(),
None,
),
(_, Some(_)) => (crate::fixed_records::hex(value), None),
_ => {
let (physical, width) = primitive(base)?;
let raw = value.get(..width)?;
let v = match physical {
Physical::Signed(_) => crate::fixed_records::read_signed(raw, false) as f64,
Physical::Unsigned(_) | Physical::Bool => {
crate::fixed_records::read_unsigned(raw, false) as f64
}
Physical::Float(4) => f32::from_le_bytes(raw.try_into().ok()?) as f64,
Physical::Float(_) => f64::from_le_bytes(raw.try_into().ok()?),
_ => return None,
};
number(Some(v))?
}
};
Some((kind.to_string(), name.to_string(), text, number))
}
fn format_number(n: f64) -> String {
if n.fract() == 0.0 && n.abs() < 1e15 {
format!("{}", n as i64)
} else {
format!("{n}")
}
}
fn level(byte: u8) -> &'static str {
match byte {
b'0' => "emergency",
b'1' => "alert",
b'2' => "critical",
b'3' => "error",
b'4' => "warning",
b'5' => "notice",
b'6' => "info",
b'7' => "debug",
_ => "unknown",
}
}
pub fn index(data: &[u8]) -> std::result::Result<Index, String> {
if !looks_like(data) {
return Err("not a ULog file: no ULog magic at the start".into());
}
if data.len() < 16 {
return Err("the header is cut short".into());
}
let mut index = Index {
version: data[7],
start_us: u64_at(data, 8).unwrap_or(0),
..Index::default()
};
let mut formats: HashMap<String, Vec<FieldDef>> = HashMap::new();
let mut subs: HashMap<u16, (String, u8, Offsets)> = HashMap::new();
let mut multi: Vec<(String, String)> = Vec::new();
let mut records = 0usize;
let mut last_time: Option<u64> = None;
let mut appended: Vec<usize> = Vec::new();
let mut end = data.len();
let mut at = 16usize;
loop {
if at >= end {
match appended.first().copied() {
Some(next) if next >= at.min(end) && next < data.len() => {
appended.remove(0);
at = next;
end = appended.first().copied().unwrap_or(data.len()).max(next);
continue;
}
_ => break,
}
}
let Some(size) = u16_at(data, at) else {
index.cut_short = at < end;
break;
};
let size = size as usize;
let Some(&kind) = data.get(at + 2) else {
index.cut_short = true;
break;
};
let body = at + 3;
if body + size > end {
if let Some(next) = find_sync(data, at + 1, end) {
index.skipped += next - at;
index.damaged += 1;
at = next;
continue;
}
index.cut_short = true;
break;
}
let payload = &data[body..body + size];
match kind {
b'B' if size >= 40 => {
if payload[8] & 1 != 0 {
for i in 0..3 {
if let Some(o) = u64_at(payload, 16 + i * 8)
&& o > 0
&& let Ok(o) = usize::try_from(o)
&& o > body + size
&& o < data.len()
{
appended.push(o);
}
}
appended.sort_unstable();
if let Some(&first) = appended.first() {
end = first;
}
}
}
b'F' => {
if formats.len() < MAX_DEFINITIONS
&& let Some((name, fields)) = parse_format(&String::from_utf8_lossy(payload))
{
formats.insert(name, fields);
}
}
b'A' if size >= 3 => {
let multi_id = payload[0];
let id = u16_at(payload, 1).unwrap_or_default();
let name = String::from_utf8_lossy(&payload[3..])
.trim_end_matches('\0')
.to_string();
if subs.len() < MAX_DEFINITIONS {
subs.insert(id, (name, multi_id, Offsets::for_file(data.len())));
}
}
b'D' if size >= 2 => {
let id = u16_at(payload, 0).unwrap_or_default();
match subs.get_mut(&id) {
Some((_, _, offsets)) if records < crate::indexed::MAX_RECORDS => {
offsets.push(body + 2);
records += 1;
}
Some(_) => index.past_limit += 1,
None => index.unsubscribed += 1,
}
}
b'I' => {
if let Some((_, name, text, _)) = key_value(payload)
&& index.info.len() < MAX_DEFINITIONS
{
index.info.push((name, text));
}
}
b'M' if size >= 1 => {
if let Some((_, name, text, _)) = key_value(&payload[1..]) {
let continued = payload[0] != 0;
match multi.iter().position(|(n, _)| *n == name) {
Some(i) if continued => {
let value = &mut multi[i].1;
if value.len() < 1 << 20 {
value.push_str(&text);
}
}
_ if multi.len() < MAX_DEFINITIONS => multi.push((name, text)),
_ => {}
}
}
}
b'P' => {
if let Some((kind, name, _, Some(value))) = key_value(payload)
&& index.parameters.len() < MAX_DEFINITIONS
{
let kind = if kind == "float" { "float" } else { "int32" };
index.parameters.push((name, kind, value, last_time));
}
}
b'L' if size >= 9 => {
let time = u64_at(payload, 1).unwrap_or_default();
last_time = Some(time);
push_logged(&mut index, time, payload[0], None, &payload[9..]);
}
b'C' if size >= 11 => {
let tag = u16_at(payload, 1);
let time = u64_at(payload, 3).unwrap_or_default();
last_time = Some(time);
push_logged(&mut index, time, payload[0], tag, &payload[11..]);
}
b'O' if size >= 2 => {
index.dropouts += 1;
index.dropout_ms += u64::from(u16_at(payload, 0).unwrap_or_default());
}
b'R' | b'S' | b'Q' | b'B' | b'A' | b'D' | b'L' | b'C' | b'O' | b'M' => {}
_ => {
match find_sync(data, at + 1, end) {
Some(next) => {
index.skipped += next - at;
index.damaged += 1;
at = next;
continue;
}
None => {
index.skipped += end - at;
index.damaged += 1;
at = end;
continue;
}
}
}
}
at = body + size;
}
for (name, value) in multi {
index.info.push((name, value));
}
let mut counts: HashMap<String, usize> = HashMap::new();
for (name, _, offsets) in subs.values() {
if !offsets.is_empty() {
*counts.entry(name.clone()).or_default() += 1;
}
}
let mut ids: Vec<u16> = subs.keys().copied().collect();
ids.sort_unstable_by_key(|id| (subs[id].0.clone(), subs[id].1));
for id in ids {
let (name, multi_id, mut offsets) = subs.remove(&id).expect("listed above");
if offsets.is_empty() {
continue;
}
let mut columns = Vec::new();
let needs = match columns_of(&name, &formats, "", 0, &mut columns, 0).and_then(|needs| {
let mut seen = std::collections::HashSet::new();
match columns.iter().find(|c| !seen.insert(c.name.clone())) {
Some(c) => Err(format!("it has two fields named {}", c.name)),
None => Ok(needs),
}
}) {
Ok(needs) => needs,
Err(why) => {
index.unread.push((name, why));
continue;
}
};
let before = offsets.len();
offsets = keep_long_enough(offsets, data, needs);
index.short += before - offsets.len();
if offsets.is_empty() {
continue;
}
offsets.shrink();
let table = if counts.get(&name).copied().unwrap_or(0) > 1 {
format!("{name}.{multi_id}")
} else {
name.clone()
};
index.names.push((table, id));
index.topics.insert(
id,
Topic {
name,
multi_id,
offsets: Arc::new(offsets),
needs,
columns,
},
);
}
Ok(index)
}
fn keep_long_enough(offsets: Offsets, data: &[u8], needs: usize) -> Offsets {
let long_enough = |at: usize| {
u16_at(data, at - 5).is_some_and(|size| (size as usize).saturating_sub(2) >= needs)
};
let mut kept = Offsets::for_file(data.len());
for i in 0..offsets.len() {
let at = offsets.get(i);
if long_enough(at) {
kept.push(at);
}
}
kept
}
fn push_logged(index: &mut Index, time: u64, level_byte: u8, tag: Option<u16>, text: &[u8]) {
if index.logged.len() >= MAX_LOGGED {
index.logged_left_out += 1;
return;
}
let text = String::from_utf8_lossy(text)
.trim_end_matches('\0')
.to_string();
index.logged.push((time, level(level_byte), tag, text));
}
fn find_sync(data: &[u8], from: usize, end: usize) -> Option<usize> {
let hay = data.get(from..end)?;
memchr::memmem::find(hay, &SYNC).map(|i| from + i + SYNC.len())
}
pub fn listed(file: &Path) -> Result<Vec<Table>> {
crate::indexed::peek::<Index>(file)
.map(|index| tables(&index))
.ok_or_else(|| eyre!("Open the log to list its tables."))
}
pub fn tables(index: &Index) -> Vec<Table> {
let mut tables: Vec<Table> = index
.names
.iter()
.map(|(name, id)| Table {
name: name.clone(),
kind: "topic".to_string(),
internal: false,
columns: index.topics[id]
.columns
.iter()
.map(|c| (c.name.to_string(), String::new()))
.collect(),
})
.collect();
let taken = |name: &str| index.names.iter().any(|(n, _)| n == name);
if !index.logged.is_empty() && !taken(LOGGED) {
tables.push(Table {
name: LOGGED.to_string(),
kind: "messages".to_string(),
internal: false,
columns: ["timestamp", "level", "tag", "message"]
.iter()
.map(|c| (c.to_string(), String::new()))
.collect(),
});
}
if !index.parameters.is_empty() && !taken(PARAMETERS) {
tables.push(Table {
name: PARAMETERS.to_string(),
kind: "parameters".to_string(),
internal: false,
columns: ["name", "type", "value", "timestamp"]
.iter()
.map(|c| (c.to_string(), String::new()))
.collect(),
});
}
tables
}
pub fn detail(index: &Index) -> Detail {
let count = crate::text_formats::count;
let mut lines = vec![
format!("Version: {}", index.version),
format!(
"Topics: {}",
count(index.names.len() as u64, "table", "tables")
),
];
if index.dropouts > 0 {
lines.push(format!(
"Dropouts: {}, {} ms in all",
index.dropouts, index.dropout_ms
));
}
let mut list: Vec<(String, MetaValue)> = index
.info
.iter()
.map(|(k, v)| (k.clone(), MetaValue::Text(v.clone())))
.collect();
let mut seen = std::collections::HashSet::new();
for (name, _, value, at) in &index.parameters {
if at.is_none() && seen.insert(name.clone()) {
list.push((
format!("param {name}"),
MetaValue::Text(format_number(*value)),
));
}
}
let total = list.len();
Detail {
tab: crate::text_formats::tab(crate::FileFormat::Ulog),
lines,
list_title: "Info and parameters",
list: crate::text_formats::capped_list(list.into_iter(), total),
first: false,
..Default::default()
}
}
pub fn notes(index: &Index) -> Vec<String> {
let group = |n: usize| crate::numfmt::group_chrome(n);
let mut notes = Vec::new();
if index.damaged > 0 {
notes.push(format!(
"{} damaged stretches skipped ({} bytes)",
group(index.damaged),
group(index.skipped)
));
}
if index.cut_short {
notes.push("log cut short mid-message".to_string());
}
if index.short > 0 {
notes.push(format!(
"{} data messages left out: too short for their topic",
group(index.short)
));
}
if index.unsubscribed > 0 {
notes.push(format!(
"{} data messages left out: no subscription",
group(index.unsubscribed)
));
}
if index.past_limit > 0 {
notes.push(format!(
"{} messages left out: past the first {}",
group(index.past_limit),
group(crate::indexed::MAX_RECORDS)
));
}
if index.logged_left_out > 0 {
notes.push(format!(
"{} logged messages left out: past the first {}",
group(index.logged_left_out),
group(MAX_LOGGED)
));
}
for (topic, why) in &index.unread {
notes.push(format!("topic {topic} not read: {why}"));
}
notes
}
pub fn indexed(path: &Path) -> Result<(Arc<Bytes>, Arc<Index>)> {
let bytes = Arc::new(Bytes::map(path).map_err(|e| in_file(path, e.into()))?);
let index = crate::indexed::cached(path, || index(bytes.as_slice()))
.map_err(|e| FileError::new(path, e))?;
Ok((bytes, index))
}
pub enum Open {
Table {
lf: Box<LazyFrame>,
opened: Box<crate::members::Opened>,
},
Several(Vec<String>),
}
pub fn open(path: &Path, wanted: Option<&str>) -> Result<Open> {
let (bytes, index) = indexed(path)?;
let tables = tables(&index);
let picked = match crate::members::pick(
tables.clone(),
wanted,
path,
" The log has no data messages.",
)? {
crate::sqlite::Pick::One(table) => table.name,
crate::sqlite::Pick::Several(tables) => {
return Ok(Open::Several(tables.into_iter().map(|t| t.name).collect()));
}
};
let mut opened = crate::members::Opened {
detail: Some(Arc::new(detail(&index))),
other_tables: crate::members::others(&tables, &picked),
notes: notes(&index)
.into_iter()
.map(|n| crate::text_formats::note(n, "the log".to_string()))
.collect(),
..Default::default()
};
let lf = if let Some((_, id)) = index.names.iter().find(|(n, _)| *n == picked) {
let topic = &index.topics[id];
let records = Arc::new(
IndexedRecords::new(bytes, topic.offsets.clone(), topic.columns.clone())
.map_err(|e| FileError::new(path, format!("table \"{picked}\": {e}")))?,
);
opened.window = Some((records.clone(), records.rows()));
records.lazy()
} else if picked == LOGGED {
let (mut time, mut lvl, mut tag, mut text) =
(Vec::new(), Vec::new(), Vec::new(), Vec::new());
for (t, l, g, m) in &index.logged {
time.push(*t as i64);
lvl.push(*l);
tag.push(*g);
text.push(m.as_str());
}
df!(
"timestamp" => Int64Chunked::from_vec("timestamp".into(), time).into_duration(TimeUnit::Microseconds).into_series(),
"level" => lvl,
"tag" => tag,
"message" => text,
)?
.lazy()
} else {
let (mut name, mut kind, mut value, mut time) =
(Vec::new(), Vec::new(), Vec::new(), Vec::new());
for (n, k, v, t) in &index.parameters {
name.push(n.as_str());
kind.push(*k);
value.push(*v);
time.push(t.map(|t| t as i64));
}
df!(
"name" => name,
"type" => kind,
"value" => value,
"timestamp" => time.into_iter().collect::<Int64Chunked>().into_duration(TimeUnit::Microseconds).into_series(),
)?
.lazy()
};
Ok(Open::Table {
lf: Box::new(lf),
opened: Box::new(opened),
})
}
fn scan(input: crate::readers::ScanIn<'_>) -> Result<crate::scan::Scan> {
let file = input.path();
Ok(match open(file, input.options.table.as_deref())? {
Open::Table { lf, opened } => {
input.report.opened = Some(Arc::new(*opened));
(*lf).into()
}
Open::Several(tables) => crate::scan::Scan::Tables {
file: file.to_path_buf(),
tables,
format: input.format,
},
})
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
#[test]
fn errors_name_the_file() {
crate::readers::bad_input::each_names_its_file(
crate::FileFormat::Ulog,
&[
("text.ulg", b"hello there, this is text", "Not a ULog file"),
("cut.ulg", b"ULog\x01\x12\x35\x01", "cut short"),
],
);
}
fn message(kind: u8, payload: &[u8]) -> Vec<u8> {
let mut out = (payload.len() as u16).to_le_bytes().to_vec();
out.push(kind);
out.extend(payload);
out
}
pub(crate) fn tiny() -> Vec<u8> {
let mut log = MAGIC.to_vec();
log.push(1);
log.extend(1_000u64.to_le_bytes());
log.extend(message(b'F', b"vec3:float x;float y;float z;"));
log.extend(message(
b'F',
b"sensor:uint64_t timestamp;uint8_t id;vec3 v;char[4] tag;int16_t[2] raw;uint8_t[3] _padding0;",
));
log.extend(message(b'F', b"status:uint64_t timestamp;bool armed;"));
let mut info = vec![b"char[6] sys_name".len() as u8];
info.extend(b"char[6] sys_name");
info.extend(b"PX4\0\0\0");
log.extend(message(b'I', &info));
let mut param = vec![b"float MPC_XY_VEL".len() as u8];
param.extend(b"float MPC_XY_VEL");
param.extend(2.5f32.to_le_bytes());
log.extend(message(b'P', ¶m));
for (multi, id, name) in [(0u8, 1u16, "sensor"), (1, 2, "sensor"), (0, 3, "status")] {
let mut a = vec![multi];
a.extend(id.to_le_bytes());
a.extend(name.as_bytes());
log.extend(message(b'A', &a));
}
let sensor = |id: u16, t: u64, x: f32, pad: bool| {
let mut d = id.to_le_bytes().to_vec();
d.extend(t.to_le_bytes());
d.push(id as u8);
for v in [x, x * 2.0, x * 3.0] {
d.extend(v.to_le_bytes());
}
d.extend(b"ab\0\0");
d.extend(7i16.to_le_bytes());
d.extend((-7i16).to_le_bytes());
if pad {
d.extend([0u8; 3]);
}
message(b'D', &d)
};
log.extend(sensor(1, 100, 1.0, true));
log.extend(sensor(2, 110, 5.0, false));
let mut l = vec![b'4'];
l.extend(120u64.to_le_bytes());
l.extend(b"low battery");
log.extend(message(b'L', &l));
log.extend([0xEE; 7]);
log.extend(message(b'S', &SYNC));
log.extend(sensor(1, 200, 2.0, false));
let mut s = 3u16.to_le_bytes().to_vec();
s.extend(150u64.to_le_bytes());
s.push(1);
log.extend(message(b'D', &s));
let cut = sensor(1, 300, 3.0, false);
log.extend(&cut[..cut.len() - 4]);
log
}
#[test]
fn topics_instances_and_text() {
let index = index(&tiny()).unwrap();
let names: Vec<&str> = index.names.iter().map(|(n, _)| n.as_str()).collect();
assert_eq!(names, ["sensor.0", "sensor.1", "status"]);
let sensor = &index.topics[&1];
assert_eq!(sensor.offsets.len(), 2);
let columns: Vec<&str> = sensor.columns.iter().map(|c| c.name.as_str()).collect();
assert_eq!(
columns,
["timestamp", "id", "v.x", "v.y", "v.z", "tag", "raw"]
);
assert_eq!(index.info, [("sys_name".to_string(), "PX4".to_string())]);
assert_eq!(index.parameters[0].0, "MPC_XY_VEL");
assert_eq!(index.logged[0].3, "low battery");
assert_eq!(index.damaged, 1);
assert!(index.cut_short);
}
#[test]
fn a_topic_decodes() {
let data = tiny();
let index = index(&data).unwrap();
let topic = &index.topics[&1];
let records = Arc::new(
IndexedRecords::new(
Arc::new(Bytes::Owned(data.clone())),
topic.offsets.clone(),
topic.columns.clone(),
)
.unwrap(),
);
let df = records.lazy().collect().unwrap();
assert_eq!(
df.column("v.y").unwrap().f32().unwrap().to_vec(),
[Some(2.0), Some(4.0)]
);
assert_eq!(
df.column("timestamp").unwrap().dtype(),
&DataType::Duration(TimeUnit::Microseconds)
);
assert_eq!(df.column("tag").unwrap().str().unwrap().get(0), Some("ab"));
assert_eq!(
df.column("raw").unwrap().dtype(),
&DataType::Array(Box::new(DataType::Int16), 2)
);
}
#[test]
fn garbage_is_bounded() {
assert!(index(b"nope").is_err());
let mut data = MAGIC.to_vec();
data.extend([1, 0, 0, 0, 0, 0, 0, 0, 0]);
data.extend([0xFF; 64]);
let index = index(&data).unwrap();
assert!(index.names.is_empty());
}
}