rabs-cas 2.1.7

Durable RABS content-addressed storage, action-cache indexing, object lifecycle, and publication transactions
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
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
//! H007 — staging directories and append journals (plan §90; crash
//! recovery for in-flight CAS writes).
//!
//! Layout under the blob-store root:
//!
//! - `staging/<op>/<attempt>/object` — each in-flight write stages in
//!   its OWN per-operation, per-attempt directory, never inside the
//!   published namespace;
//! - `journals/<op>.journal` — one append-only journal per operation
//!   recording every attempt's lifecycle: `begin` (declared identity +
//!   staging path), then exactly one of `published` / `aborted`.
//!
//! Journal records are length-framed and checksummed
//! (`len(u32 be) || payload || blake3_4(payload)`), so a crash mid-
//! append leaves a TORN TAIL that replay detects and stops at — every
//! record before the tail is trusted, nothing after it is guessed.
//! Appends are fsynced before the corresponding filesystem step is
//! considered intent-recorded (write-ahead: `begin` lands before bytes
//! stream, `published` after the atomic link, so replay can always
//! bound what the crashed process may have done).
//!
//! [`recover_operations`] scans every journal against staging/
//! published reality and resolves each non-terminal attempt:
//!
//! - staged bytes VERIFY against the declared identity → **resume**:
//!   publish through the H003 pipeline (`publish_staged`) and record
//!   `published`;
//! - staged bytes missing/partial/wrong → **clean**: remove the
//!   attempt's staging directory and record `aborted`;
//! - terminal attempts keep their outcome; their leftover staging is
//!   swept (the dead-writer orphan case H003's crash tests defer
//!   here).
//!
//! Recovery never appends after a torn tail: once every attempt is
//! resolved its outcome lives in durable reality (published objects +
//! metadata rows), so the journal is RETIRED whole. Running recovery
//! twice is a no-op.

use std::fs;
use std::io::{Read, Write};
use std::path::{Path, PathBuf};

use rabs_protocol::result_identity::TypedDigest;

use crate::blob_store::{
    BlobStoreLayout, DurabilityPolicy, PutError, PutLimits, PutOutcome, io_err, publish_staged,
    recompute_file_digest, stream_to_staging,
};
use crate::metadata_store::{RabsMetadataStore, digest_key};

/// One journal record.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum JournalRecord {
    /// An attempt began staging bytes for `declared_key`.
    Begin {
        /// Attempt id, hex.
        attempt_hex: String,
        /// Declared digest key (`domain:hex`) of the staged object.
        declared_key: String,
    },
    /// The attempt's object was published (atomic link done).
    Published {
        /// Attempt id, hex.
        attempt_hex: String,
    },
    /// The attempt was abandoned; its staging is garbage.
    Aborted {
        /// Attempt id, hex.
        attempt_hex: String,
        /// Why.
        reason: String,
    },
}

impl JournalRecord {
    fn encode(&self) -> String {
        match self {
            Self::Begin {
                attempt_hex,
                declared_key,
            } => format!("begin|{attempt_hex}|{declared_key}"),
            Self::Published { attempt_hex } => format!("published|{attempt_hex}"),
            Self::Aborted {
                attempt_hex,
                reason,
            } => format!("aborted|{attempt_hex}|{reason}"),
        }
    }

    fn decode(payload: &str) -> Option<Self> {
        let mut parts = payload.splitn(3, '|');
        let kind = parts.next()?;
        let attempt_hex = parts.next()?.to_owned();
        match kind {
            "begin" => Some(Self::Begin {
                attempt_hex,
                declared_key: parts.next()?.to_owned(),
            }),
            "published" => Some(Self::Published { attempt_hex }),
            "aborted" => Some(Self::Aborted {
                attempt_hex,
                reason: parts.next().unwrap_or("").to_owned(),
            }),
            _ => None,
        }
    }
}

/// Replay result: every intact record in order, plus whether the file
/// ended in a torn/corrupt tail.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct JournalReplay {
    /// Intact records, journal order.
    pub records: Vec<JournalRecord>,
    /// A truncated or checksum-failed tail was found (and ignored).
    pub torn_tail: bool,
}

fn record_checksum(payload: &[u8]) -> [u8; 4] {
    let digest = blake3::hash(payload);
    let mut out = [0_u8; 4];
    out.copy_from_slice(&digest.as_bytes()[..4]);
    out
}

