re_importer 0.36.1

Handles importing of Rerun data from file using importer plugins
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
use std::borrow::Cow;

use ahash::{HashMap, HashMapExt as _};
use re_log_channel::LogSender;
use re_log_types::{ApplicationId, FileSource, LogMsg};

use crate::{ImportedData, Importer as _, ImporterError, RrdImporter};

// ---

/// Imports from the given `path` using all [`crate::Importer`]s available.
///
/// A single `path` might be handled by more than one importer.
///
/// Synchronously checks whether the file exists and can be loaded. Beyond that, all
/// errors are asynchronous and handled directly by the [`crate::Importer`]s themselves
/// (i.e. they're logged).
#[cfg(not(target_arch = "wasm32"))]
pub fn import_from_path(
    settings: &crate::ImporterSettings,
    file_source: FileSource,
    path: &std::path::Path,
    // NOTE: This channel must be unbounded since we serialize all operations when running on wasm.
    tx: &LogSender,
) -> Result<(), ImporterError> {
    re_tracing::profile_function!(path.to_string_lossy());

    if !path.exists() {
        return Err(std::io::Error::new(
            std::io::ErrorKind::NotFound,
            format!("path does not exist: {path:?}"),
        )
        .into());
    }

    re_log::info!("Loading {path:?}…");

    // If no application ID was specified, we derive one from the filename.
    let application_id = settings
        .application_id
        .clone()
        .or_else(|| application_id_from_path(path));
    let settings = crate::ImporterSettings {
        // When importing a LeRobot dataset, avoid sending a `SetStoreInfo` message since the LeRobot importer handles this automatically.
        force_store_info: !re_lerobot::is_lerobot_dataset(path),
        application_id,
        ..settings.clone()
    };

    let rx = import(&settings, path, None)?;

    send(settings, file_source, rx, tx);

    Ok(())
}

/// Imports from the given `contents` using all [`crate::Importer`]s available.
///
/// A single file might be handled by more than one importer.
///
/// Synchronously checks that the file can be loaded. Beyond that, all errors are asynchronous
/// and handled directly by the [`crate::Importer`]s themselves (i.e. they're logged).
///
/// `path` is only used for informational purposes, no data is ever read from the filesystem.
pub fn import_from_file_contents(
    settings: &crate::ImporterSettings,
    file_source: FileSource,
    filepath: &std::path::Path,
    contents: std::borrow::Cow<'_, [u8]>,
    // NOTE: This channel must be unbounded since we serialize all operations when running on wasm.
    tx: &LogSender,
) -> Result<(), ImporterError> {
    re_tracing::profile_function!(filepath.to_string_lossy());

    re_log::info!("Loading {filepath:?}…");

    let application_id = settings
        .application_id
        .clone()
        .or_else(|| application_id_from_path(filepath));

    let settings = crate::ImporterSettings {
        application_id,
        ..settings.clone()
    };

    let data = import(&settings, filepath, Some(contents))?;

    send(settings, file_source, data, tx);

    Ok(())
}

// ---

fn application_id_from_path(path: &std::path::Path) -> Option<ApplicationId> {
    // `.` is not a valid application ID, so ignore line ending.
    let file_stem = path.file_stem()?;

    // Any remaining . are replaced with `_`
    let name = file_stem.to_string_lossy().replace('.', "_");

    ApplicationId::try_new(name).ok()
}

/// Prepares an adequate [`re_log_types::StoreInfo`] [`LogMsg`] given the input.
pub fn prepare_store_info(store_id: &re_log_types::StoreId, file_source: FileSource) -> LogMsg {
    re_tracing::profile_function!();

    use re_log_types::SetStoreInfo;

    let store_source = re_log_types::StoreSource::File { file_source };

    LogMsg::SetStoreInfo(SetStoreInfo {
        row_id: *re_chunk::RowId::new(),
        info: re_log_types::StoreInfo::new(store_id.clone(), store_source),
    })
}

