froe 0.12.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
//! What a reclamation run intends to do to each archive, and the
//! totals a caller previews before authorizing it.

use super::archive_certificate::certify_active_archives;
use super::providers::CertifiedReclaimSources;
use super::reclaim::{
    ArchiveRewritePolicy, ReclaimRule, analyze_standalone_segment_cleanup,
    reject_duplicate_active_segments,
};
use crate::error::{Error, Result};
use crate::segment::identifier::SegmentIdentifier;
use std::collections::HashMap;
use std::path::Path;

/// One archive's physical disposition in a standalone cleanup plan.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) enum PlannedArchiveSweep {
    /// Every entry is reclaimable; the archive can be unlinked whole.
    Remove {
        file_name: String,
        segment_count: usize,
        file_bytes: u64,
    },
    /// Enough entries are reclaimable to rewrite the survivors.
    Rewrite {
        file_name: String,
        replacement_name: String,
        segment_count: usize,
        eligible_entry_bytes: u64,
    },
    /// Reclaimable entries exist, but Oak's 25% savings gate keeps the
    /// archive byte-for-byte unchanged.
    DeferredBySavings {
        file_name: String,
        segment_count: usize,
        eligible_entry_bytes: u64,
    },
    /// Reclaimable entries exist, but the archive has exhausted the `a` to
    /// `z` rewrite namespace.
    DeferredAtLastGeneration {
        file_name: String,
        segment_count: usize,
        eligible_entry_bytes: u64,
    },
    /// Another generation pathname blocks a rewrite target or would be
    /// promoted by whole-file removal. Cleanup never truncates or promotes it
    /// implicitly; archive hygiene must classify it first.
    BlockedByOccupiedGeneration {
        file_name: String,
        occupied_name: String,
        segment_count: usize,
        eligible_entry_bytes: u64,
    },
}

impl PlannedArchiveSweep {
    pub(crate) fn file_name(&self) -> &str {
        match self {
            Self::Remove { file_name, .. }
            | Self::Rewrite { file_name, .. }
            | Self::DeferredBySavings { file_name, .. }
            | Self::DeferredAtLastGeneration { file_name, .. }
            | Self::BlockedByOccupiedGeneration { file_name, .. } => file_name,
        }
    }

    pub(crate) fn changes_disk(&self) -> bool {
        matches!(self, Self::Remove { .. } | Self::Rewrite { .. })
    }
}

/// Read-only result of the standalone FULL/retained-two mark phase.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct StandaloneSegmentCompactionPlan {
    pub(crate) archives: Vec<PlannedArchiveSweep>,
    pub(crate) marked_segments: usize,
    pub(super) reclaimable: std::collections::HashSet<SegmentIdentifier>,
}

impl StandaloneSegmentCompactionPlan {
    pub(crate) fn reclaimable_segments(&self) -> &std::collections::HashSet<SegmentIdentifier> {
        &self.reclaimable
    }
}

/// Assembles a comparable plan from a per-archive disposition map.
///
/// Sorted by file name for the same reason the directory-level planner sorts
/// its own: two plans built over the same store must compare equal whatever
/// order the archives were visited in, or the authorization check would refuse
/// runs that agree.
pub(super) fn sorted_sweep_plan(
    planned: &HashMap<String, PlannedArchiveSweep>,
    reclaimable: &std::collections::HashSet<SegmentIdentifier>,
) -> StandaloneSegmentCompactionPlan {
    let mut archives: Vec<PlannedArchiveSweep> = planned.values().cloned().collect();
    archives.sort_by(|left, right| left.file_name().cmp(right.file_name()));
    StandaloneSegmentCompactionPlan {
        archives,
        marked_segments: reclaimable.len(),
        reclaimable: reclaimable.clone(),
    }
}

/// Physical result of applying a standalone segment cleanup.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct StandaloneSegmentCompactionOutcome {
    pub(crate) rewritten_archives: usize,
    pub(crate) removed_archives: usize,
    pub(crate) removed_segments: usize,
    pub(crate) archive_bytes_before: u64,
    pub(crate) archive_bytes_after: u64,
    pub(crate) deletion_failures: Vec<DeferredFileDeletion>,
}

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct DeferredFileDeletion {
    pub(crate) file_name: String,
    pub(crate) error: String,
    pub(crate) target_was_already_absent: bool,
}