/// Append one length-framed, checksummed record to an append-only
/// journal file and fsync it (shared by the H007 lifecycle journal and
/// the H008 range journal).
pub(crate) fn append_framed(path: &Path, payload: &[u8]) -> Result<(), PutError> {
    let mut framed = Vec::with_capacity(payload.len() + 8);
    framed.extend_from_slice(&(payload.len() as u32).to_be_bytes());
    framed.extend_from_slice(payload);
    framed.extend_from_slice(&record_checksum(payload));
    let mut file = fs::OpenOptions::new()
        .create(true)
        .append(true)
        .open(path)
        .map_err(io_err("open-journal"))?;
    file.write_all(&framed).map_err(io_err("append-journal"))?;
    file.sync_all().map_err(io_err("fsync-journal"))?;
    Ok(())
}

/// Replay a framed journal, tolerating a torn/corrupt tail: returns
/// the intact payloads in order plus whether a tear was found.
pub(crate) fn replay_framed(path: &Path) -> Result<(Vec<Vec<u8>>, bool), PutError> {
    let bytes = match fs::read(path) {
        Ok(bytes) => bytes,
        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok((Vec::new(), false)),
        Err(e) => return Err(io_err("read-journal")(e)),
    };
    let mut payloads = Vec::new();
    let mut cursor = 0_usize;
    let mut torn_tail = false;
    while cursor < bytes.len() {
        let Some(header) = bytes.get(cursor..cursor + 4) else {
            torn_tail = true;
            break;
        };
        let len = u32::from_be_bytes([header[0], header[1], header[2], header[3]]) as usize;
        let Some(payload) = bytes.get(cursor + 4..cursor + 4 + len) else {
            torn_tail = true;
            break;
        };
        let Some(stored_sum) = bytes.get(cursor + 4 + len..cursor + 8 + len) else {
            torn_tail = true;
            break;
        };
        if stored_sum != record_checksum(payload) {
            torn_tail = true;
            break;
        }
        payloads.push(payload.to_vec());
        cursor += 8 + len;
    }
    Ok((payloads, torn_tail))
}

fn journals_dir(layout: &BlobStoreLayout) -> PathBuf {
    layout.root().join("journals")
}

fn journal_path(layout: &BlobStoreLayout, op_hex: &str) -> PathBuf {
    journals_dir(layout).join(format!("{op_hex}.journal"))
}

fn op_staging_dir(layout: &BlobStoreLayout, op_hex: &str) -> PathBuf {
    layout.root().join("staging").join(op_hex)
}

fn u128_hex(v: u128) -> String {
    format!("{v:032x}")
}

/// Handle to one operation's staging + journal.
#[derive(Debug, Clone)]
pub struct StagingJournal {
    layout: BlobStoreLayout,
    op_hex: String,
}

impl StagingJournal {
    /// Open (creating the journal directory if needed) the journal for
    /// operation `op`.
    ///
    /// # Errors
    /// [`PutError::Io`] when the directory cannot be created.
    pub fn open(layout: &BlobStoreLayout, op: u128) -> Result<Self, PutError> {
        fs::create_dir_all(journals_dir(layout)).map_err(io_err("create-journals-dir"))?;
        Ok(Self {
            layout: layout.clone(),
            op_hex: u128_hex(op),
        })
    }

    /// The attempt's staging directory (`staging/<op>/<attempt>/`).
    #[must_use]
    pub fn attempt_dir(&self, attempt: u128) -> PathBuf {
        op_staging_dir(&self.layout, &self.op_hex).join(u128_hex(attempt))
    }

    /// Append one record and fsync the journal file.
    ///
    /// # Errors
    /// [`PutError::Io`] on append/sync failure.
    pub fn append(&self, record: &JournalRecord) -> Result<(), PutError> {
        append_framed(
            &journal_path(&self.layout, &self.op_hex),
            record.encode().as_bytes(),
        )
    }

    /// Replay the journal, tolerating a torn tail.
    ///
    /// # Errors
    /// [`PutError::Io`] when the journal exists but cannot be read.
    pub fn replay(&self) -> Result<JournalReplay, PutError> {
        replay_file(&journal_path(&self.layout, &self.op_hex))
    }
}

fn replay_file(path: &Path) -> Result<JournalReplay, PutError> {
    let (payloads, mut torn_tail) = replay_framed(path)?;
    let mut records = Vec::new();
    for payload in payloads {
        let Some(record) = std::str::from_utf8(&payload)
            .ok()
            .and_then(JournalRecord::decode)
        else {
            // An intact frame that does not decode is corruption at the
            // record layer: treated exactly like a torn tail — trust
            // nothing from here on.
            torn_tail = true;
            break;
        };
        records.push(record);
    }
    Ok(JournalReplay { records, torn_tail })
}

