Skip to main content

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}