/// The observed physical result of one archive sweep attempt.
///
/// `newly_unavailable` is populated only by the mutation branch that proved
/// its unlink or higher-generation publication completed. Callers must use
/// this set, rather than the earlier plan, when filtering graph edges in a
/// later rewrite.
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(super) enum ArchiveSweepDisposition {
    #[default]
    Unchanged,
    Removed,
    Rewritten,
}

#[derive(Debug, Default)]
pub(super) struct ArchiveSweepOutcome {
    pub(super) disposition: ArchiveSweepDisposition,
    pub(super) deletion_failures: Vec<DeferredFileDeletion>,
    pub(super) newly_unavailable: std::collections::HashSet<SegmentIdentifier>,
}

/// Everything a post-compaction reclaim pass needs, including the plan it is
/// authorized to carry out.
#[derive(Clone, Copy)]
pub(crate) struct GenerationReclaimRequest<'sources> {
    /// The generation predicate this pass applies.
    pub(crate) rule: ReclaimRule,
    /// Which archives holding reclaimable segments may be rewritten.
    pub(crate) rewrite_policy: ArchiveRewritePolicy,
    /// A proof the caller already certified these sources under the lock.
    pub(crate) certified_sources: Option<&'sources CertifiedReclaimSources>,
    /// The confirmed plan. When present, the pass replans from disk and
    /// refuses — before it unlinks anything — if the two disagree, so a run
    /// can never mutate an archive its operator did not authorize.
    pub(crate) expected: Option<&'sources StandaloneSegmentCompactionPlan>,
}

/// What a reclaim pass did, so a caller can report it rather than guess.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct SegmentSweepOutcome {
    /// Archives unlinked whole.
    pub(crate) removed_archives: usize,
    /// Archives republished at the next generation letter.
    pub(crate) rewritten_archives: usize,
    /// Segments the sweep made unavailable, by either disposition.
    pub(crate) removed_segments: usize,
    /// Planned unlinks that did not happen, each with its reason.
    pub(crate) deletion_failures: Vec<DeferredFileDeletion>,
}

/// Generations a froe maintenance run retains behind its reference. One
/// value, read by every phase of the run through [`ReclaimRule`].
///
/// This is the value Oak's own offline tooling uses:
/// `SegmentGCOptions.setOffline()` sets `retainedGenerations = 1`
/// (`docs/analysis/write-compaction.md` §4, lines 486 and 531;
/// `write-backup-restore-recovery.md` line 238). The online default of two
/// protects against two things froe does not have: a concurrent writer, and a
/// head that reuses records in place from an older generation. froe's
/// exclusive `repo.lock` excludes the first; a deep copy that rewrites the
/// head into the reference generation excludes the second.
///
/// At this value head safety no longer follows from the predicate — the
/// predicate only decides which *older* generations are reclaimable, and
/// whether the head reaches one is a property of the store. It is proved per
/// run instead, by `validate_reclaim_reference_invariant`, which re-evaluates
/// this exact rule over the head's transitive closure and refuses before any
/// mutation.
pub(crate) const RETAINED_GENERATIONS: i32 = 1;