/// Journaled `put_if_absent`: write-ahead `begin`, stage under the
/// attempt's own directory, verify + publish through the H003
/// pipeline, then record the terminal outcome. Refusals record
/// `aborted` and clean the attempt's staging directory.
///
/// # Errors
/// A typed [`PutError`]; refusals publish nothing.
#[allow(clippy::too_many_arguments)]
pub fn put_if_absent_journaled(
    layout: &BlobStoreLayout,
    store: &mut dyn RabsMetadataStore,
    op: u128,
    attempt: u128,
    declared: &TypedDigest,
    reader: &mut dyn Read,
    limits: PutLimits,
    durability: DurabilityPolicy,
) -> Result<PutOutcome, PutError> {
    let journal = StagingJournal::open(layout, op)?;
    let attempt_hex = u128_hex(attempt);
    let dir = journal.attempt_dir(attempt);
    fs::create_dir_all(&dir).map_err(io_err("create-attempt-dir"))?;
    // Write-ahead intent BEFORE any object bytes exist.
    journal.append(&JournalRecord::Begin {
        attempt_hex: attempt_hex.clone(),
        declared_key: digest_key(declared),
    })?;

    let staging = dir.join("object");
    let outcome = stage_verify_publish(
        layout, store, declared, reader, limits, durability, &staging,
    );
    match &outcome {
        Ok(_) => {
            journal.append(&JournalRecord::Published {
                attempt_hex: attempt_hex.clone(),
            })?;
        }
        Err(e) => {
            journal.append(&JournalRecord::Aborted {
                attempt_hex: attempt_hex.clone(),
                reason: format!("{e:?}"),
            })?;
        }
    }
    let _ = fs::remove_dir_all(&dir);
    outcome
}

fn stage_verify_publish(
    layout: &BlobStoreLayout,
    store: &mut dyn RabsMetadataStore,
    declared: &TypedDigest,
    reader: &mut dyn Read,
    limits: PutLimits,
    durability: DurabilityPolicy,
    staging: &Path,
) -> Result<PutOutcome, PutError> {
    let digests = match stream_to_staging(staging, reader, limits) {
        Ok(digests) => digests,
        Err(e) => {
            let _ = fs::remove_file(staging);
            return Err(e);
        }
    };
    if digests.atp_content_id != *declared {
        let computed = digest_key(&digests.atp_content_id);
        let _ = fs::remove_file(staging);
        return Err(PutError::DeclaredDigestMismatch {
            declared: digest_key(declared),
            computed,
        });
    }
    publish_staged(layout, store, declared, staging, durability)
}

/// What recovery did to one attempt.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RecoveryAction {
    /// Staged bytes verified against the declared identity and were
    /// published.
    Resumed {
        /// Published path.
        path: String,
    },
    /// Staging was missing/partial/wrong (or the attempt was already
    /// terminal); leftovers removed.
    Cleaned,
}

/// The recovery product.
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct RecoveryReport {
    /// Per (op hex, attempt hex) resolution of non-terminal attempts.
    pub resolved: Vec<(String, String, RecoveryAction)>,
    /// Journals that ended in a torn/corrupt tail (rewritten).
    pub torn_journals: Vec<String>,
    /// Journals fully terminal and retired this pass.
    pub retired_journals: Vec<String>,
}

