falcon_mdf/multi_ops.rs
1//! Operations that span several channels or several measurements:
2//! [`Mf4File::filter`](crate::Mf4File::filter), [`concatenate`] and [`stack`].
3//!
4//! All three hand back decoded [`SignalSeries`] in memory. None of them writes
5//! a file: a concatenated or stacked measurement has a structure our writer
6//! cannot express, and producing a half-correct file would be worse than
7//! producing none.
8//!
9//! # How the time axes are lined up
10//!
11//! An MF4 file records absolute time in its header and relative time in its
12//! master channel, so two files recorded an hour apart both start their master
13//! at (or near) zero. Combining them without looking at the headers would pile
14//! both recordings on top of each other at t = 0. [`TimeAlignment::StartTime`],
15//! the default, therefore shifts every file by its header start time relative
16//! to the earliest file's — the same rule asammdf's `sync=True` applies in
17//! `MDF.concatenate` and `MDF.stack`. [`TimeAlignment::AsRecorded`] is the
18//! escape hatch for files whose headers are not trustworthy.
19
20use std::sync::Arc;
21
22use crate::error::{Mf4Error, Result};
23use crate::file::Mf4File;
24use crate::model::{Channel, SignalValues};
25use crate::time_ops::SignalSeries;
26
27/// How a channel is picked out of a file by [`Mf4File::filter`].
28///
29/// `ChannelSelector::from("Speed")` is the common case; the other two variants
30/// exist for names that several channels share.
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub enum ChannelSelector {
33 /// The channel with this name, which must be unique in the file.
34 ///
35 /// A name carried by more than one channel is an error rather than a
36 /// silent pick of the first: a master channel called `t` exists in every
37 /// group, and quietly filtering the wrong group's copy of it is precisely
38 /// the kind of plausible-but-wrong answer this crate tries not to give.
39 /// Disambiguate with [`ChannelSelector::NameInGroup`] or
40 /// [`ChannelSelector::Position`].
41 Name(String),
42 /// The channel with this name inside one specific channel group.
43 NameInGroup {
44 /// Channel name.
45 name: String,
46 /// Index of the data group.
47 data_group: usize,
48 /// Index of the channel group within that data group.
49 channel_group: usize,
50 },
51 /// The channel at an exact position, regardless of its name.
52 Position {
53 /// Index of the data group.
54 data_group: usize,
55 /// Index of the channel group within that data group.
56 channel_group: usize,
57 /// Index of the channel within that channel group.
58 index: usize,
59 },
60}
61
62impl From<&str> for ChannelSelector {
63 fn from(name: &str) -> Self {
64 ChannelSelector::Name(name.to_string())
65 }
66}
67
68impl From<String> for ChannelSelector {
69 fn from(name: String) -> Self {
70 ChannelSelector::Name(name)
71 }
72}
73
74impl ChannelSelector {
75 /// Resolves this selector against a file.
76 pub(crate) fn resolve<'a>(&self, file: &'a Mf4File) -> Result<&'a Channel> {
77 match self {
78 ChannelSelector::Name(name) => {
79 let found = file.find_channels(name);
80 match found.len() {
81 0 => Err(Mf4Error::ChannelNotFound { name: name.clone() }),
82 1 => Ok(found[0]),
83 n => Err(Mf4Error::parse_error(format!(
84 "channel name '{name}' is carried by {n} channels; select it with \
85 ChannelSelector::NameInGroup or ChannelSelector::Position"
86 ))),
87 }
88 }
89 ChannelSelector::NameInGroup {
90 name,
91 data_group,
92 channel_group,
93 } => {
94 let group = group_at(file, *data_group, *channel_group)?;
95 group.find_channel(name).ok_or_else(|| {
96 Mf4Error::parse_error(format!(
97 "channel '{name}' not found in data group {data_group}, channel group \
98 {channel_group}"
99 ))
100 })
101 }
102 ChannelSelector::Position {
103 data_group,
104 channel_group,
105 index,
106 } => {
107 let group = group_at(file, *data_group, *channel_group)?;
108 group.channels.get(*index).ok_or_else(|| {
109 Mf4Error::parse_error(format!(
110 "channel index {index} is out of range for data group {data_group}, \
111 channel group {channel_group}, which has {} channels",
112 group.channels.len()
113 ))
114 })
115 }
116 }
117 }
118}
119
120fn group_at(
121 file: &Mf4File,
122 data_group: usize,
123 channel_group: usize,
124) -> Result<&crate::model::ChannelGroup> {
125 let dg = file.data_groups().get(data_group).ok_or_else(|| {
126 Mf4Error::parse_error(format!(
127 "data group {data_group} is out of range; the file has {}",
128 file.data_groups().len()
129 ))
130 })?;
131 dg.channel_groups.get(channel_group).ok_or_else(|| {
132 Mf4Error::parse_error(format!(
133 "channel group {channel_group} is out of range for data group {data_group}, which \
134 has {}",
135 dg.channel_groups.len()
136 ))
137 })
138}
139
140/// How the time axes of several measurements are lined up before combining
141/// them.
142#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
143pub enum TimeAlignment {
144 /// Shift every file by its header start time relative to the earliest
145 /// file's, so a later recording lands after an earlier one.
146 ///
147 /// Equivalent to asammdf's `sync=True`, which is its default.
148 #[default]
149 StartTime,
150 /// Use each file's master channel exactly as recorded, applying no shift.
151 ///
152 /// Equivalent to asammdf's `sync=False`.
153 AsRecorded,
154}
155
156/// One series produced by [`stack`], tagged with the file it came from.
157///
158/// Stacking keeps same-named channels from different files apart, so the name
159/// alone does not identify a series; `file_index` does.
160#[derive(Debug, Clone, PartialEq)]
161pub struct StackedSeries {
162 /// Position of the source file in the slice handed to [`stack`].
163 pub file_index: usize,
164 /// The decoded series, with its timestamps already shifted.
165 pub series: SignalSeries,
166}
167
168/// Joins several measurements end to end along time.
169///
170/// The files must have the same internal structure: the same number of channel
171/// groups, and the same set of channel names in each. A channel present in one
172/// file and absent from another is an error, not a gap to be filled — the two
173/// files describe different measurements, and inventing samples for the
174/// missing one would fabricate data. This is what asammdf's `MDF.concatenate`
175/// does too.
176///
177/// Channels sharing a name across files are joined into one series: that is
178/// the point of concatenating. Their order within a group may differ between
179/// files; they are matched by name, and the result follows the first file's
180/// order.
181///
182/// # Time
183///
184/// Each file is shifted by its header start time relative to the earliest
185/// file's (see [`TimeAlignment`]). If, after that shift, a file would still
186/// start at or before the previous file's last timestamp, it is pushed forward
187/// to begin one sample interval after it — the result is monotonic, and no two
188/// files overlap. The sample interval used is the second minus the first
189/// timestamp of the file being appended, or 1 ms when it has fewer than two
190/// samples. asammdf applies the identical rule.
191///
192/// # Returns
193///
194/// One [`SignalSeries`] per non-master channel, in the first file's group and
195/// channel order. The channel metadata is the first file's.
196///
197/// # Example
198///
199/// ```no_run
200/// # use falcon_mdf::{Mf4File, TimeAlignment, multi_ops};
201/// let a = Mf4File::open("run1.mf4")?;
202/// let b = Mf4File::open("run2.mf4")?;
203/// let joined = multi_ops::concatenate(&[&a, &b], TimeAlignment::StartTime)?;
204/// println!("{} samples", joined[0].len());
205/// # Ok::<(), falcon_mdf::error::Mf4Error>(())
206/// ```
207pub fn concatenate(files: &[&Mf4File], alignment: TimeAlignment) -> Result<Vec<SignalSeries>> {
208 if files.is_empty() {
209 return Err(Mf4Error::parse_error("concatenate needs at least one file"));
210 }
211
212 let offsets = start_time_offsets(files, alignment);
213 let layouts: Vec<Vec<GroupLayout<'_>>> = files.iter().map(|f| group_layouts(f)).collect();
214 check_same_structure(&layouts)?;
215
216 let mut out: Vec<SignalSeries> = Vec::new();
217
218 for (group_index, first_group) in layouts[0].iter().enumerate() {
219 if first_group.channels.is_empty() {
220 continue;
221 }
222
223 // One accumulator per channel of this group, in the first file's order.
224 let mut timestamps: Vec<Vec<f64>> = vec![Vec::new(); first_group.channels.len()];
225 let mut values: Vec<Option<SignalValues>> = vec![None; first_group.channels.len()];
226 let mut validity: Vec<Option<Vec<bool>>> = vec![None; first_group.channels.len()];
227 let mut counts: Vec<usize> = vec![0; first_group.channels.len()];
228 let mut last_timestamp: Option<f64> = None;
229
230 for (file_index, file) in files.iter().enumerate() {
231 let group = &layouts[file_index][group_index];
232 // Matched by name, so a file that lists the same channels in a
233 // different order still lines up. Matches are consumed: a group
234 // holding two channels of the same name would otherwise resolve
235 // both of the first file's slots to the same channel of the
236 // second, silently duplicating one and dropping the other.
237 let mut taken = vec![false; group.channels.len()];
238 let selectors: Vec<&Channel> = first_group
239 .channels
240 .iter()
241 .map(|ch| {
242 let at = group
243 .channels
244 .iter()
245 .enumerate()
246 .find(|(i, other)| other.name == ch.name && !taken[*i])
247 .map(|(i, _)| i)
248 .expect("the structure check matched the names as multisets");
249 taken[at] = true;
250 group.channels[at]
251 })
252 .collect();
253
254 let series = file.series_for(&selectors)?;
255
256 // Every channel of a group shares its master, so the shifted time
257 // axis is computed once from the first of them and applied to all.
258 // A group with no samples contributes an empty axis and leaves
259 // `last_timestamp` alone — it is not a recording that ended, so
260 // nothing should be continued from it — but its channels are still
261 // merged below, so a channel no file has samples for still reports
262 // its own sample type rather than a stand-in.
263 let master: Vec<f64> = match series.first() {
264 Some(first) if !first.timestamps.is_empty() => {
265 let mut master: Vec<f64> = first
266 .timestamps
267 .iter()
268 .map(|t| t + offsets[file_index])
269 .collect();
270 if let Some(last) = last_timestamp {
271 if last >= master[0] {
272 let delta = if master.len() >= 2 {
273 master[1] - master[0]
274 } else {
275 0.001
276 };
277 let shift = last + delta - master[0];
278 for t in &mut master {
279 *t += shift;
280 }
281 }
282 }
283 last_timestamp = master.last().copied();
284 master
285 }
286 _ => Vec::new(),
287 };
288
289 for (slot, s) in series.into_iter().enumerate() {
290 timestamps[slot].extend_from_slice(&master);
291 let added = s.values.len();
292 append_validity(
293 &mut validity[slot],
294 counts[slot],
295 s.validity.as_deref(),
296 added,
297 );
298 counts[slot] += added;
299 match &mut values[slot] {
300 Some(acc) => append_values(acc, &s.values)?,
301 none => *none = Some(s.values),
302 }
303 }
304 }
305
306 for (slot, channel) in first_group.channels.iter().enumerate() {
307 // Always `Some` in practice: the file list is non-empty and every
308 // file yields one series per slot, even a file with no samples.
309 let vals = values[slot].take().unwrap_or(SignalValues::F64(Vec::new()));
310 out.push(SignalSeries::new(
311 (*channel).clone(),
312 std::mem::take(&mut timestamps[slot]),
313 vals,
314 validity[slot].take(),
315 )?);
316 }
317 }
318
319 Ok(out)
320}
321
322/// Combines several measurements side by side: every file's channels present
323/// together, rather than one file's samples after another's.
324///
325/// Nothing is merged and nothing is resampled. A channel name that appears in
326/// two files yields two series, each keeping its own file's samples and its own
327/// file's sample rate; [`StackedSeries::file_index`] says which file each came
328/// from. Files need not have the same structure — a channel present in only one
329/// of them simply appears once.
330///
331/// The "common time base" is a common *origin*, not a common raster: under
332/// [`TimeAlignment::StartTime`] each file's timestamps are shifted by its
333/// header start time relative to the earliest file's, so t = 0 means the same
334/// instant for all of them. asammdf's `MDF.stack` does exactly this. To put
335/// stacked series on one raster as well, resample them afterwards with
336/// [`SignalSeries::resample`].
337///
338/// # Example
339///
340/// ```no_run
341/// # use falcon_mdf::{Mf4File, TimeAlignment, multi_ops};
342/// let a = Mf4File::open("engine.mf4")?;
343/// let b = Mf4File::open("chassis.mf4")?;
344/// for s in multi_ops::stack(&[&a, &b], TimeAlignment::StartTime)? {
345/// println!("file {}: {}", s.file_index, s.series.name());
346/// }
347/// # Ok::<(), falcon_mdf::error::Mf4Error>(())
348/// ```
349pub fn stack(files: &[&Mf4File], alignment: TimeAlignment) -> Result<Vec<StackedSeries>> {
350 if files.is_empty() {
351 return Err(Mf4Error::parse_error("stack needs at least one file"));
352 }
353
354 let offsets = start_time_offsets(files, alignment);
355 let mut out = Vec::new();
356
357 for (file_index, file) in files.iter().enumerate() {
358 let offset = offsets[file_index];
359 for group in group_layouts(file) {
360 if group.channels.is_empty() {
361 continue;
362 }
363 // Every channel of a group sits on the same master axis, so the
364 // shifted axis is built once and shared: shifting each series'
365 // own timestamps would copy the axis per channel.
366 let series_list = file.series_for(&group.channels)?;
367 let shifted = series_list.first().map(|first| {
368 Arc::new(
369 first
370 .timestamps()
371 .iter()
372 .map(|t| t + offset)
373 .collect::<Vec<f64>>(),
374 )
375 });
376 for mut series in series_list {
377 if let Some(shifted) = &shifted {
378 if shifted.len() == series.timestamps.len() {
379 series.timestamps = Arc::clone(shifted);
380 out.push(StackedSeries { file_index, series });
381 continue;
382 }
383 }
384 for t in Arc::make_mut(&mut series.timestamps) {
385 *t += offset;
386 }
387 out.push(StackedSeries { file_index, series });
388 }
389 }
390 }
391
392 Ok(out)
393}
394
395/// The non-master channels of one channel group.
396struct GroupLayout<'a> {
397 channels: Vec<&'a Channel>,
398}
399
400/// Lists every channel group's non-master channels, in file order.
401fn group_layouts(file: &Mf4File) -> Vec<GroupLayout<'_>> {
402 file.data_groups()
403 .iter()
404 .flat_map(|dg| dg.channel_groups.iter())
405 .map(|cg| GroupLayout {
406 // The master is the time axis, carried on every series already;
407 // emitting it as a channel of its own would duplicate it.
408 channels: cg.channels.iter().filter(|ch| !ch.is_master()).collect(),
409 })
410 .collect()
411}
412
413/// Offsets, in seconds, from the earliest file's header start time.
414///
415/// Clamped at zero so that a file whose header claims to predate the earliest
416/// one cannot pull samples backwards; asammdf clamps the same way.
417fn start_time_offsets(files: &[&Mf4File], alignment: TimeAlignment) -> Vec<f64> {
418 match alignment {
419 TimeAlignment::AsRecorded => vec![0.0; files.len()],
420 TimeAlignment::StartTime => {
421 let starts: Vec<i64> = files.iter().map(|f| f.start_time().timestamp_ns).collect();
422 let oldest = starts.iter().copied().min().unwrap_or(0);
423 // `saturating_sub`: a header start time is whatever the file says
424 // it is, and two of them far enough apart overflow an i64 of
425 // nanoseconds. A wrapped difference would place a file at a
426 // plausible but wrong offset, which is the failure mode to avoid.
427 starts
428 .iter()
429 .map(|&ns| (ns.saturating_sub(oldest) as f64 / 1e9).max(0.0))
430 .collect()
431 }
432 }
433}
434
435/// Rejects files that do not describe the same measurement.
436fn check_same_structure(layouts: &[Vec<GroupLayout<'_>>]) -> Result<()> {
437 let first = &layouts[0];
438 for (file_index, layout) in layouts.iter().enumerate().skip(1) {
439 if layout.len() != first.len() {
440 return Err(Mf4Error::parse_error(format!(
441 "cannot concatenate: file 0 has {} channel groups but file {file_index} has {}",
442 first.len(),
443 layout.len()
444 )));
445 }
446 for (group_index, (a, b)) in first.iter().zip(layout).enumerate() {
447 let mut want: Vec<&str> = a.channels.iter().map(|ch| ch.name.as_str()).collect();
448 let mut got: Vec<&str> = b.channels.iter().map(|ch| ch.name.as_str()).collect();
449 want.sort_unstable();
450 got.sort_unstable();
451 if want != got {
452 return Err(Mf4Error::parse_error(format!(
453 "cannot concatenate: channel group {group_index} holds {want:?} in file 0 \
454 but {got:?} in file {file_index}"
455 )));
456 }
457 }
458 }
459 Ok(())
460}
461
462/// Extends an accumulated validity mask with `added` samples' worth of `next`.
463///
464/// A file without invalidation bits contributes valid samples, so mixing a file
465/// that has a mask with one that does not yields a mask covering both rather
466/// than dropping the one that existed.
467fn append_validity(
468 acc: &mut Option<Vec<bool>>,
469 acc_len: usize,
470 next: Option<&[bool]>,
471 added: usize,
472) {
473 match (acc.as_mut(), next) {
474 (None, None) => {}
475 (None, Some(v)) => {
476 let mut mask = vec![true; acc_len];
477 mask.extend_from_slice(v);
478 *acc = Some(mask);
479 }
480 (Some(mask), Some(v)) => mask.extend_from_slice(v),
481 (Some(mask), None) => mask.extend(std::iter::repeat_n(true, added)),
482 }
483}
484
485/// Appends `next` to `acc` in place, refusing to mix sample representations.
486///
487/// The refusal matters: two files whose copies of a channel decode to different
488/// widths are not the same channel, and gluing their bytes together would
489/// produce samples that parse but mean nothing.
490fn append_values(acc: &mut SignalValues, next: &SignalValues) -> Result<()> {
491 fn mismatch(acc: &SignalValues, next: &SignalValues) -> Mf4Error {
492 Mf4Error::parse_error(format!(
493 "cannot concatenate {:?} samples onto {:?} samples",
494 next.kind(),
495 acc.kind()
496 ))
497 }
498
499 match (acc, next) {
500 (SignalValues::U8(a), SignalValues::U8(b)) => a.extend_from_slice(b),
501 (SignalValues::U16(a), SignalValues::U16(b)) => a.extend_from_slice(b),
502 (SignalValues::U32(a), SignalValues::U32(b)) => a.extend_from_slice(b),
503 (SignalValues::U64(a), SignalValues::U64(b)) => a.extend_from_slice(b),
504 (SignalValues::I8(a), SignalValues::I8(b)) => a.extend_from_slice(b),
505 (SignalValues::I16(a), SignalValues::I16(b)) => a.extend_from_slice(b),
506 (SignalValues::I32(a), SignalValues::I32(b)) => a.extend_from_slice(b),
507 (SignalValues::I64(a), SignalValues::I64(b)) => a.extend_from_slice(b),
508 (SignalValues::F32(a), SignalValues::F32(b)) => a.extend_from_slice(b),
509 (SignalValues::F64(a), SignalValues::F64(b)) => a.extend_from_slice(b),
510 (SignalValues::Str(a), SignalValues::Str(b)) => a.extend_from_slice(b),
511 (SignalValues::CanopenDate(a), SignalValues::CanopenDate(b)) => a.extend_from_slice(b),
512 (SignalValues::CanopenTime(a), SignalValues::CanopenTime(b)) => a.extend_from_slice(b),
513 (SignalValues::Complex { re, im }, SignalValues::Complex { re: re_b, im: im_b }) => {
514 re.extend_from_slice(re_b);
515 im.extend_from_slice(im_b);
516 }
517 (
518 SignalValues::Bytes { data, width },
519 SignalValues::Bytes {
520 data: data_b,
521 width: width_b,
522 },
523 ) => {
524 if width != width_b {
525 return Err(Mf4Error::parse_error(format!(
526 "cannot concatenate {width_b}-byte samples onto {width}-byte samples"
527 )));
528 }
529 data.extend_from_slice(data_b);
530 }
531 (
532 SignalValues::VarBytes { data, starts },
533 SignalValues::VarBytes {
534 data: data_b,
535 starts: starts_b,
536 },
537 ) => {
538 let base = data.len();
539 data.extend_from_slice(data_b);
540 starts.extend(starts_b.iter().skip(1).map(|&s| s + base));
541 }
542 (
543 SignalValues::Array {
544 values,
545 elements_per_sample,
546 },
547 SignalValues::Array {
548 values: values_b,
549 elements_per_sample: per_b,
550 },
551 ) => {
552 if elements_per_sample != per_b {
553 return Err(Mf4Error::parse_error(format!(
554 "cannot concatenate {per_b}-element samples onto \
555 {elements_per_sample}-element samples"
556 )));
557 }
558 values.extend_from_slice(values_b);
559 }
560 (
561 SignalValues::ArrayVarLen { values, starts },
562 SignalValues::ArrayVarLen {
563 values: values_b,
564 starts: starts_b,
565 },
566 ) => {
567 let base = values.len();
568 values.extend_from_slice(values_b);
569 starts.extend(starts_b.iter().skip(1).map(|&s| s + base));
570 }
571 (acc, next) => return Err(mismatch(acc, next)),
572 }
573 Ok(())
574}