/// Plans Oak's standalone cleanup predicate: FULL GC, the current committed
/// head generation as reference, and one retained generation. `protected`
/// is a conservative keep-veto for journal history; it never makes a segment
/// reclaimable and therefore cannot weaken Oak's head/checkpoint safety.
pub(crate) fn plan_standalone_segment_cleanup(
    directory: &Path,
    repository: &crate::store::Repository,
    rule: ReclaimRule,
    current_head_segment: SegmentIdentifier,
    protected: &std::collections::HashSet<SegmentIdentifier>,
    rewrite_policy: ArchiveRewritePolicy,
    observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<StandaloneSegmentCompactionPlan> {
    reject_duplicate_active_segments(repository.archives())?;
    certify_active_archives(repository, repository.archives())?;
    analyze_standalone_segment_cleanup(
        directory,
        repository.archives(),
        rule,
        current_head_segment,
        protected,
        rewrite_policy,
        observer,
    )
}

/// Segments and bytes a plan's actionable dispositions would physically
/// free. Deferred and blocked archives free nothing and contribute nothing.
///
/// A whole-file removal frees the archive's own size — index and trailers
/// included — while a rewrite frees only the entry bytes it drops.
pub(crate) fn plan_reclaimed_totals(plan: &StandaloneSegmentCompactionPlan) -> (usize, u64) {
    let mut segments = 0usize;
    let mut bytes = 0u64;
    for archive in &plan.archives {
        let (archive_segments, archive_bytes) = match archive {
            PlannedArchiveSweep::Remove {
                segment_count,
                file_bytes,
                ..
            } => (*segment_count, *file_bytes),
            PlannedArchiveSweep::Rewrite {
                segment_count,
                eligible_entry_bytes,
                ..
            } => (*segment_count, *eligible_entry_bytes),
            PlannedArchiveSweep::DeferredBySavings { .. }
            | PlannedArchiveSweep::DeferredAtLastGeneration { .. }
            | PlannedArchiveSweep::BlockedByOccupiedGeneration { .. } => continue,
        };
        segments = segments.saturating_add(archive_segments);
        bytes = bytes.saturating_add(archive_bytes);
    }
    (segments, bytes)
}

/// Replans the same sweep with the journal-history keep-veto lifted, to
/// price what retiring that history would actually release.
///
/// Reusing the real mark and sweep rather than reasoning about the veto
/// separately is the point: the veto holds bulk segments only indirectly —
/// a vetoed data segment keeps seeding its references — and it interacts
/// with the 25% rewrite gate, since releasing more of an archive can push
/// it over the threshold. Only the sweep itself accounts for both, so any
/// hand-rolled estimate would understate the price of the veto, badly on a
/// store whose history holds inline binaries.
///
/// The caller has already certified these archives for the vetoed plan;
/// this is the mark and sweep alone.
pub(crate) fn measure_unvetoed_reclamation(
    directory: &Path,
    repository: &crate::store::Repository,
    rule: ReclaimRule,
    current_head_segment: SegmentIdentifier,
    rewrite_policy: ArchiveRewritePolicy,
    observer: &mut dyn crate::progress::ProgressObserver,
) -> Result<(usize, u64)> {
    let unvetoed = analyze_standalone_segment_cleanup(
        directory,
        repository.archives(),
        rule,
        current_head_segment,
        &std::collections::HashSet::new(),
        rewrite_policy,
        observer,
    )?;
    Ok(plan_reclaimed_totals(&unvetoed))
}

/// Segment identifiers that the actionable archive dispositions in `plan`
/// would make unavailable. Deferred and blocked archives contribute none.
/// Duplicate identifiers have already been rejected while constructing the
/// plan, so each identifier has exactly one physical active copy.
pub(crate) fn planned_unavailable_segments(
    directory: &Path,
    plan: &StandaloneSegmentCompactionPlan,
) -> Result<std::collections::HashSet<SegmentIdentifier>> {
    let actionable: std::collections::HashSet<&str> = plan
        .archives
        .iter()
        .filter(|archive| archive.changes_disk())
        .map(PlannedArchiveSweep::file_name)
        .collect();
    let archives = crate::store::open_all_archives(directory)?;
    let mut unavailable = std::collections::HashSet::new();
    for archive in archives {
        if !actionable.contains(archive.file_name()) {
            continue;
        }
        let Some(index) = archive.index() else {
            return Err(Error::InvalidFormat {
                details: format!(
                    "cleanup planned to mutate recovered archive {}, which has no valid index",
                    archive.file_name()
                ),
            });
        };
        unavailable.extend(
            index
                .entries()
                .iter()
                .map(|entry| entry.segment_identifier)
                .filter(|identifier| plan.reclaimable.contains(identifier)),
        );
    }
    Ok(unavailable)
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::writer::compaction::CompactionKind;

    use crate::writer::store_writer::reclaim::*;

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

    use crate::writer::store_writer::test_support::*;
    use std::collections::HashSet;

    #[test]
    fn dangling_future_boundary_is_shared_across_archive_files() {
        let directory = TestDirectory::new("dangling-future-cross-archive");
        let older = data_identifier(50);
        let root = data_identifier(51);
        let future = data_identifier(52);
        let current = generation(9, 9, true);
        write_test_archive(
            &directory,
            "data00000a.tar",
            &[
                TestArchiveEntry::new(older, 1, current),
                TestArchiveEntry::new(root, 1, current),
            ],
        );
        write_test_archive(
            &directory,
            "data00001a.tar",
            &[TestArchiveEntry::new(future, 1, current)],
        );
        write_manifest(&directory);

        let plan = plan_cleanup_from_directory(&directory.path, current, root, &HashSet::new())
            .expect("plan across archives");
        assert_eq!(plan.marked_segments, 1);
        assert_eq!(plan.reclaimable_segments(), &HashSet::from([future]));
        assert!(matches!(
            plan.archives.as_slice(),
            [PlannedArchiveSweep::Remove {
                file_name,
                segment_count: 1,
                ..
            }] if file_name == "data00001a.tar"
        ));
        assert!(!plan.reclaimable_segments().contains(&root));
        assert!(!plan.reclaimable_segments().contains(&older));
    }
    /// The authorization the post-compaction sweep never had: a confirmed plan
    /// it must match before it unlinks anything. A disagreement is answered by
    /// refusing, not by explaining afterwards.
    #[test]
    fn a_post_compaction_sweep_refuses_an_unconfirmed_disposition_before_it_unlinks() {
        let (directory, live) = reclaimable_base_fixture("sweep-authorization");
        let before: Vec<String> = crate::store::list_archive_file_names(&directory.path)
            .expect("list the archives before");
        assert!(!before.is_empty());

        // A plan naming a disposition this store cannot produce.
        let impossible = StandaloneSegmentCompactionPlan {
            archives: vec![PlannedArchiveSweep::Remove {
                file_name: "data09999a.tar".to_owned(),
                segment_count: 1,
                file_bytes: 1,
            }],
            marked_segments: 1,
            reclaimable: std::collections::HashSet::new(),
        };

        let mut store = WritableRepository::open(&directory.path).expect("open for the sweep");
        let error = store
            .reclaim_old_generations_with(GenerationReclaimRequest {
                rule: ReclaimRule {
                    reference: live,
                    kind: CompactionKind::Full,
                    retained_generations: RETAINED_GENERATIONS,
                },
                rewrite_policy: ArchiveRewritePolicy::EveryReclaimableArchive,
                certified_sources: None,
                expected: Some(&impossible),
            })
            .expect_err("an unconfirmed disposition must be refused");
        store.close().expect("close after the refusal");

        assert!(
            error.to_string().contains("changed after confirmation"),
            "the refusal names its reason: {error}"
        );
        assert_eq!(
            crate::store::list_archive_file_names(&directory.path)
                .expect("list the archives after"),
            before,
            "a refused sweep unlinks nothing"
        );
    }
    /// The rewrite policy now reaches the code that acts on it. Before this,
    /// a run under Oak's savings heuristic predicted a deferral and then
    /// rewrote the archive anyway.
    #[test]
    fn the_savings_gate_reaches_the_post_compaction_sweep() {
        for (policy, expect_unchanged) in [
            (ArchiveRewritePolicy::OakSavingsGate, true),
            (ArchiveRewritePolicy::EveryReclaimableArchive, false),
        ] {
            let (directory, live) = reclaimable_base_fixture(&format!("sweep-policy-{policy:?}"));
            let before = crate::store::list_archive_file_names(&directory.path).expect("list");

            let mut store = WritableRepository::open(&directory.path).expect("open for the sweep");
            let outcome = store
                .reclaim_old_generations_with(GenerationReclaimRequest {
                    rule: ReclaimRule {
                        reference: live,
                        kind: CompactionKind::Full,
                        retained_generations: RETAINED_GENERATIONS,
                    },
                    rewrite_policy: policy,
                    certified_sources: None,
                    expected: None,
                })
                .expect("the sweep runs");
            store.close().expect("close after the sweep");

            let after = crate::store::list_archive_file_names(&directory.path).expect("list");
            if expect_unchanged {
                assert_eq!(
                    after, before,
                    "Oak's heuristic leaves an archive it will not repay: {outcome:?}"
                );
                assert_eq!(outcome, SegmentSweepOutcome::default());
            } else {
                assert_ne!(
                    after, before,
                    "the default policy reclaims what the heuristic declined"
                );
                assert!(
                    outcome.removed_archives + outcome.rewritten_archives > 0,
                    "and reports what it did: {outcome:?}"
                );
            }
        }
    }
}