/// Scan every journal, resolve every non-terminal attempt
/// (resume-or-clean), sweep terminal attempts' staging leftovers,
/// rewrite torn journals, and retire fully-terminal ones. Idempotent.
///
/// # Errors
/// [`PutError`] on filesystem/store failures (individual attempt
/// resolutions that legitimately refuse — e.g. collision incidents —
/// are recorded as `aborted`, not surfaced as errors).
pub fn recover_operations(
    layout: &BlobStoreLayout,
    store: &mut dyn RabsMetadataStore,
    durability: DurabilityPolicy,
) -> Result<RecoveryReport, PutError> {
    let mut report = RecoveryReport::default();
    let journals = journals_dir(layout);
    fs::create_dir_all(&journals).map_err(io_err("create-journals-dir"))?;
    let mut journal_files: Vec<PathBuf> = fs::read_dir(&journals)
        .map_err(io_err("read-journals-dir"))?
        .filter_map(|entry| {
            let path = entry.ok()?.path();
            (path.extension().is_some_and(|e| e == "journal")).then_some(path)
        })
        .collect();
    journal_files.sort();

    for path in journal_files {
        let op_hex = path
            .file_stem()
            .map(|s| s.to_string_lossy().into_owned())
            .unwrap_or_default();
        let replay = replay_file(&path)?;
        if replay.torn_tail {
            report.torn_journals.push(op_hex.clone());
        }

        // Fold to per-attempt final state, preserving begin metadata.
        let mut attempts: Vec<(String, Option<String>, bool)> = Vec::new(); // (attempt, declared_key if open, terminal)
        for record in &replay.records {
            match record {
                JournalRecord::Begin {
                    attempt_hex,
                    declared_key,
                } => {
                    if !attempts.iter().any(|(a, _, _)| a == attempt_hex) {
                        attempts.push((attempt_hex.clone(), Some(declared_key.clone()), false));
                    }
                }
                JournalRecord::Published { attempt_hex }
                | JournalRecord::Aborted { attempt_hex, .. } => {
                    if let Some(entry) = attempts.iter_mut().find(|(a, _, _)| a == attempt_hex) {
                        entry.2 = true;
                    }
                }
            }
        }

        for (attempt_hex, declared_key, terminal) in &attempts {
            let attempt_dir = op_staging_dir(layout, &op_hex).join(attempt_hex);
            if *terminal {
                // Dead writer's leftovers (H003's deferred orphans).
                let _ = fs::remove_dir_all(&attempt_dir);
                continue;
            }
            let staged = attempt_dir.join("object");
            let action = resolve_open_attempt(layout, store, declared_key, &staged, durability)?;
            let _ = fs::remove_dir_all(&attempt_dir);
            report
                .resolved
                .push((op_hex.clone(), attempt_hex.clone(), action));
        }

        // Every attempt is now terminal: retire the journal and the
        // operation's staging directory. (Rewrite is unnecessary — the
        // resolved state is fully reflected in durable reality.)
        let _ = fs::remove_dir_all(op_staging_dir(layout, &op_hex));
        fs::remove_file(&path).map_err(io_err("retire-journal"))?;
        report.retired_journals.push(op_hex);
    }
    Ok(report)
}