/// Imports data at `path` using all available [`crate::Importer`]s.
///
/// On success, returns a channel with all the [`ImportedData`]:
/// - On native, this is filled asynchronously from other threads.
/// - On wasm, this is pre-filled synchronously.
///
/// There is only one way this function can return an error: not a single [`crate::Importer`]
/// (whether it is builtin, custom or external) was capable of loading the data, in which case
/// [`ImporterError::Incompatible`] will be returned.
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn import(
    settings: &crate::ImporterSettings,
    path: &std::path::Path,
    contents: Option<std::borrow::Cow<'_, [u8]>>,
) -> Result<crossbeam::channel::Receiver<ImportedData>, ImporterError> {
    re_tracing::profile_function!(path.display().to_string());

    // On native we run importers in parallel so this needs to become static.
    let contents: Option<std::sync::Arc<std::borrow::Cow<'static, [u8]>>> =
        contents.map(|contents| std::sync::Arc::new(Cow::Owned(contents.into_owned())));

    let rx_importer = {
        let (tx_importer, rx_importer) = crossbeam::channel::bounded(1024);

        let any_compatible_importer = {
            #[derive(Debug, PartialEq, Eq)]
            struct CompatibleImporterFound;
            let (tx_feedback, rx_feedback) =
                crossbeam::channel::bounded::<CompatibleImporterFound>(128);

            // When loading a file type with native support (.rrd, .mcap, .png, …)
            // then we don't need the overhead and noise of external importers:
            // See <https://github.com/rerun-io/rerun/issues/6530>.
            let importers = {
                use rayon::iter::Either;

                use crate::Importer as _;

                let extension = crate::extension(path);
                if crate::is_supported_file_extension(&extension) {
                    Either::Left(
                        crate::iter_importers()
                            .filter(|importer| importer.name() != crate::ExternalImporter.name()),
                    )
                } else {
                    // We need to use an external importer
                    Either::Right(crate::iter_importers())
                }
            };

            for importer in importers {
                let importer = std::sync::Arc::clone(&importer);

                let settings = settings.clone();
                let path = path.to_owned();
                let contents = contents.clone(); // arc

                let tx_importer = tx_importer.clone();
                let tx_feedback = tx_feedback.clone();

                rayon::spawn(move || {
                    re_tracing::profile_scope!("inner", importer.name());

                    if let Some(contents) = contents.as_deref() {
                        let contents = Cow::Borrowed(contents.as_ref());

                        if let Err(err) = importer.import_from_file_contents(
                            &settings,
                            path.clone(),
                            contents,
                            tx_importer,
                        ) {
                            if err.is_incompatible() {
                                return;
                            }
                            re_log::error!(?path, importer = importer.name(), %err, "Failed to import data");
                        }
                    } else if let Err(err) =
                        importer.import_from_path(&settings, path.clone(), tx_importer)
                    {
                        if err.is_incompatible() {
                            return;
                        }
                        re_log::error!(?path, importer = importer.name(), %err, "Failed to import data from file");
                    }

                    re_log::debug!(
                        importer = importer.name(),
                        ?path,
                        "compatible importer found"
                    );
                    re_quota_channel::send_crossbeam(&tx_feedback, CompatibleImporterFound).ok();
                });
            }

            re_tracing::profile_wait!("compatible_importer");

            drop(tx_feedback);

            rx_feedback.recv() == Ok(CompatibleImporterFound)
        };

        // Implicitly closing `tx_importer`!

        any_compatible_importer.then_some(rx_importer)
    };

    if let Some(rx_importer) = rx_importer {
        Ok(rx_importer)
    } else {
        Err(ImporterError::Incompatible(path.to_owned()))
    }
}

/// Imports data at `path` using all available [`crate::Importer`]s.
///
/// On success, returns a channel (pre-filled synchronously) with all the [`ImportedData`].
///
/// There is only one way this function can return an error: not a single [`crate::Importer`]
/// (whether it is builtin, custom or external) was capable of loading the data, in which case
/// [`ImporterError::Incompatible`] will be returned.
#[cfg(target_arch = "wasm32")]
#[expect(clippy::needless_pass_by_value)]
pub(crate) fn import(
    settings: &crate::ImporterSettings,
    path: &std::path::Path,
    contents: Option<std::borrow::Cow<'_, [u8]>>,
) -> Result<crossbeam::channel::Receiver<ImportedData>, ImporterError> {
    re_tracing::profile_function!(path.display().to_string());

    let rx_importer = {
        let (tx_importer, rx_importer) = crossbeam::channel::unbounded();

        let any_compatible_importer = crate::iter_importers().any(|importer| {
            if let Some(contents) = contents.as_deref() {
                let settings = settings.clone();
                let tx_importer = tx_importer.clone();
                let path = path.to_owned();
                let contents = Cow::Borrowed(contents);

                if let Err(err) = importer.import_from_file_contents(&settings, path.clone(), contents, tx_importer) {
                    if err.is_incompatible() {
                        return false;
                    }
                    re_log::error!(?path, importer = importer.name(), %err, "Failed to import data from file");
                }

                true
            } else {
                false
            }
        });

        // Implicitly closing `tx_importer`!

        any_compatible_importer.then_some(rx_importer)
    };

    if let Some(rx_importer) = rx_importer {
        Ok(rx_importer)
    } else {
        Err(ImporterError::Incompatible(path.to_owned()))
    }
}

