froe 0.10.0

Reader and offline maintenance toolkit for Apache Jackrabbit Oak segment-tar (TarMK) repositories: parse archives and records, extract node data, compact, back up, and recover.
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
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
//! Recovering an archive number whose newest generation cannot be
//! opened, by scanning the segments still readable in it.

use super::file_identity::preserve_file_metadata;
use super::providers::ArchiveSegmentsProvider;
use super::providers::read_blob_identifiers;
use super::repair::{AuthorizeVersionTwoWrite, VersionTwoAlreadyEstablished};
use super::startup::install_target_generation;
use crate::content::provider::SegmentProvider as _;
use crate::error::{Error, Result};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::parsed_segment::ParsedSegment;
use crate::tar_archive::archive::TarArchiveReader;
use crate::tar_archive::file_name::ArchiveFileName;
use crate::writer::segment_builder::GarbageCollectionGeneration;
use crate::writer::tar_writer::TarArchiveWriter;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;

/// Picks the generation letter of one archive number to write against:
/// newest letter first, the first valid index wins. Also reports whether
/// any letter held bytes at all.
///
/// Zero-length letters are skipped exactly as the read path skips them
/// (`crate::store::open_archives_newest_valid_first`): a writer creates its
/// next archive lazily, and an empty file is that creation's race window —
/// or what it leaves behind when it is killed inside it. Opening one yields
/// no segments, so recovering the number would rebuild it as an archive
/// with no entries, which is not a file `TarArchiveWriter` ever creates.
pub(super) fn select_writable_generation(
    directory: &Path,
    generations: &[ArchiveFileName],
) -> (Option<TarArchiveReader>, bool) {
    let mut any_nonempty = false;
    for candidate in generations.iter().rev() {
        let path = directory.join(&candidate.file_name);
        if std::fs::metadata(&path).is_ok_and(|metadata| metadata.len() == 0) {
            continue;
        }
        any_nonempty = true;
        if let Ok(reader) = TarArchiveReader::open(&path)
            && !reader.is_recovered()
        {
            return (Some(reader), any_nonempty);
        }
    }
    (None, any_nonempty)
}