/// Resolve one open attempt: resume iff the staged bytes verify
/// against the declared identity; clean otherwise. Publication
/// refusals (e.g. a collision incident) resolve as Cleaned — the
/// incident machinery has already preserved the evidence.
fn resolve_open_attempt(
    layout: &BlobStoreLayout,
    store: &mut dyn RabsMetadataStore,
    declared_key: &Option<String>,
    staged: &Path,
    durability: DurabilityPolicy,
) -> Result<RecoveryAction, PutError> {
    let Some(declared_key) = declared_key else {
        return Ok(RecoveryAction::Cleaned);
    };
    if !staged.exists() {
        return Ok(RecoveryAction::Cleaned);
    }
    let recomputed = match recompute_file_digest(staged) {
        Ok(digest) => digest,
        Err(_) => return Ok(RecoveryAction::Cleaned),
    };
    if digest_key(&recomputed) != *declared_key {
        return Ok(RecoveryAction::Cleaned);
    }
    match publish_staged(layout, store, &recomputed, staged, durability) {
        Ok(PutOutcome::Stored { path } | PutOutcome::IdempotentDuplicate { path }) => {
            Ok(RecoveryAction::Resumed { path })
        }
        Err(PutError::CollisionIncident { .. }) => Ok(RecoveryAction::Cleaned),
        Err(e) => Err(e),
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::digest_set::{DigestRequest, digest_set};
    use crate::metadata_store::{RusqliteEngine, SqlMetadataStore};
    use std::sync::atomic::{AtomicU64, Ordering};

    static DIR_COUNTER: AtomicU64 = AtomicU64::new(0);

    fn fresh_layout(tag: &str) -> BlobStoreLayout {
        let n = DIR_COUNTER.fetch_add(1, Ordering::SeqCst);
        let root = std::env::temp_dir().join(format!("rabs-h007-{}-{tag}-{n}", std::process::id()));
        fs::create_dir_all(&root).unwrap();
        BlobStoreLayout::open(&root).unwrap()
    }

    fn store() -> SqlMetadataStore<RusqliteEngine> {
        SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap()
    }

    fn id_of(bytes: &[u8]) -> TypedDigest {
        digest_set(bytes, DigestRequest::default(), None)
            .unwrap()
            .atp_content_id
    }

    #[test]
    fn h007_journaled_put_records_lifecycle_and_cleans_staging() {
        let layout = fresh_layout("basic");
        let mut store = store();
        let bytes = b"journaled object".to_vec();
        let declared = id_of(&bytes);

        let outcome = put_if_absent_journaled(
            &layout,
            &mut store,
            5,
            9,
            &declared,
            &mut bytes.as_slice(),
            PutLimits::default(),
            DurabilityPolicy::FULL,
        )
        .unwrap();
        assert!(matches!(outcome, PutOutcome::Stored { .. }));

        let journal = StagingJournal::open(&layout, 5).unwrap();
        let replay = journal.replay().unwrap();
        assert!(!replay.torn_tail);
        assert_eq!(
            replay.records,
            vec![
                JournalRecord::Begin {
                    attempt_hex: format!("{:032x}", 9),
                    declared_key: digest_key(&declared),
                },
                JournalRecord::Published {
                    attempt_hex: format!("{:032x}", 9),
                },
            ]
        );
        assert!(!journal.attempt_dir(9).exists());

        // A refused put records `aborted`.
        let wrong = id_of(b"different bytes");
        let refused = put_if_absent_journaled(
            &layout,
            &mut store,
            5,
            10,
            &wrong,
            &mut bytes.as_slice(),
            PutLimits::default(),
            DurabilityPolicy::FULL,
        );
        assert!(matches!(
            refused,
            Err(PutError::DeclaredDigestMismatch { .. })
        ));
        let replay = journal.replay().unwrap();
        assert!(matches!(
            replay.records.last(),
            Some(JournalRecord::Aborted { attempt_hex, .. }) if *attempt_hex == format!("{:032x}", 10)
        ));

        // Recovery over a fully-terminal journal just retires it.
        let report = recover_operations(&layout, &mut store, DurabilityPolicy::FULL).unwrap();
        assert!(report.resolved.is_empty());
        assert_eq!(report.retired_journals, vec![format!("{:032x}", 5)]);
        assert!(journal.replay().unwrap().records.is_empty());
    }

    #[test]
    fn h007_crash_mid_stage_recovery_resumes_complete_and_cleans_partial() {
        let layout = fresh_layout("recover");
        let mut store = store();

        // Simulate a writer that died mid-operation: journal has
        // `begin` for two attempts; attempt A staged COMPLETE bytes,
        // attempt B staged a PARTIAL prefix. No terminal records.
        let complete = b"complete staged object".to_vec();
        let declared_a = id_of(&complete);
        let journal = StagingJournal::open(&layout, 7).unwrap();
        let dir_a = journal.attempt_dir(1);
        let dir_b = journal.attempt_dir(2);
        fs::create_dir_all(&dir_a).unwrap();
        fs::create_dir_all(&dir_b).unwrap();
        journal
            .append(&JournalRecord::Begin {
                attempt_hex: format!("{:032x}", 1),
                declared_key: digest_key(&declared_a),
            })
            .unwrap();
        journal
            .append(&JournalRecord::Begin {
                attempt_hex: format!("{:032x}", 2),
                declared_key: digest_key(&declared_a),
            })
            .unwrap();
        fs::write(dir_a.join("object"), &complete).unwrap();
        fs::write(dir_b.join("object"), &complete[..5]).unwrap();

        let report = recover_operations(&layout, &mut store, DurabilityPolicy::FULL).unwrap();
        assert_eq!(report.resolved.len(), 2);
        let action_a = &report
            .resolved
            .iter()
            .find(|(_, a, _)| *a == format!("{:032x}", 1))
            .unwrap()
            .2;
        let RecoveryAction::Resumed { path } = action_a else {
            panic!("complete staging must RESUME, got {action_a:?}");
        };
        assert_eq!(fs::read(path).unwrap(), complete);
        assert!(store.object_located(&declared_a).unwrap());
        let action_b = &report
            .resolved
            .iter()
            .find(|(_, a, _)| *a == format!("{:032x}", 2))
            .unwrap()
            .2;
        assert_eq!(*action_b, RecoveryAction::Cleaned);

        // Staging fully swept, journal retired, second pass a no-op.
        assert!(!op_staging_dir(&layout, &format!("{:032x}", 7)).exists());
        let again = recover_operations(&layout, &mut store, DurabilityPolicy::FULL).unwrap();
        assert!(again.resolved.is_empty() && again.retired_journals.is_empty());
    }

    #[test]
    fn h007_torn_tail_is_detected_and_never_trusted() {
        let layout = fresh_layout("torn");
        let mut store = store();
        let bytes = b"torn tail object".to_vec();
        let declared = id_of(&bytes);
        let journal = StagingJournal::open(&layout, 3).unwrap();
        journal
            .append(&JournalRecord::Begin {
                attempt_hex: format!("{:032x}", 1),
                declared_key: digest_key(&declared),
            })
            .unwrap();
        journal
            .append(&JournalRecord::Published {
                attempt_hex: format!("{:032x}", 1),
            })
            .unwrap();

        // Crash mid-append: a truncated frame, then (separately) a
        // checksum-corrupt frame.
        let path = journals_dir(&layout).join(format!("{:032x}.journal", 3));
        let intact = fs::read(&path).unwrap();
        for tail in [
            vec![0, 0, 0, 42, b'p', b'a', b'r'], // truncated payload
            {
                let payload = b"aborted|deadbeef|x".to_vec();
                let mut frame = (payload.len() as u32).to_be_bytes().to_vec();
                frame.extend_from_slice(&payload);
                frame.extend_from_slice(&[0, 0, 0, 0]); // wrong checksum
                frame
            },
        ] {
            let mut torn = intact.clone();
            torn.extend_from_slice(&tail);
            fs::write(&path, &torn).unwrap();
            let replay = journal.replay().unwrap();
            assert!(replay.torn_tail, "tail {tail:?} must be detected");
            assert_eq!(replay.records.len(), 2, "intact prefix fully trusted");
        }

        // Recovery over the torn journal reports it and still resolves
        // cleanly (all recorded attempts are terminal).
        let report = recover_operations(&layout, &mut store, DurabilityPolicy::FULL).unwrap();
        assert_eq!(report.torn_journals, vec![format!("{:032x}", 3)]);
        assert_eq!(report.retired_journals, vec![format!("{:032x}", 3)]);
    }

    #[test]
    fn h007_recovery_is_reconstructible_at_every_crash_boundary() {
        // Crash the journaled put at each lifecycle boundary by
        // REPLAYING the exact on-disk states it passes through, and
        // assert recovery resolves every one without partial exposure.
        let bytes = b"boundary object".to_vec();
        struct Boundary {
            name: &'static str,
            stage_bytes: Option<&'static [u8]>, // staged file content at crash
            begin_recorded: bool,
        }
        let boundaries = [
            Boundary {
                name: "after-begin-no-bytes",
                stage_bytes: None,
                begin_recorded: true,
            },
            Boundary {
                name: "after-partial-stage",
                stage_bytes: Some(b"boundary"),
                begin_recorded: true,
            },
            Boundary {
                name: "after-full-stage",
                stage_bytes: Some(b"boundary object"),
                begin_recorded: true,
            },
        ];
        for boundary in boundaries {
            let layout = fresh_layout(boundary.name);
            let mut store = store();
            let declared = id_of(&bytes);
            let journal = StagingJournal::open(&layout, 11).unwrap();
            let dir = journal.attempt_dir(1);
            fs::create_dir_all(&dir).unwrap();
            if boundary.begin_recorded {
                journal
                    .append(&JournalRecord::Begin {
                        attempt_hex: format!("{:032x}", 1),
                        declared_key: digest_key(&declared),
                    })
                    .unwrap();
            }
            if let Some(staged) = boundary.stage_bytes {
                fs::write(dir.join("object"), staged).unwrap();
            }

            let report = recover_operations(&layout, &mut store, DurabilityPolicy::FULL).unwrap();
            assert_eq!(report.resolved.len(), 1, "{}", boundary.name);
            match &report.resolved[0].2 {
                RecoveryAction::Resumed { path } => {
                    assert_eq!(
                        fs::read(path).unwrap(),
                        bytes,
                        "{}: resumed object must be complete",
                        boundary.name
                    );
                    assert!(store.object_located(&declared).unwrap());
                }
                RecoveryAction::Cleaned => {
                    assert!(
                        !store.object_located(&declared).unwrap(),
                        "{}: cleaned attempt must not be recorded",
                        boundary.name
                    );
                }
            }
            assert!(!dir.exists(), "{}: staging swept", boundary.name);
            assert!(
                journal.replay().unwrap().records.is_empty(),
                "{}: journal retired",
                boundary.name
            );
        }
    }
}