/// Forwards the data in `rx_importer` to `tx`, taking care of necessary conversions, if any.
///
/// Runs asynchronously from another thread on native, synchronously on wasm.
pub(crate) fn send(
    settings: crate::ImporterSettings,
    file_source: FileSource,
    rx_importer: crossbeam::channel::Receiver<ImportedData>,
    tx: &LogSender,
) {
    spawn({
        re_tracing::profile_function!();

        #[derive(Default, Debug)]
        struct Tracked {
            is_rrd_or_rbl: bool,
            already_has_store_info: bool,
        }

        let mut store_info_tracker: HashMap<re_log_types::StoreId, Tracked> = HashMap::new();

        let tx = tx.clone();
        move || {
            // ## Ignoring channel errors
            //
            // Not our problem whether or not the other end has hung up, but we still want to
            // poll the channel in any case so as to make sure that the data producer
            // doesn't get stuck.
            for data in rx_importer {
                let importer_name = data.importer_name().clone();
                let msg = match data.into_log_msg() {
                    Ok(msg) => {
                        let store_info = match &msg {
                            LogMsg::SetStoreInfo(set_store_info) => {
                                Some((set_store_info.info.store_id.clone(), true))
                            }
                            LogMsg::ArrowMsg(store_id, _arrow_msg) => {
                                Some((store_id.clone(), false))
                            }
                            LogMsg::BlueprintActivationCommand(_) => None,
                        };

                        if let Some((store_id, store_info_created)) = store_info {
                            let tracked = store_info_tracker.entry(store_id).or_default();
                            tracked.is_rrd_or_rbl =
                                *importer_name == RrdImporter::name(&RrdImporter);
                            tracked.already_has_store_info |= store_info_created;
                        }

                        msg
                    }
                    Err(err) => {
                        re_log::error!(%err, "Couldn't serialize component data");
                        continue;
                    }
                };
                tx.send(msg.into()).ok();
            }

            for (store_id, tracked) in store_info_tracker {
                let is_a_preexisting_recording =
                    Some(&store_id) == settings.opened_store_id.as_ref();

                // Never try to send custom store info for RRDs and RBLs, they always have their own, and
                // it's always right.
                let should_force_store_info = settings.force_store_info && !tracked.is_rrd_or_rbl;

                let should_send_new_store_info = should_force_store_info
                    || (!tracked.already_has_store_info && !is_a_preexisting_recording);

                if should_send_new_store_info {
                    let store_info = prepare_store_info(&store_id, file_source.clone());
                    tx.send(store_info.into()).ok();
                }
            }

            tx.quit(None).ok();
        }
    });
}

// NOTE:
// - On native, we parallelize using `rayon`.
// - On wasm, we serialize everything, which works because the data-loading channels are unbounded.

#[cfg(not(target_arch = "wasm32"))]
fn spawn<F>(f: F)
where
    F: FnOnce() + Send + 'static,
{
    if 1 < rayon::current_num_threads() {
        rayon::spawn(f);
    } else {
        // Avoids a deadlock when send-channel gets full.
        // We usually only use `-j1` for profiling the main application; not data loading.
        std::thread::Builder::new()
            .name("importer".to_owned())
            .spawn(f)
            .expect("Failed to spawn a thread");
    }
}

#[cfg(target_arch = "wasm32")]
fn spawn<F>(f: F)
where
    F: FnOnce(),
{
    f();
}

#[cfg(test)]
mod tests {
    #[test]
    fn test_application_id_from_path() {
        use super::*;

        assert_eq!(
            application_id_from_path(std::path::Path::new("foo/bar/baz.rrd")),
            Some(ApplicationId::try_new("baz").unwrap())
        );
        assert_eq!(
            application_id_from_path(std::path::Path::new("foo/bar/baz.thing.jpg")),
            Some(ApplicationId::try_new("baz_thing").unwrap())
        );
    }
}