/// Opens the winning generation letter of each archive number, deleting
/// the losers, and reports one completed archive number at a time.
pub(super) fn open_archive_numbers_for_writing(
    directory: &Path,
    by_number: std::collections::BTreeMap<u32, Vec<ArchiveFileName>>,
    observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<Vec<TarArchiveReader>> {
    let archive_numbers = by_number.len();
    let mut archives = Vec::new();
    for (opened, (_, mut generations)) in by_number.into_iter().enumerate() {
        observer.step_advanced(crate::progress::count(opened));
        generations.sort_by_key(|name| name.file_generation);
        let (winner, any_nonempty) = select_writable_generation(directory, &generations);
        match winner {
            Some(reader) => {
                // Delete every other generation letter of this number.
                for stale in &generations {
                    if stale.file_name != reader.file_name() {
                        std::fs::remove_file(directory.join(&stale.file_name))?;
                    }
                }
                archives.push(reader);
            }
            // Every letter of this number is empty, so there is nothing to
            // recover and nothing to serve. Nothing is deleted here: the
            // number simply contributes no archive, which either frees it
            // for the next write to fill — the same thing the writer that
            // created it was about to do — or leaves the files for
            // cleanup's stale-archive task to remove under its own
            // plan-and-confirm contract. Reuse can only ever land on a
            // zero-byte file, because a single non-empty letter sends the
            // whole number down the recovery path instead.
            None if !any_nonempty => {}
            None => {
                archives.push(recover_archive_number(
                    directory,
                    &generations,
                    &mut VersionTwoAlreadyEstablished,
                )?);
            }
        }
    }
    observer.step_advanced(crate::progress::count(archive_numbers));
    Ok(archives)
}

/// The refusal an archive number earns when it holds bytes but no segment
/// the recovery scan can read. Naming the files matters more than usual
/// here: the operator has to decide whether to move them aside or keep them
/// as evidence, and neither the number nor an errno tells them which files
/// are involved.
///
/// The remedy deliberately does not name cleanup. `plan_stale_archives`
/// marks only *zero-byte* letters of an unindexed number stale; a non-empty
/// letter is preserved with a warning, precisely because it may still hold
/// unrecovered bytes. Telling the operator to run cleanup here would send
/// them to a command that will decline to act.
pub(super) fn unrecoverable_archive_number_refusal(generations: &[ArchiveFileName]) -> Error {
    let names: Vec<&str> = generations
        .iter()
        .map(|generation| generation.file_name.as_str())
        .collect();
    Error::InvalidFormat {
        details: format!(
            "archive number {} has no valid index and no recoverable segment in {}; \
             refusing to replace it with an empty archive. Cleanup preserves this file \
             rather than removing it, so opening the store for writing needs it moved \
             aside — keep it, it is the only copy of whatever it holds",
            generations.first().map_or(0, |first| first.archive_number),
            names.join(", ")
        ),
    }
}

/// Recovers one archive number with no valid index: scans every letter in
/// ascending order (later letters overwrite duplicates), rebuilds the
/// recovered segments as a fresh archive, and only after that archive is
/// written, fsynced, and re-validated are the originals retired to
/// `.bak` names and the replacement installed under the lowest letter's
/// file name. A failure before installation leaves every original in
/// place; a failure during installation rolls back best-effort (see
/// [`install_recovered_archive`]).
/// Gives a rebuilt archive the ownership and mode of the archive it replaces,
/// rather than the process umask.
///
/// Every other replacement path in maintenance does this; without it a store
/// whose archives are group-owned and setgid silently loses both on the one
/// file that was rewritten, and a later cleanup's apply-identity preflight
/// reads the wrong metadata. A target that does not exist yet has nothing to
/// inherit, which is not an error.
pub(super) fn inherit_replaced_archive_metadata(
    directory: &Path,
    target_name: &str,
    temporary_path: &Path,
) -> Result<()> {
    let Ok(source_metadata) = std::fs::metadata(directory.join(target_name)) else {
        return Ok(());
    };
    let staged = std::fs::OpenOptions::new()
        .write(true)
        .open(temporary_path)?;
    preserve_file_metadata(&staged, &source_metadata)
}

/// Charges the caller's version-2 price at the last instant before a rebuilt
/// archive becomes visible.
///
/// The staged rebuild already exists, is durable, and has re-opened with a
/// valid index; nothing version-2 is visible yet. If authorization fails the
/// staging file is removed like every other pre-install failure, so the
/// number is left exactly as it was found.
pub(super) fn authorize_before_install(
    authorize: &mut dyn AuthorizeVersionTwoWrite,
    temporary_path: &Path,
) -> Result<()> {
    if let Err(error) = authorize.authorize() {
        let _ = std::fs::remove_file(temporary_path);
        return Err(error);
    }
    Ok(())
}

pub(super) fn recover_archive_number(
    directory: &Path,
    generations: &[ArchiveFileName],
    authorize: &mut dyn AuthorizeVersionTwoWrite,
) -> Result<TarArchiveReader> {
    let recovered = scan_recoverable_segments(directory, generations);
    // A non-empty file that yields no segment is residue this function
    // cannot act on: writing the replacement would produce an archive with
    // no entries, which `TarArchiveWriter` never creates at all, and the
    // re-open below would then fail on a missing path with a bare errno.
    // Refuse with the file names instead, and say what clears them.
    if recovered.is_empty() {
        return Err(unrecoverable_archive_number_refusal(generations));
    }

    // Parse every segment once — data *and* bulk, so blob identifier
    // strings whose block lists spill into bulk segments resolve too.
    // The parsed structures also back the provider that resolves blob
    // identifier strings across the recovered segments of this archive
    // number.
    let mut parsed_segments: HashMap<SegmentIdentifier, Arc<ParsedSegment>> = HashMap::new();
    for (identifier, bytes) in &recovered {
        parsed_segments.insert(
            *identifier,
            Arc::new(ParsedSegment::parse(*identifier, bytes)?),
        );
    }
    let provider = ArchiveSegmentsProvider {
        segments: recovered
            .iter()
            .filter_map(|(identifier, bytes)| {
                parsed_segments
                    .get(identifier)
                    .map(|parsed| (*identifier, (Arc::clone(parsed), bytes.as_slice())))
            })
            .collect(),
    };

    // Build the replacement beside the originals; nothing is renamed or
    // deleted until it exists, is durable, and re-opens with a valid
    // index.
    let target_name = &install_target_generation(directory, generations).file_name;
    let temporary_name = format!("{target_name}.recovering");
    let temporary_path = directory.join(&temporary_name);
    let _ = std::fs::remove_file(&temporary_path);
    let write_replacement =
        || -> Result<()> {
            let mut writer = TarArchiveWriter::new(directory, &temporary_name);
            for (identifier, bytes) in &recovered {
                let (generation, references, binary_references) =
                    if let Some(parsed) = parsed_segments.get(identifier) {
                        // Fail closed when a blob identifier cannot be
                        // resolved: publishing an incomplete catalog would let
                        // AEM's blob garbage collection delete a
                        // still-referenced binary.
                        let segment = provider.segment(*identifier)?;
                        let binary_references = read_blob_identifiers(&provider, &segment)
                            .map_err(|error| Error::InvalidFormat {
                                details: format!(
                                    "cannot rebuild the binary references catalog while \
                                     recovering {target_name}: an external blob identifier in \
                                     segment {identifier} does not resolve within the recovered \
                                     segments ({error}); refusing to publish an incomplete \
                                     catalog, which could let blob garbage collection delete \
                                     referenced binaries"
                                ),
                            })?;
                        (
                            GarbageCollectionGeneration {
                                generation: parsed.generation,
                                full_generation: parsed.full_generation,
                                is_compacted: parsed.is_compacted,
                            },
                            parsed.referenced_segments.clone(),
                            binary_references,
                        )
                    } else {
                        (
                            GarbageCollectionGeneration {
                                generation: 0,
                                full_generation: 0,
                                is_compacted: false,
                            },
                            Vec::new(),
                            Vec::new(),
                        )
                    };
                writer.write_segment(
                    *identifier,
                    bytes,
                    generation,
                    &references,
                    &binary_references,
                )?;
            }
            writer.close()?;
            Ok(())
        };
    if let Err(error) = write_replacement() {
        let _ = std::fs::remove_file(&temporary_path);
        return Err(error);
    }
    if let Err(error) = inherit_replaced_archive_metadata(directory, target_name, &temporary_path) {
        let _ = std::fs::remove_file(&temporary_path);
        return Err(error);
    }
    crate::writer::compaction::fsync_directory(directory);
    match TarArchiveReader::open(&temporary_path) {
        Ok(validated) if !validated.is_recovered() => drop(validated),
        Ok(_) => {
            let _ = std::fs::remove_file(&temporary_path);
            return Err(Error::InvalidFormat {
                details: format!("the rebuilt archive {temporary_name} failed index validation"),
            });
        }
        Err(error) => {
            let _ = std::fs::remove_file(&temporary_path);
            return Err(error);
        }
    }
    authorize_before_install(authorize, &temporary_path)?;
    install_recovered_archive(directory, generations, target_name, &temporary_path)
}

/// Scans every generation letter of one archive number in ascending
/// order — later letters overwrite duplicate segments — returning the
/// recovered segments in scan order.
/// Whether a rebuild of this archive number would find anything to rebuild
/// from — the same question [`scan_recoverable_segments`] answers, without
/// materializing an answer nobody wants.
///
/// The scan copies every segment's bytes into owned buffers, so asking it
/// this question allocates the whole archive to discard it, and does so in
/// the survey that now runs before every repair. `segment_count()` on a
/// recovery-scanned reader is the length of that same scan's entry list,
/// read straight off the memory map. It is also exactly what the cleanup
/// side reads off its already-open readers, so both callers now derive the
/// predicate the same way and cannot drift apart.
pub(super) fn any_recoverable_segment(directory: &Path, generations: &[ArchiveFileName]) -> bool {
    generations.iter().any(|generation| {
        TarArchiveReader::open(&directory.join(&generation.file_name))
            .is_ok_and(|reader| reader.segment_count() > 0)
    })
}

pub(super) fn scan_recoverable_segments(
    directory: &Path,
    generations: &[ArchiveFileName],
) -> Vec<(SegmentIdentifier, Vec<u8>)> {
    let mut recovered: Vec<(SegmentIdentifier, Vec<u8>)> = Vec::new();
    let mut positions: HashMap<SegmentIdentifier, usize> = HashMap::new();
    for generation in generations {
        let path = directory.join(&generation.file_name);
        if let Ok(reader) = TarArchiveReader::open(&path) {
            for identifier in reader.segment_identifiers() {
                if let Some(bytes) = reader.segment_data(identifier) {
                    if let Some(&position) = positions.get(&identifier) {
                        recovered[position].1 = bytes.to_vec();
                    } else {
                        positions.insert(identifier, recovered.len());
                        recovered.push((identifier, bytes.to_vec()));
                    }
                }
            }
        }
    }
    recovered
}

/// Retires the original generation letters to `.bak` names and installs
/// the validated replacement under the target name. The target's own
/// original is preserved through a hard link (or, on filesystems without
/// hard links, a full copy), so a `.tar` under the target name exists at
/// every instant; the other letters are plain renames. An *error* at any
/// step — including the final re-open — rolls every completed step back,
/// normally leaving the originals under their own names. The rollback is
/// best effort: a rollback rename that itself fails cannot be recovered
/// further, is dropped in favor of reporting the primary error, and can
/// leave a mix of `.bak` and installed states — as can a *crash*
/// mid-installation, the inherent limit of multi-file replacement. The
/// `.bak` copies always preserve the original bytes for manual repair.
pub(super) fn install_recovered_archive(
    directory: &Path,
    generations: &[ArchiveFileName],
    target_name: &str,
    temporary_path: &Path,
) -> Result<TarArchiveReader> {
    let mut renamed: Vec<(PathBuf, PathBuf)> = Vec::new();
    let mut target_backup: Option<PathBuf> = None;
    let roll_back = |renamed: &[(PathBuf, PathBuf)]| {
        for (original, backup) in renamed.iter().rev() {
            let _ = std::fs::rename(backup, original);
        }
    };
    for generation in generations {
        let path = directory.join(&generation.file_name);
        // Zero-length letters hold nothing to preserve and are not archives;
        // retiring one would only manufacture an empty `.bak`. They stay for
        // the stale-archive task, which plans and confirms their removal.
        if std::fs::metadata(&path).is_ok_and(|metadata| metadata.len() == 0) {
            continue;
        }
        let backup = backup_path(directory, &generation.file_name);
        if generation.file_name == *target_name {
            // The target keeps its directory entry: the backup is a
            // second link (or a copy) of the same content, never a
            // rename away.
            if std::fs::hard_link(&path, &backup).is_err()
                && let Err(error) = std::fs::copy(&path, &backup)
            {
                roll_back(&renamed);
                return Err(error.into());
            }
            target_backup = Some(backup);
        } else if let Err(error) = std::fs::rename(&path, &backup) {
            roll_back(&renamed);
            return Err(error.into());
        } else {
            renamed.push((path, backup));
        }
    }
    let target_path = directory.join(target_name);
    if let Err(error) = std::fs::rename(temporary_path, &target_path) {
        if let Some(backup) = &target_backup {
            let _ = std::fs::remove_file(backup);
        }
        roll_back(&renamed);
        return Err(error.into());
    }
    crate::writer::compaction::fsync_directory(directory);
    match TarArchiveReader::open(&target_path) {
        Ok(reader) => Ok(reader),
        Err(error) => {
            // The replacement was validated before installation, so a
            // failing re-open is environmental (for example an I/O
            // error). Restore the original atomically from its backup
            // link and undo the other renames before reporting.
            if let Some(backup) = &target_backup {
                let _ = std::fs::rename(backup, &target_path);
            }
            roll_back(&renamed);
            crate::writer::compaction::fsync_directory(directory);
            Err(error)
        }
    }
}

/// The first free `.bak` name for a damaged archive: `name.bak`, then
/// `name.2.bak`, `name.3.bak`, …
pub(super) fn backup_path(directory: &Path, file_name: &str) -> PathBuf {
    let first = directory.join(format!("{file_name}.bak"));
    if !first.exists() {
        return first;
    }
    let mut counter = 2u32;
    loop {
        let candidate = directory.join(format!("{file_name}.{counter}.bak"));
        if !candidate.exists() {
            return candidate;
        }
        counter += 1;
    }
}

#[cfg(test)]
mod tests {
    use crate::content::provider::SegmentProvider;
    use crate::store::Repository;

    use crate::writer::store_writer::repository::*;

    use crate::writer::store_writer::test_support::*;

    /// An empty archive file is what a writer killed inside its own lazy
    /// next-archive creation leaves behind — the read path has always
    /// skipped it. Before this was mirrored here, the number fell to
    /// `recover_archive_number`, which rebuilt it as an archive with no
    /// entries; `TarArchiveWriter` never creates a file for that, so the
    /// re-open failed on a missing path and every froe write command
    /// reported a bare `No such file or directory`.
    #[test]
    fn an_empty_archive_file_does_not_break_opening_for_writing() {
        let directory = TestDirectory::new("empty-archive-open");
        {
            let store = WritableRepository::open(&directory.path).expect("bootstrap");
            store.close().expect("close");
        }
        let empty = directory.path.join("data00009a.tar");
        std::fs::write(&empty, b"").expect("create the empty archive");

        let store = WritableRepository::open(&directory.path).expect("write open must succeed");
        let head = store.head();
        assert!(store.segment(head.segment).is_ok(), "head still resolves");
        store.close().expect("close");

        Repository::open(&directory.path).expect("the reader still opens");
    }
    /// A non-empty letter that yields no recoverable segment is residue
    /// `recover_archive_number` cannot act on. It must say so, not surface
    /// the `ENOENT` of the replacement it declined to write.
    #[test]
    fn an_unrecoverable_archive_is_refused_with_its_file_name() {
        let directory = TestDirectory::new("unrecoverable-archive");
        {
            let store = WritableRepository::open(&directory.path).expect("bootstrap");
            store.close().expect("close");
        }
        // 512-byte blocks that are neither a valid index nor a parseable
        // tar entry: the scan recovers nothing from them.
        let junk = directory.path.join("data00009a.tar");
        std::fs::write(&junk, vec![0x5au8; 4096]).expect("write junk archive");

        let message = match WritableRepository::open(&directory.path) {
            Ok(_) => panic!("opening for writing must refuse an unrecoverable archive"),
            Err(error) => error.to_string(),
        };
        assert!(
            message.contains("data00009a.tar"),
            "the refusal names the unusable file: {message}"
        );
        assert!(
            message.contains("no recoverable segment"),
            "the refusal states why: {message}"
        );
        assert!(
            !message.contains("No such file or directory"),
            "the refusal must not surface a bare errno: {message}"
        );
        assert!(junk.exists(), "the refusal leaves the file in place");
    }

    #[test]
    fn archives_without_an_index_are_recovered_with_backups() {
        let directory = TestDirectory::new("write-recovery");
        {
            let store = WritableRepository::open(&directory.path).expect("bootstrap");
            store.close().expect("close");
        }
        // Truncate the archive's trailers, leaving only entry data.
        let path = directory.path.join("data00000a.tar");
        let full = std::fs::read(&path).expect("read");
        // Find the first trailer: the '.brf' entry header.
        let trailer_start = full
            .windows(4)
            .position(|window| window == b".brf")
            .map(|position| (position / 512) * 512)
            .expect("brf trailer present");
        let mut truncated = full[..trailer_start].to_vec();
        truncated.extend_from_slice(&[0u8; 1024]);
        std::fs::write(&path, &truncated).expect("truncate");

        {
            let store = WritableRepository::open(&directory.path).expect("recovering open");
            let head = store.head();
            assert!(
                store.segment(head.segment).is_ok(),
                "head segment recovered"
            );
            store.close().expect("close");
        }
        assert!(
            directory.path.join("data00000a.tar.bak").exists(),
            "the damaged archive is backed up"
        );
        let repository = Repository::open(&directory.path).expect("reader opens");
        assert!(
            !repository
                .archives()
                .iter()
                .any(crate::tar_archive::archive::TarArchiveReader::is_recovered),
            "the regenerated archive has a valid index"
        );
        repository.content_root().expect("content root resolves");
    }
}