Skip to main content

fxfs/object_store/
journal.rs

1// Copyright 2021 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5//! The journal is implemented as an ever extending file which contains variable length records
6//! that describe mutations to be applied to various objects.  The journal file consists of
7//! blocks, with a checksum at the end of each block, but otherwise it can be considered a
8//! continuous stream.
9//!
10//! The checksum is seeded with the checksum from the previous block.  To free space in the
11//! journal, records are replaced with sparse extents when it is known they are no longer
12//! needed to mount.
13//!
14//! At mount time, the journal is replayed: the mutations are applied into memory.
15//! Eventually, a checksum failure will indicate no more records exist to be replayed,
16//! at which point the mount can continue and the journal will be extended from that point with
17//! further mutations as required.
18
19mod bootstrap_handle;
20mod checksum_list;
21mod reader;
22pub mod super_block;
23mod writer;
24
25use crate::checksum::{Checksum, Checksums};
26use crate::errors::FxfsError;
27use crate::filesystem::{
28    ApplyContext, ApplyMode, FlushReason, ForceMajor, FxFilesystem, SyncOptions,
29};
30use crate::hooks::HooksHandle;
31use crate::log::*;
32use crate::lsm_tree::types::LayerIterator;
33use crate::object_handle::{ObjectHandle as _, ReadObjectHandle};
34use crate::object_store::allocator::Allocator;
35use crate::object_store::data_object_handle::OverwriteOptions;
36use crate::object_store::extent_record::{ExtentMode, ExtentValue};
37use crate::object_store::graveyard::Graveyard;
38use crate::object_store::journal::bootstrap_handle::BootstrapObjectHandle;
39use crate::object_store::journal::checksum_list::ChecksumList;
40use crate::object_store::journal::reader::{JournalReader, ReadResult};
41use crate::object_store::journal::super_block::{
42    SuperBlockHeader, SuperBlockInstance, SuperBlockManager,
43};
44use crate::object_store::journal::writer::JournalWriter;
45use crate::object_store::object_manager::ObjectManager;
46use crate::object_store::object_record::{AttributeKey, ObjectKey, ObjectKeyData, ObjectValue};
47use crate::object_store::transaction::{
48    AllocatorMutation, LockKey, Mutation, MutationV56, MutationV57, MutationV59,
49    ObjectMutationIterator, ObjectStoreMutation, Options, ReservationOptions,
50    TRANSACTION_MAX_JOURNAL_USAGE, Transaction, lock_keys,
51};
52use crate::object_store::{
53    AssocObj, AttributeId, DataObjectHandle, Extent, HandleOptions, HandleOwner, INVALID_OBJECT_ID,
54    Item, ItemRef, NewChildStoreOptions, ObjectStore, ReservedId,
55};
56use crate::range::RangeExt;
57use crate::round::round_div;
58use crate::serialized_types::{LATEST_VERSION, Migrate, Version, Versioned, migrate_to_version};
59use anyhow::{Context, Error, anyhow, bail, ensure};
60use core::iter::Iterator;
61use event_listener::Event;
62use fprint::TypeFingerprint;
63use fuchsia_inspect::NumericProperty;
64use fuchsia_sync::Mutex;
65use futures::FutureExt as _;
66use futures::future::poll_fn;
67use rustc_hash::FxHashMap as HashMap;
68use serde::{Deserialize, Serialize};
69use static_assertions::const_assert;
70use std::clone::Clone;
71use std::collections::HashSet;
72use std::num::NonZero;
73use std::ops::{Bound, Range};
74use std::sync::atomic::{AtomicBool, Ordering};
75use std::sync::{Arc, OnceLock};
76use std::task::{Poll, Waker};
77use storage_device::Device;
78use storage_units::BlockSize;
79
80// The journal file is written to in blocks of this size.
81pub const BLOCK_SIZE: BlockSize = BlockSize::SIZE_4KIB;
82
83// The journal file is extended by this amount when necessary.
84const CHUNK_SIZE: u64 = 131_072;
85const_assert!(CHUNK_SIZE > TRANSACTION_MAX_JOURNAL_USAGE);
86
87// See the comment for the `reclaim_size` member of Inner.
88pub const DEFAULT_RECLAIM_SIZE: u64 = 524_288;
89
90// Temporary space that should be reserved for the journal.  For example: space that is currently
91// used in the journal file but cannot be deallocated yet because we are flushing.
92pub const RESERVED_SPACE: u64 = 1_048_576;
93
94// Whenever the journal is replayed (i.e. the system is unmounted and remounted), we reset the
95// journal stream, at which point any half-complete transactions are discarded.  We indicate a
96// journal reset by XORing the previous block's checksum with this mask, and using that value as a
97// seed for the next journal block.
98const RESET_XOR: u64 = 0xffffffffffffffff;
99
100// To keep track of offsets within a journal file, we need both the file offset and the check-sum of
101// the preceding block, since the check-sum of the preceding block is an input to the check-sum of
102// every block.
103pub type JournalCheckpoint = JournalCheckpointV32;
104
105#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize, TypeFingerprint)]
106pub struct JournalCheckpointV32 {
107    pub file_offset: u64,
108
109    // Starting check-sum for block that contains file_offset i.e. the checksum for the previous
110    // block.
111    pub checksum: Checksum,
112
113    // If versioned, the version of elements stored in the journal. e.g. JournalRecord version.
114    // This can change across reset events so we store it along with the offset and checksum to
115    // know which version to deserialize.
116    pub version: Version,
117}
118
119pub type JournalRecord = JournalRecordV59;
120
121#[allow(clippy::large_enum_variant)]
122#[derive(Clone, Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
123#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
124pub enum JournalRecordV59 {
125    EndBlock,
126    Mutation {
127        object_id: u64,
128        mutation: MutationV59,
129    },
130    /// Commits records in the transaction.
131    Commit,
132    /// Discard all mutations with offsets greater than or equal to the given offset.
133    Discard(u64),
134    /// Indicates the device was flushed at the given journal offset.
135    /// Note that this really means that at this point in the journal offset, we can be certain that
136    /// there's no remaining buffered data in the block device; the buffers and the disk contents
137    /// are consistent.
138    /// We insert one of these records *after* a flush along with the *next* transaction to go
139    /// through.  If that never comes (either due to graceful or hard shutdown), the journal reset
140    /// on the next mount will serve the same purpose and count as a flush, although it is necessary
141    /// to defensively flush the device before replaying the journal (if possible, i.e. not
142    /// read-only) in case the block device connection was reused.
143    DidFlushDevice(u64),
144    /// Checksums for a data range written by this transaction. A transaction is only valid if these
145    /// checksums are right. The range is the device offset the checksums are for.
146    ///
147    /// A boolean indicates whether this range is being written to for the first time. For overwrite
148    /// extents, we only check the checksums for a block if it has been written to for the first
149    /// time since the last flush, because otherwise we can't roll it back anyway so it doesn't
150    /// matter. For copy-on-write extents, the bool is always true.
151    DataChecksums(Range<u64>, crate::checksum::ChecksumsV38, bool),
152}
153
154#[allow(clippy::large_enum_variant)]
155#[derive(Migrate, Clone, Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
156#[migrate_to_version(JournalRecordV59)]
157pub enum JournalRecordV57 {
158    EndBlock,
159    Mutation { object_id: u64, mutation: MutationV57 },
160    Commit,
161    Discard(u64),
162    DidFlushDevice(u64),
163    DataChecksums(Range<u64>, crate::checksum::ChecksumsV38, bool),
164}
165
166#[allow(clippy::large_enum_variant)]
167#[derive(Migrate, Clone, Debug, Serialize, Deserialize, TypeFingerprint, Versioned)]
168#[migrate_to_version(JournalRecordV57)]
169pub enum JournalRecordV56 {
170    EndBlock,
171    Mutation { object_id: u64, mutation: MutationV56 },
172    Commit,
173    Discard(u64),
174    DidFlushDevice(u64),
175    DataChecksums(Range<u64>, crate::checksum::ChecksumsV38, bool),
176}
177
178pub(super) fn journal_handle_options() -> HandleOptions {
179    HandleOptions { skip_journal_checks: true, ..Default::default() }
180}
181
182/// The journal records a stream of mutations that are to be applied to other objects.  At mount
183/// time, these records can be replayed into memory.  It provides a way to quickly persist changes
184/// without having to make a large number of writes; they can be deferred to a later time (e.g.
185/// when a sufficient number have been queued).  It also provides support for transactions, the
186/// ability to have mutations that are to be applied atomically together.
187pub struct Journal {
188    objects: Arc<ObjectManager>,
189    hooks: Arc<HooksHandle>,
190    handle: OnceLock<DataObjectHandle<ObjectStore>>,
191    super_block_manager: SuperBlockManager,
192    inner: Mutex<Inner>,
193    writer_mutex: Mutex<()>,
194    sync_mutex: futures::lock::Mutex<()>,
195    trace: AtomicBool,
196
197    // This event is used when we are waiting for a compaction to free up journal space.
198    reclaim_event: Event,
199}
200
201struct Inner {
202    super_block_header: SuperBlockHeader,
203
204    // The offset that we can zero the journal up to now that it is no longer needed.
205    zero_offset: Option<u64>,
206
207    // The journal offset that we most recently flushed to the device.
208    device_flushed_offset: u64,
209
210    // If true, indicates a DidFlushDevice record is pending.
211    needs_did_flush_device: bool,
212
213    // The writer for the journal.
214    writer: JournalWriter,
215
216    // Set when a reset is encountered during a read.
217    // Used at write pre_commit() time to ensure we write a version first thing after a reset.
218    output_reset_version: bool,
219
220    // Waker for the flush task.
221    flush_waker: Option<Waker>,
222
223    // Indicates the journal has been terminated.
224    terminate: bool,
225
226    // Latched error indicating reason for journal termination if not graceful.
227    terminate_reason: Option<Error>,
228
229    // Disable compactions.
230    disable_compactions: bool,
231
232    // When true, compactions are paused.
233    compactions_paused: bool,
234
235    // True if compactions are running.
236    compaction_running: bool,
237
238    // Waker for the sync task for when it's waiting for the flush task to finish.
239    sync_waker: Option<Waker>,
240
241    // The last offset we flushed to the journal file.
242    flushed_offset: u64,
243
244    // The last offset that should be considered valid in the journal file.  Most of the time, this
245    // will be the same as `flushed_offset` but at mount time, this could be less and will only be
246    // up to the end of the last valid transaction; it won't include transactions that follow that
247    // have been discarded.
248    valid_to: u64,
249
250    // If, after replaying, we have to discard a number of mutations (because they don't validate),
251    // this offset specifies where we need to discard back to.  This is so that when we next replay,
252    // we ignore those mutations and continue with new good mutations.
253    discard_offset: Option<u64>,
254
255    // In the steady state, the journal should fluctuate between being approximately half of this
256    // number and this number.  New super-blocks will be written every time about half of this
257    // amount is written to the journal.
258    reclaim_size: u64,
259
260    image_builder_mode: Option<SuperBlockInstance>,
261
262    // If true and `needs_barrier`, issue a pre-barrier on the first device write of each journal
263    // write (which happens in multiples of `BLOCK_SIZE`). This ensures that all the corresponding
264    // data writes make it to disk before the journal gets written to.
265    barriers_enabled: bool,
266
267    // If true, indicates that data write requests have been made to the device since the last
268    // journal write.
269    needs_barrier: bool,
270
271    // True if a compaction is being forced for reasons other than the journal being full.
272    forced_compaction: bool,
273}
274
275impl Inner {
276    fn terminate(&mut self, reason: Option<Error>) {
277        self.terminate = true;
278
279        if let Some(err) = reason {
280            error!(error:? = err; "Terminating journal");
281            // Log previous error if one was already set, otherwise latch the error.
282            if let Some(prev_err) = self.terminate_reason.as_ref() {
283                error!(error:? = prev_err; "Journal previously terminated");
284            } else {
285                self.terminate_reason = Some(err);
286            }
287        }
288
289        if let Some(waker) = self.flush_waker.take() {
290            waker.wake();
291        }
292        if let Some(waker) = self.sync_waker.take() {
293            waker.wake();
294        }
295    }
296}
297
298pub struct JournalOptions {
299    /// In the steady state, the journal should fluctuate between being approximately half of this
300    /// number and this number.  New super-blocks will be written every time about half of this
301    /// amount is written to the journal.
302    pub reclaim_size: u64,
303
304    // If true, issue a pre-barrier on the first device write of each journal write (which happens
305    // in multiples of `BLOCK_SIZE`). This ensures that all the corresponding data writes make it
306    // to disk before the journal gets written to.
307    pub barriers_enabled: bool,
308}
309
310impl Default for JournalOptions {
311    fn default() -> Self {
312        JournalOptions { reclaim_size: DEFAULT_RECLAIM_SIZE, barriers_enabled: false }
313    }
314}
315
316struct JournaledTransactions {
317    transactions: Vec<JournaledTransaction>,
318    device_flushed_offset: u64,
319}
320
321#[derive(Debug, Default)]
322pub struct JournaledTransaction {
323    pub checkpoint: JournalCheckpoint,
324    pub root_parent_mutations: Vec<Mutation>,
325    pub root_mutations: Vec<Mutation>,
326    /// List of (store_object_id, mutation).
327    pub non_root_mutations: Vec<(u64, Mutation)>,
328    pub end_offset: u64,
329    pub checksums: Vec<JournaledChecksums>,
330
331    /// Records offset + 1 of the matching begin_flush transaction. The +1 is because we want to
332    /// ignore the begin flush transaction; we don't need or want to replay it.
333    pub end_flush: Option<(/* store_id: */ u64, /* begin offset: */ u64)>,
334
335    /// The volume which was deleted in this transaction, if any.
336    pub volume_deleted: Option</* store_id: */ u64>,
337}
338
339impl JournaledTransaction {
340    fn new(checkpoint: JournalCheckpoint) -> Self {
341        Self { checkpoint, ..Default::default() }
342    }
343}
344
345const VOLUME_DELETED: u64 = u64::MAX;
346
347#[derive(Debug)]
348pub struct JournaledChecksums {
349    pub device_range: Range<u64>,
350    pub checksums: Checksums,
351    pub first_write: bool,
352}
353
354/// Handles for journal-like objects have some additional functionality to manage their extents,
355/// since during replay we need to add extents as we find them.
356pub trait JournalHandle: ReadObjectHandle {
357    /// The end offset of the last extent in the JournalHandle.  Used only for validating extents
358    /// (which will be skipped if None is returned).
359    /// Note this is equivalent in value to ReadObjectHandle::get_size, when present.
360    fn end_offset(&self) -> Option<u64>;
361    /// Adds an extent to the current end of the journal stream.
362    /// `added_offset` is the offset into the journal of the transaction which added this extent,
363    /// used for discard_extents.
364    fn push_extent(&mut self, added_offset: u64, device_range: Range<u64>);
365    /// Discards all extents which were added in a transaction at offset >= |discard_offset|.
366    fn discard_extents(&mut self, discard_offset: u64);
367}
368
369// Provide a stub implementation for DataObjectHandle so we can use it in
370// Journal::read_transactions.  Manual extent management is a NOP (which is OK since presumably the
371// DataObjectHandle already knows where its extents live).
372impl<S: HandleOwner> JournalHandle for DataObjectHandle<S> {
373    fn end_offset(&self) -> Option<u64> {
374        None
375    }
376    fn push_extent(&mut self, _added_offset: u64, _device_range: Range<u64>) {
377        // NOP
378    }
379    fn discard_extents(&mut self, _discard_offset: u64) {
380        // NOP
381    }
382}
383
384#[fxfs_trace::trace]
385impl Journal {
386    pub fn new(
387        objects: Arc<ObjectManager>,
388        options: JournalOptions,
389        hooks: Arc<HooksHandle>,
390    ) -> Journal {
391        let starting_checksum = rand::random_range(1..u64::MAX);
392        Journal {
393            objects: objects,
394            hooks,
395            handle: OnceLock::new(),
396            super_block_manager: SuperBlockManager::new(),
397            inner: Mutex::new(Inner {
398                super_block_header: SuperBlockHeader::default(),
399                zero_offset: None,
400                device_flushed_offset: 0,
401                needs_did_flush_device: false,
402                writer: JournalWriter::new(BLOCK_SIZE, starting_checksum),
403                output_reset_version: false,
404                flush_waker: None,
405                terminate: false,
406                terminate_reason: None,
407                disable_compactions: false,
408                compactions_paused: false,
409                compaction_running: false,
410                sync_waker: None,
411                flushed_offset: 0,
412                valid_to: 0,
413                discard_offset: None,
414                reclaim_size: options.reclaim_size,
415                image_builder_mode: None,
416                barriers_enabled: options.barriers_enabled,
417                needs_barrier: false,
418                forced_compaction: false,
419            }),
420            writer_mutex: Mutex::new(()),
421            sync_mutex: futures::lock::Mutex::new(()),
422            trace: AtomicBool::new(false),
423            reclaim_event: Event::new(),
424        }
425    }
426
427    pub fn set_trace(&self, trace: bool) {
428        let old_value = self.trace.swap(trace, Ordering::Relaxed);
429        if trace != old_value {
430            info!(trace; "J: trace");
431        }
432    }
433
434    pub fn set_image_builder_mode(&self, mode: Option<SuperBlockInstance>) {
435        self.inner.lock().image_builder_mode = mode;
436        if let Some(instance) = mode {
437            *self.super_block_manager.next_instance.lock() = instance;
438        }
439    }
440
441    pub fn image_builder_mode(&self) -> Option<SuperBlockInstance> {
442        self.inner.lock().image_builder_mode
443    }
444
445    #[cfg(feature = "migration")]
446    pub fn set_filesystem_uuid(&self, uuid: &[u8; 16]) -> Result<(), Error> {
447        ensure!(
448            self.inner.lock().image_builder_mode.is_some(),
449            "Can only set filesystem uuid in image builder mode."
450        );
451        self.inner.lock().super_block_header.guid.0 = uuid::Uuid::from_bytes(*uuid);
452        Ok(())
453    }
454
455    pub(crate) async fn read_superblocks(
456        &self,
457        device: Arc<dyn Device>,
458        block_size: BlockSize,
459    ) -> Result<(SuperBlockHeader, ObjectStore), Error> {
460        self.super_block_manager.load(device, block_size).await
461    }
462
463    /// Used during replay to validate a mutation.  This should return false if the mutation is not
464    /// valid and should not be applied.  This could be for benign reasons: e.g. the device flushed
465    /// data out-of-order, or because of a malicious actor.
466    fn validate_mutation(
467        &self,
468        mutation: &Mutation,
469        block_size: BlockSize,
470        device_size: u64,
471    ) -> bool {
472        match mutation {
473            Mutation::ObjectStore(ObjectStoreMutation {
474                item:
475                    Item {
476                        key:
477                            ObjectKey {
478                                data: ObjectKeyData::Attribute(_, AttributeKey::Extent(extent)),
479                                ..
480                            },
481                        value: ObjectValue::Extent(ExtentValue::Some { device_offset, mode, .. }),
482                        ..
483                    },
484                ..
485            }) => {
486                if extent.is_empty() || !block_size.is_aligned(extent) {
487                    return false;
488                }
489                let len = extent.length().unwrap();
490                if let ExtentMode::Cow(checksums) = mode {
491                    if checksums.len() > 0 {
492                        if len % checksums.len() as u64 != 0 {
493                            return false;
494                        }
495                        if !block_size.is_aligned(len / checksums.len() as u64) {
496                            return false;
497                        }
498                    }
499                }
500                if !block_size.is_aligned(device_offset)
501                    || *device_offset >= device_size
502                    || device_size - *device_offset < len
503                {
504                    return false;
505                }
506            }
507            Mutation::ObjectStore(_) => {}
508            Mutation::EncryptedObjectStore(_) => {}
509            Mutation::Allocator(AllocatorMutation::Allocate { device_range, owner_object_id }) => {
510                return !device_range.is_empty()
511                    && *owner_object_id != INVALID_OBJECT_ID
512                    && device_range.end <= device_size;
513            }
514            Mutation::Allocator(AllocatorMutation::Deallocate {
515                device_range,
516                owner_object_id,
517            }) => {
518                return !device_range.is_empty()
519                    && *owner_object_id != INVALID_OBJECT_ID
520                    && device_range.end <= device_size;
521            }
522            Mutation::Allocator(AllocatorMutation::MarkForDeletion(owner_object_id)) => {
523                return *owner_object_id != INVALID_OBJECT_ID;
524            }
525            Mutation::Allocator(AllocatorMutation::SetLimit { owner_object_id, .. }) => {
526                return *owner_object_id != INVALID_OBJECT_ID;
527            }
528            Mutation::BeginFlush => {}
529            Mutation::EndFlush => {}
530            Mutation::DeleteVolume => {}
531            Mutation::UpdateBorrowed(_) => {}
532            Mutation::UpdateMutationsKey(_) => {}
533            Mutation::CreateInternalDir(owner_object_id) => {
534                return *owner_object_id != INVALID_OBJECT_ID;
535            }
536        }
537        true
538    }
539
540    // Assumes that `mutation` has been validated.
541    fn update_checksum_list(
542        &self,
543        journal_offset: u64,
544        mutation: &Mutation,
545        checksum_list: &mut ChecksumList,
546    ) -> Result<(), Error> {
547        match mutation {
548            Mutation::ObjectStore(_) => {}
549            Mutation::Allocator(AllocatorMutation::Deallocate { device_range, .. }) => {
550                checksum_list.mark_deallocated(journal_offset, device_range.clone().into());
551            }
552            _ => {}
553        }
554        Ok(())
555    }
556
557    /// Reads the latest super-block, and then replays journaled records.
558    #[trace]
559    pub async fn replay(
560        &self,
561        filesystem: Arc<FxFilesystem>,
562        on_new_allocator: Option<Box<dyn Fn(Arc<Allocator>) + Send + Sync>>,
563    ) -> Result<(), Error> {
564        let block_size = filesystem.block_size();
565
566        let (super_block, root_parent) =
567            self.super_block_manager.load(filesystem.device(), block_size).await?;
568
569        let root_parent = Arc::new(ObjectStore::attach_filesystem(root_parent, filesystem.clone()));
570
571        self.objects.set_root_parent_store(root_parent.clone());
572        let allocator =
573            Arc::new(Allocator::new(filesystem.clone(), super_block.allocator_object_id));
574        if let Some(on_new_allocator) = on_new_allocator {
575            on_new_allocator(allocator.clone());
576        }
577        self.objects.set_allocator(allocator.clone());
578        self.objects.set_borrowed_metadata_space(super_block.borrowed_metadata_space);
579        self.objects.set_last_end_offset(super_block.super_block_journal_file_offset);
580        {
581            let mut inner = self.inner.lock();
582            inner.super_block_header = super_block.clone();
583        }
584
585        let device = filesystem.device();
586
587        let mut handle;
588        {
589            let root_parent_layer = root_parent.tree().mutable_layer();
590            let mut iter = root_parent_layer.seek(Bound::Included(&ObjectKey::attribute(
591                super_block.journal_object_id,
592                AttributeId::DATA,
593                AttributeKey::Extent(Extent::search_key_from_offset(
594                    BLOCK_SIZE.align_down(super_block.journal_checkpoint.file_offset),
595                )),
596            )));
597            let start_offset = if let Some(ItemRef {
598                key:
599                    ObjectKey {
600                        data:
601                            ObjectKeyData::Attribute(AttributeId::DATA, AttributeKey::Extent(extent)),
602                        ..
603                    },
604                ..
605            }) = iter.get()
606            {
607                extent.start
608            } else {
609                0
610            };
611            handle = BootstrapObjectHandle::new_with_start_offset(
612                super_block.journal_object_id,
613                device.clone(),
614                start_offset,
615            );
616            while let Some(item) = iter.get() {
617                if !match item.into() {
618                    Some((
619                        object_id,
620                        AttributeId::DATA,
621                        extent,
622                        ExtentValue::Some { device_offset, .. },
623                    )) if object_id == super_block.journal_object_id => {
624                        if let Some(end_offset) = handle.end_offset() {
625                            if extent.start != end_offset {
626                                bail!(anyhow!(FxfsError::Inconsistent).context(format!(
627                                    "Unexpected journal extent {:?}, expected start: {}",
628                                    item, end_offset
629                                )));
630                            }
631                        }
632                        handle.push_extent(
633                            0, // We never discard extents from the root parent store.
634                            *device_offset
635                                ..*device_offset + extent.length().context("Invalid extent")?,
636                        );
637                        true
638                    }
639                    _ => false,
640                } {
641                    break;
642                }
643                iter.advance().await.context("Failed to advance root parent store iterator")?;
644            }
645        }
646
647        let mut reader = JournalReader::new(handle, &super_block.journal_checkpoint);
648        let JournaledTransactions { mut transactions, device_flushed_offset } = self
649            .read_transactions(&mut reader, None, INVALID_OBJECT_ID)
650            .await
651            .context("Reading transactions for replay")?;
652
653        // Validate all the mutations.
654        let mut checksum_list = ChecksumList::new(device_flushed_offset);
655        let mut valid_to = reader.journal_file_checkpoint().file_offset;
656        let device_size = device.size();
657        'bad_replay: for JournaledTransaction {
658            checkpoint,
659            root_parent_mutations,
660            root_mutations,
661            non_root_mutations,
662            checksums,
663            ..
664        } in &transactions
665        {
666            for JournaledChecksums { device_range, checksums, first_write } in checksums {
667                checksum_list
668                    .push(
669                        checkpoint.file_offset,
670                        device_range.clone(),
671                        checksums.maybe_as_ref().context("Malformed checksums")?,
672                        *first_write,
673                    )
674                    .context("Pushing journal checksum records to checksum list")?;
675            }
676            for mutation in root_parent_mutations
677                .iter()
678                .chain(root_mutations)
679                .chain(non_root_mutations.iter().map(|(_, m)| m))
680            {
681                if !self.validate_mutation(mutation, block_size, device_size) {
682                    info!(mutation:?; "Stopping replay at bad mutation");
683                    valid_to = checkpoint.file_offset;
684                    break 'bad_replay;
685                }
686                self.update_checksum_list(checkpoint.file_offset, &mutation, &mut checksum_list)?;
687            }
688        }
689
690        // Validate the checksums. Note if barriers are enabled, there will be no checksums in
691        // practice to verify.
692        let valid_to = checksum_list
693            .verify(device.as_ref(), valid_to)
694            .await
695            .context("Failed to validate checksums")?;
696
697        // Apply the mutations...
698
699        let mut last_checkpoint = reader.journal_file_checkpoint();
700        let mut journal_offsets = super_block.journal_file_offsets.clone();
701
702        // Start with the root-parent mutations, and also determine the journal offsets for all
703        // other objects.
704        for (
705            index,
706            JournaledTransaction {
707                checkpoint,
708                root_parent_mutations,
709                end_flush,
710                volume_deleted,
711                ..
712            },
713        ) in transactions.iter_mut().enumerate()
714        {
715            if checkpoint.file_offset >= valid_to {
716                last_checkpoint = checkpoint.clone();
717
718                // Truncate the transactions so we don't need to worry about them on the next pass.
719                transactions.truncate(index);
720                break;
721            }
722
723            let context = ApplyContext { mode: ApplyMode::Replay, checkpoint: checkpoint.clone() };
724            for mutation in root_parent_mutations.drain(..) {
725                self.objects
726                    .apply_mutation(
727                        super_block.root_parent_store_object_id,
728                        mutation,
729                        &context,
730                        AssocObj::None,
731                    )
732                    .context("Failed to replay root parent store mutations")?;
733            }
734
735            if let Some((object_id, journal_offset)) = end_flush {
736                journal_offsets.insert(*object_id, *journal_offset);
737            }
738
739            if let Some(object_id) = volume_deleted {
740                journal_offsets.insert(*object_id, VOLUME_DELETED);
741            }
742        }
743
744        // Now we can open the root store.
745        let root_store = ObjectStore::open(&root_parent, super_block.root_store_object_id, None)
746            .await
747            .context("Unable to open root store")?;
748
749        ensure!(
750            !root_store.is_encrypted(),
751            anyhow!(FxfsError::Inconsistent).context("Root store is encrypted")
752        );
753        self.objects.set_root_store(root_store);
754
755        let root_store_offset =
756            journal_offsets.get(&super_block.root_store_object_id).copied().unwrap_or(0);
757
758        // Now replay the root store mutations.
759        for JournaledTransaction { checkpoint, root_mutations, .. } in &mut transactions {
760            if checkpoint.file_offset < root_store_offset {
761                continue;
762            }
763
764            let context = ApplyContext { mode: ApplyMode::Replay, checkpoint: checkpoint.clone() };
765            for mutation in root_mutations.drain(..) {
766                self.objects
767                    .apply_mutation(
768                        super_block.root_store_object_id,
769                        mutation,
770                        &context,
771                        AssocObj::None,
772                    )
773                    .context("Failed to replay root store mutations")?;
774            }
775        }
776
777        // Now we can open the allocator.
778        allocator.open().await.context("Failed to open allocator")?;
779
780        // Now replay all other mutations.
781        for JournaledTransaction { checkpoint, non_root_mutations, end_offset, .. } in transactions
782        {
783            self.objects
784                .replay_mutations(
785                    non_root_mutations,
786                    &journal_offsets,
787                    &ApplyContext { mode: ApplyMode::Replay, checkpoint },
788                    end_offset,
789                )
790                .await
791                .context("Failed to replay mutations")?;
792        }
793
794        allocator.on_replay_complete().await.context("Failed to complete replay for allocator")?;
795
796        let discarded_to =
797            if last_checkpoint.file_offset != reader.journal_file_checkpoint().file_offset {
798                Some(reader.journal_file_checkpoint().file_offset)
799            } else {
800                None
801            };
802
803        // Configure the journal writer so that we can continue.
804        {
805            if last_checkpoint.file_offset < super_block.super_block_journal_file_offset {
806                return Err(anyhow!(FxfsError::Inconsistent).context(format!(
807                    "journal replay cut short; journal finishes at {}, but super-block was \
808                     written at {}",
809                    last_checkpoint.file_offset, super_block.super_block_journal_file_offset
810                )));
811            }
812            let handle = ObjectStore::open_object(
813                &root_parent,
814                super_block.journal_object_id,
815                journal_handle_options(),
816                None,
817            )
818            .await
819            .with_context(|| {
820                format!(
821                    "Failed to open journal file (object id: {})",
822                    super_block.journal_object_id
823                )
824            })?;
825            let _ = self.handle.set(handle);
826            let mut inner = self.inner.lock();
827            reader.skip_to_end_of_block();
828            let mut writer_checkpoint = reader.journal_file_checkpoint();
829
830            // Make sure we don't accidentally use the reader from now onwards.
831            std::mem::drop(reader);
832
833            // Reset the stream to indicate that we've remounted the journal.
834            writer_checkpoint.checksum ^= RESET_XOR;
835            writer_checkpoint.version = LATEST_VERSION;
836            inner.flushed_offset = writer_checkpoint.file_offset;
837
838            // When we open the filesystem as writable, we flush the device.
839            inner.device_flushed_offset = inner.flushed_offset;
840
841            inner.writer.seek(writer_checkpoint);
842            inner.output_reset_version = true;
843            inner.valid_to = last_checkpoint.file_offset;
844            if last_checkpoint.file_offset < inner.flushed_offset {
845                inner.discard_offset = Some(last_checkpoint.file_offset);
846            }
847        }
848
849        self.objects
850            .on_replay_complete()
851            .await
852            .context("Failed to complete replay for object manager")?;
853
854        info!(checkpoint = last_checkpoint.file_offset, discarded_to; "replay complete");
855        Ok(())
856    }
857
858    async fn read_transactions(
859        &self,
860        reader: &mut JournalReader,
861        end_offset: Option<u64>,
862        object_id_filter: u64,
863    ) -> Result<JournaledTransactions, Error> {
864        let mut transactions = Vec::new();
865        let (mut device_flushed_offset, root_parent_store_object_id, root_store_object_id) = {
866            let super_block = &self.inner.lock().super_block_header;
867            (
868                super_block.super_block_journal_file_offset,
869                super_block.root_parent_store_object_id,
870                super_block.root_store_object_id,
871            )
872        };
873        let mut current_transaction = None;
874        let mut begin_flush_offsets = HashMap::default();
875        let mut stores_deleted = HashSet::new();
876        loop {
877            // Cache the checkpoint before we deserialize a record.
878            let checkpoint = reader.journal_file_checkpoint();
879            if let Some(end_offset) = end_offset {
880                if checkpoint.file_offset >= end_offset {
881                    break;
882                }
883            }
884            let result =
885                reader.deserialize().await.context("Failed to deserialize journal record")?;
886            match result {
887                ReadResult::Reset(_) => {
888                    if current_transaction.is_some() {
889                        current_transaction = None;
890                        transactions.pop();
891                    }
892                    let offset = reader.journal_file_checkpoint().file_offset;
893                    if offset > device_flushed_offset {
894                        device_flushed_offset = offset;
895                    }
896                }
897                ReadResult::Some(record) => {
898                    match record {
899                        JournalRecord::EndBlock => {
900                            reader.skip_to_end_of_block();
901                        }
902                        JournalRecord::Mutation { object_id, mutation } => {
903                            let current_transaction = match current_transaction.as_mut() {
904                                None => {
905                                    transactions.push(JournaledTransaction::new(checkpoint));
906                                    current_transaction = transactions.last_mut();
907                                    current_transaction.as_mut().unwrap()
908                                }
909                                Some(transaction) => transaction,
910                            };
911
912                            if stores_deleted.contains(&object_id) {
913                                bail!(
914                                    anyhow!(FxfsError::Inconsistent)
915                                        .context("Encountered mutations for deleted store")
916                                );
917                            }
918
919                            match &mutation {
920                                Mutation::BeginFlush => {
921                                    begin_flush_offsets.insert(
922                                        object_id,
923                                        current_transaction.checkpoint.file_offset,
924                                    );
925                                }
926                                Mutation::EndFlush => {
927                                    if let Some(offset) = begin_flush_offsets.remove(&object_id) {
928                                        if let Some(deleted_volume) =
929                                            &current_transaction.volume_deleted
930                                        {
931                                            if *deleted_volume == object_id {
932                                                bail!(anyhow!(FxfsError::Inconsistent).context(
933                                                    "Multiple EndFlush/DeleteVolume mutations in a \
934                                                    single transaction for the same object"
935                                                ));
936                                            }
937                                        }
938                                        // The +1 is because we don't want to replay the transaction
939                                        // containing the begin flush; we don't need or want to
940                                        // replay it.
941                                        if current_transaction
942                                            .end_flush
943                                            .replace((object_id, offset + 1))
944                                            .is_some()
945                                        {
946                                            bail!(anyhow!(FxfsError::Inconsistent).context(
947                                                "Multiple EndFlush mutations in a \
948                                                 single transaction"
949                                            ));
950                                        }
951                                    }
952                                }
953                                Mutation::DeleteVolume => {
954                                    if let Some((flushed_object, _)) =
955                                        &current_transaction.end_flush
956                                    {
957                                        if *flushed_object == object_id {
958                                            bail!(anyhow!(FxfsError::Inconsistent).context(
959                                                "Multiple EndFlush/DeleteVolume mutations in a \
960                                                    single transaction for the same object"
961                                            ));
962                                        }
963                                    }
964                                    if current_transaction
965                                        .volume_deleted
966                                        .replace(object_id)
967                                        .is_some()
968                                    {
969                                        bail!(anyhow!(FxfsError::Inconsistent).context(
970                                            "Multiple DeleteVolume mutations in a single \
971                                             transaction"
972                                        ));
973                                    }
974                                    stores_deleted.insert(object_id);
975                                }
976                                _ => {}
977                            }
978
979                            // If this mutation doesn't need to be applied, don't bother adding it
980                            // to the transaction.
981                            if (object_id_filter == INVALID_OBJECT_ID
982                                || object_id_filter == object_id)
983                                && self.should_apply(object_id, &current_transaction.checkpoint)
984                            {
985                                if object_id == root_parent_store_object_id {
986                                    current_transaction.root_parent_mutations.push(mutation);
987                                } else if object_id == root_store_object_id {
988                                    current_transaction.root_mutations.push(mutation);
989                                } else {
990                                    current_transaction
991                                        .non_root_mutations
992                                        .push((object_id, mutation));
993                                }
994                            }
995                        }
996                        JournalRecord::DataChecksums(device_range, checksums, first_write) => {
997                            let current_transaction = match current_transaction.as_mut() {
998                                None => {
999                                    transactions.push(JournaledTransaction::new(checkpoint));
1000                                    current_transaction = transactions.last_mut();
1001                                    current_transaction.as_mut().unwrap()
1002                                }
1003                                Some(transaction) => transaction,
1004                            };
1005                            current_transaction.checksums.push(JournaledChecksums {
1006                                device_range,
1007                                checksums,
1008                                first_write,
1009                            });
1010                        }
1011                        JournalRecord::Commit => {
1012                            if let Some(&mut JournaledTransaction {
1013                                ref checkpoint,
1014                                ref root_parent_mutations,
1015                                ref mut end_offset,
1016                                ..
1017                            }) = current_transaction.take()
1018                            {
1019                                for mutation in root_parent_mutations {
1020                                    // Snoop the mutations for any that might apply to the journal
1021                                    // file so that we can pass them to the reader so that it can
1022                                    // read the journal file.
1023                                    if let Mutation::ObjectStore(ObjectStoreMutation {
1024                                        item:
1025                                            Item {
1026                                                key:
1027                                                    ObjectKey {
1028                                                        object_id,
1029                                                        data:
1030                                                            ObjectKeyData::Attribute(
1031                                                                AttributeId::DATA,
1032                                                                AttributeKey::Extent(extent),
1033                                                            ),
1034                                                        ..
1035                                                    },
1036                                                value:
1037                                                    ObjectValue::Extent(ExtentValue::Some {
1038                                                        device_offset,
1039                                                        ..
1040                                                    }),
1041                                                ..
1042                                            },
1043                                        ..
1044                                    }) = mutation
1045                                    {
1046                                        // Add the journal extents we find on the way to our
1047                                        // reader.
1048                                        let handle = reader.handle();
1049                                        if *object_id != handle.object_id() {
1050                                            continue;
1051                                        }
1052                                        if let Some(end_offset) = handle.end_offset() {
1053                                            if extent.start != end_offset {
1054                                                bail!(anyhow!(FxfsError::Inconsistent).context(
1055                                                    format!(
1056                                                        "Unexpected journal extent {:?} -> {}, \
1057                                                           expected start: {}",
1058                                                        *extent, device_offset, end_offset,
1059                                                    )
1060                                                ));
1061                                            }
1062                                        }
1063                                        handle.push_extent(
1064                                            checkpoint.file_offset,
1065                                            *device_offset
1066                                                ..*device_offset
1067                                                    + extent.length().context("Invalid extent")?,
1068                                        );
1069                                    }
1070                                }
1071                                *end_offset = reader.journal_file_checkpoint().file_offset;
1072                            }
1073                        }
1074                        JournalRecord::Discard(offset) => {
1075                            if offset == 0 {
1076                                bail!(
1077                                    anyhow!(FxfsError::Inconsistent)
1078                                        .context("Invalid offset for Discard")
1079                                );
1080                            }
1081                            if let Some(transaction) = current_transaction.as_ref() {
1082                                if transaction.checkpoint.file_offset < offset {
1083                                    // Odd, but OK.
1084                                    continue;
1085                                }
1086                            }
1087                            current_transaction = None;
1088                            while let Some(transaction) = transactions.last() {
1089                                if transaction.checkpoint.file_offset < offset {
1090                                    break;
1091                                }
1092                                transactions.pop();
1093                            }
1094                            reader.handle().discard_extents(offset);
1095                        }
1096                        JournalRecord::DidFlushDevice(offset) => {
1097                            if offset > device_flushed_offset {
1098                                device_flushed_offset = offset;
1099                            }
1100                        }
1101                    }
1102                }
1103                // This is expected when we reach the end of the journal stream.
1104                ReadResult::ChecksumMismatch => break,
1105            }
1106        }
1107
1108        // Discard any uncommitted transaction.
1109        if current_transaction.is_some() {
1110            transactions.pop();
1111        }
1112
1113        Ok(JournaledTransactions { transactions, device_flushed_offset })
1114    }
1115
1116    /// Creates an empty filesystem with the minimum viable objects (including a root parent and
1117    /// root store but no further child stores).
1118    pub async fn init_empty(&self, filesystem: Arc<FxFilesystem>) -> Result<(), Error> {
1119        // The following constants are only used at format time. When mounting, the recorded values
1120        // in the superblock should be used.  The root parent store does not have a parent, but
1121        // needs an object ID to be registered with ObjectManager, so it cannot collide (i.e. have
1122        // the same object ID) with any objects in the root store that use the journal to track
1123        // mutations.
1124        const INIT_ROOT_PARENT_STORE_OBJECT_ID: u64 = 3;
1125        const INIT_ROOT_STORE_OBJECT_ID: u64 = 4;
1126        const INIT_ALLOCATOR_OBJECT_ID: u64 = 5;
1127
1128        info!(device_size = filesystem.device().size(); "Formatting");
1129
1130        let checkpoint = JournalCheckpoint {
1131            version: LATEST_VERSION,
1132            ..self.inner.lock().writer.journal_file_checkpoint()
1133        };
1134
1135        let mut current_generation = 1;
1136        if filesystem.options().image_builder_mode.is_some() {
1137            // Note that in non-image_builder_mode we write both superblocks when we format
1138            // (in FxFilesystemBuilder::open). In image_builder_mode we only write once at the end
1139            // as part of finalize(), which is why we must make sure the generation we write is
1140            // newer than any existing generation.
1141
1142            // Note: This should is the *filesystem* block size, not the device block size which
1143            // is currently always 4096 (https://fxbug.dev/42063349)
1144            let block_size = filesystem.block_size();
1145            match self.read_superblocks(filesystem.device(), block_size).await {
1146                Ok((super_block, _)) => {
1147                    log::info!(
1148                        "Found existing superblock with generation {}. Bumping by 1.",
1149                        super_block.generation
1150                    );
1151                    current_generation = super_block.generation.wrapping_add(1);
1152                }
1153                Err(_) => {
1154                    // TODO(https://fxbug.dev/463757813): It's not unusual to fail to read
1155                    // superblocks when we're formatting a new filesystem but we should probably
1156                    // fail the format if we get an IO error.
1157                }
1158            }
1159        }
1160
1161        let root_parent = ObjectStore::new_empty(
1162            None,
1163            INIT_ROOT_PARENT_STORE_OBJECT_ID,
1164            filesystem.clone(),
1165            None,
1166        );
1167        self.objects.set_root_parent_store(root_parent.clone());
1168
1169        let allocator = Arc::new(Allocator::new(filesystem.clone(), INIT_ALLOCATOR_OBJECT_ID));
1170        self.objects.set_allocator(allocator.clone());
1171        self.objects.init_metadata_reservation()?;
1172
1173        let journal_handle;
1174        let super_block_a_handle;
1175        let super_block_b_handle;
1176        let root_store;
1177        let mut transaction = root_parent
1178            .new_transaction(
1179                lock_keys![],
1180                Options { skip_journal_checks: true, ..Default::default() },
1181            )
1182            .await?;
1183        root_store = root_parent
1184            .new_child_store(
1185                &mut transaction,
1186                NewChildStoreOptions { object_id: INIT_ROOT_STORE_OBJECT_ID, ..Default::default() },
1187                None,
1188            )
1189            .await
1190            .context("new_child_store")?;
1191        self.objects.set_root_store(root_store.clone());
1192
1193        allocator.create(&mut transaction).await?;
1194
1195        // Create the super-block objects...
1196        super_block_a_handle = ObjectStore::create_object_with_id(
1197            &root_store,
1198            &mut transaction,
1199            ReservedId::new(&root_store, NonZero::new(SuperBlockInstance::A.object_id()).unwrap()),
1200            HandleOptions::default(),
1201            None,
1202        )
1203        .context("create super block")?;
1204        root_store.update_last_object_id(SuperBlockInstance::A.object_id());
1205        super_block_a_handle
1206            .extend(&mut transaction, SuperBlockInstance::A.first_extent())
1207            .await
1208            .context("extend super block")?;
1209        super_block_b_handle = ObjectStore::create_object_with_id(
1210            &root_store,
1211            &mut transaction,
1212            ReservedId::new(&root_store, NonZero::new(SuperBlockInstance::B.object_id()).unwrap()),
1213            HandleOptions::default(),
1214            None,
1215        )
1216        .context("create super block")?;
1217        root_store.update_last_object_id(SuperBlockInstance::B.object_id());
1218        super_block_b_handle
1219            .extend(&mut transaction, SuperBlockInstance::B.first_extent())
1220            .await
1221            .context("extend super block")?;
1222
1223        // the journal object...
1224        journal_handle = ObjectStore::create_object(
1225            &root_parent,
1226            &mut transaction,
1227            journal_handle_options(),
1228            None,
1229        )
1230        .await
1231        .context("create journal")?;
1232        if self.inner.lock().image_builder_mode.is_none() {
1233            let mut file_range = 0..self.chunk_size();
1234            journal_handle
1235                .preallocate_range(&mut transaction, &mut file_range)
1236                .await
1237                .context("preallocate journal")?;
1238            if file_range.start < file_range.end {
1239                bail!("preallocate_range returned too little space");
1240            }
1241        }
1242
1243        // Write the root store object info.
1244        root_store.create(&mut transaction).await?;
1245
1246        // The root parent graveyard.
1247        root_parent.set_graveyard_directory_object_id(
1248            Graveyard::create(&mut transaction, &root_parent).await?,
1249        );
1250
1251        transaction.commit().await?;
1252
1253        self.inner.lock().super_block_header = SuperBlockHeader::new(
1254            current_generation,
1255            root_parent.store_object_id(),
1256            root_parent.graveyard_directory_object_id(),
1257            root_store.store_object_id(),
1258            allocator.object_id(),
1259            journal_handle.object_id(),
1260            checkpoint,
1261            /* earliest_version: */ LATEST_VERSION,
1262        );
1263
1264        // Initialize the journal writer.
1265        let _ = self.handle.set(journal_handle);
1266        Ok(())
1267    }
1268
1269    /// Normally we allocate the journal when creating the filesystem.
1270    /// This is used image_builder_mode when journal allocation is done last.
1271    pub async fn allocate_journal(&self) -> Result<(), Error> {
1272        let handle = self.handle.get().unwrap();
1273        let mut transaction = handle
1274            .store()
1275            .new_transaction(
1276                lock_keys![LockKey::object(handle.store().store_object_id(), handle.object_id()),],
1277                Options { skip_journal_checks: true, ..Default::default() },
1278            )
1279            .await?;
1280        let mut file_range = 0..self.chunk_size();
1281        self.handle
1282            .get()
1283            .unwrap()
1284            .preallocate_range(&mut transaction, &mut file_range)
1285            .await
1286            .context("preallocate journal")?;
1287        if file_range.start < file_range.end {
1288            bail!("preallocate_range returned too little space");
1289        }
1290        transaction.commit().await?;
1291        Ok(())
1292    }
1293
1294    pub async fn init_superblocks(&self) -> Result<(), Error> {
1295        // Overwrite both superblocks.
1296        for _ in 0..2 {
1297            self.write_super_block().await?;
1298        }
1299        Ok(())
1300    }
1301
1302    /// Takes a snapshot of all journaled transactions which affect |object_id| since its last
1303    /// flush.
1304    /// The caller is responsible for locking; it must ensure that the journal is not trimmed during
1305    /// this call.  For example, a Flush lock could be held on the object in question (assuming that
1306    /// object has data to flush and is registered with ObjectManager).
1307    pub async fn read_transactions_for_object(
1308        &self,
1309        object_id: u64,
1310    ) -> Result<Vec<JournaledTransaction>, Error> {
1311        let handle = self.handle.get().expect("No journal handle");
1312        // Reopen the handle since JournalReader needs an owned handle.
1313        let handle = ObjectStore::open_object(
1314            handle.owner(),
1315            handle.object_id(),
1316            journal_handle_options(),
1317            None,
1318        )
1319        .await?;
1320
1321        let checkpoint = match self.objects.journal_checkpoint(object_id) {
1322            Some(checkpoint) => checkpoint,
1323            None => return Ok(vec![]),
1324        };
1325        let mut reader = JournalReader::new(handle, &checkpoint);
1326        // Record the current end offset and only read to there, so we don't accidentally read any
1327        // partially flushed blocks.
1328        let end_offset = self.inner.lock().valid_to;
1329        Ok(self.read_transactions(&mut reader, Some(end_offset), object_id).await?.transactions)
1330    }
1331
1332    /// Commits a transaction.  This is not thread safe; the caller must take appropriate locks.
1333    pub async fn commit(&self, transaction: &mut Transaction<'_>) -> Result<u64, Error> {
1334        if transaction.is_empty() {
1335            return Ok(self.inner.lock().writer.journal_file_checkpoint().file_offset);
1336        }
1337
1338        self.pre_commit(transaction).await?;
1339        Ok(self.write_and_apply_mutations(transaction))
1340    }
1341
1342    // Before we commit, we might need to extend the journal or write pending records to the
1343    // journal.
1344    async fn pre_commit(&self, _transaction: &Transaction<'_>) -> Result<(), Error> {
1345        let handle;
1346
1347        let (size, zero_offset) = {
1348            let mut inner = self.inner.lock();
1349
1350            // If this is the first write after a RESET, we need to output version first.
1351            if std::mem::take(&mut inner.output_reset_version) {
1352                LATEST_VERSION.serialize_into(&mut inner.writer)?;
1353            }
1354
1355            if let Some(discard_offset) = inner.discard_offset {
1356                JournalRecord::Discard(discard_offset).serialize_into(&mut inner.writer)?;
1357                inner.discard_offset = None;
1358            }
1359
1360            if inner.needs_did_flush_device {
1361                let offset = inner.device_flushed_offset;
1362                JournalRecord::DidFlushDevice(offset).serialize_into(&mut inner.writer)?;
1363                inner.needs_did_flush_device = false;
1364            }
1365
1366            handle = match self.handle.get() {
1367                None => return Ok(()),
1368                Some(x) => x,
1369            };
1370
1371            let file_offset = inner.writer.journal_file_checkpoint().file_offset;
1372
1373            let size = handle.get_size();
1374            let size = if file_offset + self.chunk_size() > size { Some(size) } else { None };
1375
1376            if size.is_none()
1377                && inner.zero_offset.is_none()
1378                && !self.objects.needs_borrow_for_journal(file_offset)
1379            {
1380                return Ok(());
1381            }
1382
1383            (size, inner.zero_offset)
1384        };
1385
1386        let mut transaction = handle
1387            .new_transaction_with_options(Options {
1388                skip_journal_checks: true,
1389                reservation: ReservationOptions::BorrowedMetadataAndData,
1390                ..Default::default()
1391            })
1392            .await?;
1393        if let Some(size) = size {
1394            handle
1395                .preallocate_range(&mut transaction, &mut (size..size + self.chunk_size()))
1396                .await?;
1397        }
1398        if let Some(zero_offset) = zero_offset {
1399            handle.zero(&mut transaction, 0..zero_offset).await?;
1400        }
1401
1402        // We can't use regular transaction commit, because that can cause re-entrancy issues, so
1403        // instead we just apply the transaction directly here.
1404        self.write_and_apply_mutations(&mut transaction);
1405
1406        let mut inner = self.inner.lock();
1407
1408        // Make sure the transaction to extend the journal made it to the journal within the old
1409        // size, since otherwise, it won't be possible to replay.
1410        if let Some(size) = size {
1411            assert!(inner.writer.journal_file_checkpoint().file_offset < size);
1412        }
1413
1414        if inner.zero_offset == zero_offset {
1415            inner.zero_offset = None;
1416        }
1417
1418        Ok(())
1419    }
1420
1421    // Determines whether a mutation at the given checkpoint should be applied.  During replay, not
1422    // all records should be applied because the object store or allocator might already contain the
1423    // mutation.  After replay, that obviously isn't the case and we want to apply all mutations.
1424    fn should_apply(&self, object_id: u64, journal_file_checkpoint: &JournalCheckpoint) -> bool {
1425        let super_block_header = &self.inner.lock().super_block_header;
1426        let offset = super_block_header
1427            .journal_file_offsets
1428            .get(&object_id)
1429            .cloned()
1430            .unwrap_or(super_block_header.super_block_journal_file_offset);
1431        journal_file_checkpoint.file_offset >= offset
1432    }
1433
1434    /// Flushes previous writes to the device and then writes out a new super-block.
1435    /// Callers must ensure that we do not make concurrent calls.
1436    async fn write_super_block(&self) -> Result<(), Error> {
1437        let root_parent_store = self.objects.root_parent_store();
1438
1439        // We need to flush previous writes to the device since the new super-block we are writing
1440        // relies on written data being observable, and we also need to lock the root parent store
1441        // so that no new entries are written to it whilst we are writing the super-block, and for
1442        // that we use the write lock.
1443        let old_layers;
1444        let old_super_block_offset;
1445        let mut new_super_block_header;
1446        let checkpoint;
1447        let borrowed;
1448
1449        {
1450            let _sync_guard = debug_assert_not_too_long!(self.sync_mutex.lock());
1451            {
1452                let _write_guard = self.writer_mutex.lock();
1453                (checkpoint, borrowed) = self.pad_to_block()?;
1454                old_layers = super_block::compact_root_parent(&*root_parent_store)?;
1455            }
1456            self.flush_device(checkpoint.file_offset)
1457                .await
1458                .context("flush failed when writing superblock")?;
1459        }
1460
1461        new_super_block_header = self.inner.lock().super_block_header.clone();
1462
1463        old_super_block_offset = new_super_block_header.journal_checkpoint.file_offset;
1464
1465        let (journal_file_offsets, min_checkpoint) = self.objects.journal_file_offsets();
1466
1467        new_super_block_header.generation = new_super_block_header.generation.wrapping_add(1);
1468        new_super_block_header.super_block_journal_file_offset = checkpoint.file_offset;
1469        new_super_block_header.journal_checkpoint = min_checkpoint.unwrap_or(checkpoint);
1470        new_super_block_header.journal_checkpoint.version = LATEST_VERSION;
1471        new_super_block_header.journal_file_offsets = journal_file_offsets;
1472        new_super_block_header.borrowed_metadata_space = borrowed;
1473
1474        self.super_block_manager
1475            .save(
1476                new_super_block_header.clone(),
1477                self.objects.root_parent_store().filesystem(),
1478                old_layers,
1479            )
1480            .await?;
1481        {
1482            let mut inner = self.inner.lock();
1483            inner.super_block_header = new_super_block_header;
1484            inner.zero_offset = Some(BLOCK_SIZE.align_down(old_super_block_offset));
1485        }
1486
1487        Ok(())
1488    }
1489
1490    /// Flushes any buffered journal data to the device.  Note that this does not flush the device
1491    /// unless the flush_device option is set, in which case data should have been persisted to
1492    /// lower layers.  If a precondition is supplied, it is evaluated and the sync will be skipped
1493    /// if it returns false.  This allows callers to check a condition whilst a lock is held.  If a
1494    /// sync is performed, this function returns the checkpoint that was flushed and the amount of
1495    /// borrowed metadata space at the point it was flushed.
1496    pub async fn sync(
1497        &self,
1498        options: SyncOptions<'_>,
1499    ) -> Result<Option<(JournalCheckpoint, u64)>, Error> {
1500        let _guard = debug_assert_not_too_long!(self.sync_mutex.lock());
1501
1502        let (checkpoint, borrowed) = {
1503            if let Some(precondition) = options.precondition {
1504                if !precondition() {
1505                    return Ok(None);
1506                }
1507            }
1508
1509            // This guard is required so that we don't insert an EndBlock record in the middle of a
1510            // transaction.
1511            let _guard = self.writer_mutex.lock();
1512
1513            self.pad_to_block()?
1514        };
1515
1516        if options.flush_device {
1517            self.flush_device(checkpoint.file_offset).await.context("sync: flush failed")?;
1518        }
1519
1520        Ok(Some((checkpoint, borrowed)))
1521    }
1522
1523    // Returns the checkpoint as it was prior to padding.  This is done because the super block
1524    // needs to record where the last transaction ends and it's the next transaction that pays the
1525    // price of the padding.
1526    fn pad_to_block(&self) -> Result<(JournalCheckpoint, u64), Error> {
1527        let mut inner = self.inner.lock();
1528        let checkpoint = inner.writer.journal_file_checkpoint();
1529        if checkpoint.file_offset % BLOCK_SIZE != 0 {
1530            JournalRecord::EndBlock.serialize_into(&mut inner.writer)?;
1531            inner.writer.pad_to_block()?;
1532            if let Some(waker) = inner.flush_waker.take() {
1533                waker.wake();
1534            }
1535        }
1536        Ok((checkpoint, self.objects.borrowed_metadata_space()))
1537    }
1538
1539    async fn flush_device(&self, checkpoint_offset: u64) -> Result<(), Error> {
1540        assert!(
1541            self.inner.lock().image_builder_mode.is_none(),
1542            "flush_device called in image builder mode"
1543        );
1544        debug_assert_not_too_long!(poll_fn(|ctx| {
1545            let mut inner = self.inner.lock();
1546            if inner.flushed_offset >= checkpoint_offset {
1547                Poll::Ready(Ok(()))
1548            } else if inner.terminate {
1549                let context = inner
1550                    .terminate_reason
1551                    .as_ref()
1552                    .map(|e| format!("Journal closed with error: {:?}", e))
1553                    .unwrap_or_else(|| "Journal closed".to_string());
1554                Poll::Ready(Err(anyhow!(FxfsError::JournalFlushError).context(context)))
1555            } else {
1556                inner.sync_waker = Some(ctx.waker().clone());
1557                Poll::Pending
1558            }
1559        }))?;
1560
1561        let needs_flush = self.inner.lock().device_flushed_offset < checkpoint_offset;
1562        if needs_flush {
1563            let trace = self.trace.load(Ordering::Relaxed);
1564            if trace {
1565                info!("J: start flush device");
1566            }
1567            self.handle.get().unwrap().flush_device().await?;
1568            if trace {
1569                info!("J: end flush device");
1570            }
1571
1572            // We need to write a DidFlushDevice record at some point, but if we are in the
1573            // process of shutting down the filesystem, we want to leave the journal clean to
1574            // avoid there being log messages complaining about unwritten journal data, so we
1575            // queue it up so that the next transaction will trigger this record to be written.
1576            // If we are shutting down, that will never happen but since the DidFlushDevice
1577            // message is purely advisory (it reduces the number of checksums we have to verify
1578            // during replay), it doesn't matter if it isn't written.
1579            {
1580                let mut inner = self.inner.lock();
1581                inner.device_flushed_offset = checkpoint_offset;
1582                inner.needs_did_flush_device = true;
1583            }
1584
1585            // Tell the allocator that we flushed the device so that it can now start using
1586            // space that was deallocated.
1587            self.objects.allocator().did_flush_device(checkpoint_offset);
1588            if trace {
1589                info!("J: did flush device");
1590            }
1591        }
1592
1593        Ok(())
1594    }
1595
1596    /// Returns a copy of the super-block header.
1597    pub fn super_block_header(&self) -> SuperBlockHeader {
1598        self.inner.lock().super_block_header.clone()
1599    }
1600
1601    /// Waits for there to be sufficient space in the journal.
1602    pub async fn check_journal_space(&self) -> Result<(), Error> {
1603        loop {
1604            let listener = {
1605                let inner = self.inner.lock();
1606                if inner.terminate {
1607                    // If the flush error is set, this will never make progress, since we can't
1608                    // extend the journal any more.
1609                    let context = inner
1610                        .terminate_reason
1611                        .as_ref()
1612                        .map(|e| format!("Journal closed with error: {:?}", e))
1613                        .unwrap_or_else(|| "Journal closed".to_string());
1614                    return Err(anyhow!(FxfsError::JournalFlushError).context(context));
1615                }
1616                if self.objects.last_end_offset()
1617                    - inner.super_block_header.journal_checkpoint.file_offset
1618                    < inner.reclaim_size
1619                {
1620                    return Ok(());
1621                }
1622                if inner.image_builder_mode.is_some() {
1623                    return Ok(());
1624                }
1625                if inner.disable_compactions {
1626                    return Err(
1627                        anyhow!(FxfsError::JournalFlushError).context("Compactions disabled")
1628                    );
1629                }
1630                self.reclaim_event.listen()
1631            };
1632            self.hooks.on_waiting_for_journal_space();
1633            debug_assert_not_too_long!(listener);
1634        }
1635    }
1636
1637    fn chunk_size(&self) -> u64 {
1638        CHUNK_SIZE
1639    }
1640
1641    fn write_and_apply_mutations(&self, transaction: &mut Transaction<'_>) -> u64 {
1642        let checkpoint_before;
1643        let checkpoint_after;
1644        {
1645            let _guard = self.writer_mutex.lock();
1646            checkpoint_before = {
1647                let mut inner = self.inner.lock();
1648                if transaction.includes_write() {
1649                    inner.needs_barrier = true;
1650                }
1651                let checkpoint = inner.writer.journal_file_checkpoint();
1652                let mut iter = transaction.mutations().iter().peekable();
1653                while let Some(store_iter) = ObjectMutationIterator::new(&mut iter) {
1654                    let object_id = store_iter.object_id();
1655                    self.objects.write_mutations(
1656                        object_id,
1657                        store_iter,
1658                        Writer(object_id, &mut inner.writer),
1659                    );
1660                }
1661                checkpoint
1662            };
1663            let maybe_mutation =
1664                self.objects.apply_transaction(transaction, &checkpoint_before).expect(
1665                    "apply_transaction should not fail in live mode; \
1666                     filesystem will be in an inconsistent state",
1667                );
1668            checkpoint_after = {
1669                let mut inner = self.inner.lock();
1670                if let Some(mutation) = maybe_mutation {
1671                    inner
1672                        .writer
1673                        .write_record(&JournalRecord::Mutation { object_id: 0, mutation })
1674                        .unwrap();
1675                }
1676                for (device_range, checksums, first_write) in
1677                    transaction.take_checksums().into_iter()
1678                {
1679                    inner
1680                        .writer
1681                        .write_record(&JournalRecord::DataChecksums(
1682                            device_range,
1683                            Checksums::fletcher(checksums),
1684                            first_write,
1685                        ))
1686                        .unwrap();
1687                }
1688                inner.writer.write_record(&JournalRecord::Commit).unwrap();
1689
1690                inner.writer.journal_file_checkpoint()
1691            };
1692        }
1693        self.objects.did_commit_transaction(
1694            transaction,
1695            &checkpoint_before,
1696            checkpoint_after.file_offset,
1697        );
1698
1699        if let Some(waker) = self.inner.lock().flush_waker.take() {
1700            waker.wake();
1701        }
1702
1703        checkpoint_before.file_offset
1704    }
1705
1706    /// This task will flush journal data to the device when there is data that needs flushing, and
1707    /// trigger compactions when short of journal space.  It will return after the terminate method
1708    /// has been called, or an error is encountered with either flushing or compaction.
1709    pub async fn flush_task(self: Arc<Self>) {
1710        let mut flush_fut = None;
1711        let mut compact_fut = None;
1712        let mut flush_error = false;
1713        poll_fn(|ctx| {
1714            loop {
1715                {
1716                    let mut inner = self.inner.lock();
1717                    if flush_fut.is_none() && !flush_error && self.handle.get().is_some() {
1718                        let flushable = inner.writer.flushable_bytes();
1719                        if flushable > 0 {
1720                            flush_fut = Some(Box::pin(self.flush(flushable)));
1721                        }
1722                    }
1723                    if inner.terminate && flush_fut.is_none() && compact_fut.is_none() {
1724                        return Poll::Ready(());
1725                    }
1726                    // `journal_bytes` refers to bytes in the journal that haven't yet been flushed
1727                    // to the layer files. It increases with each transaction and decreases when
1728                    // compactions complete. The flush task is woken whenever a transaction is
1729                    // committed and we should see this metric updated regularly.
1730                    let journal_bytes = self.objects.last_end_offset()
1731                        - inner.super_block_header.journal_checkpoint.file_offset;
1732                    fxfs_trace::counter!("journal-bytes", 0, "total" => journal_bytes);
1733                    // The / 2 is here because after compacting, we cannot reclaim the space until
1734                    // the _next_ time we flush the device since the super-block is not guaranteed
1735                    // to persist until then.
1736                    if compact_fut.is_none()
1737                        && !inner.terminate
1738                        && !inner.disable_compactions
1739                        && !inner.compactions_paused
1740                        && inner.image_builder_mode.is_none()
1741                        && journal_bytes > inner.reclaim_size / 2
1742                    {
1743                        compact_fut = Some(Box::pin(self.compact()));
1744                        inner.compaction_running = true;
1745                    }
1746                    inner.flush_waker = Some(ctx.waker().clone());
1747                }
1748                let mut pending = true;
1749                if let Some(fut) = flush_fut.as_mut() {
1750                    if let Poll::Ready(result) = fut.poll_unpin(ctx) {
1751                        if let Err(e) = result {
1752                            self.inner.lock().terminate(Some(e.context("Flush error")));
1753                            self.reclaim_event.notify(usize::MAX);
1754                            flush_error = true;
1755                        }
1756                        flush_fut = None;
1757                        pending = false;
1758                    }
1759                }
1760                if let Some(fut) = compact_fut.as_mut() {
1761                    if let Poll::Ready(result) = fut.poll_unpin(ctx) {
1762                        let mut inner = self.inner.lock();
1763                        if let Err(e) = result {
1764                            inner.terminate(Some(e.context("Compaction error")));
1765                        }
1766                        compact_fut = None;
1767                        inner.compaction_running = false;
1768                        self.reclaim_event.notify(usize::MAX);
1769                        pending = false;
1770                        fxfs_trace::counter!(
1771                            "journal-bytes",
1772                            0,
1773                            "total" => self.objects.last_end_offset()
1774                                - inner.super_block_header.journal_checkpoint.file_offset
1775                        );
1776                    }
1777                }
1778                if pending {
1779                    return Poll::Pending;
1780                }
1781            }
1782        })
1783        .await;
1784    }
1785
1786    /// Returns a yielder that can be used for compactions.
1787    pub fn get_compaction_yielder(&self) -> CompactionYielder<'_> {
1788        CompactionYielder::new(self)
1789    }
1790
1791    async fn flush(&self, amount: usize) -> Result<(), Error> {
1792        let handle = self.handle.get().unwrap();
1793        let mut buf = handle.allocate_buffer(amount).await;
1794        let (offset, len, barrier_on_first_write) = {
1795            let mut inner = self.inner.lock();
1796            let offset = inner.writer.take_flushable(buf.as_mut());
1797            let barrier_on_first_write = inner.needs_barrier && inner.barriers_enabled;
1798            // Reset `needs_barrier` before instead of after the overwrite in case a txn commit
1799            // that contains data happens during the overwrite.
1800            inner.needs_barrier = false;
1801            (offset, buf.len() as u64, barrier_on_first_write)
1802        };
1803        self.handle
1804            .get()
1805            .unwrap()
1806            .overwrite(
1807                offset,
1808                buf.as_mut(),
1809                OverwriteOptions { barrier_on_first_write, ..Default::default() },
1810            )
1811            .await?;
1812
1813        let mut inner = self.inner.lock();
1814        if let Some(waker) = inner.sync_waker.take() {
1815            waker.wake();
1816        }
1817        inner.flushed_offset = offset + len;
1818        inner.valid_to = inner.flushed_offset;
1819        Ok(())
1820    }
1821
1822    fn should_force_major_compaction(&self) -> ForceMajor {
1823        // required_reservation is the total amount of space we have set aside for metadata.  That
1824        // space is divided into borrowed space (which is used by transactions like file deletions
1825        // which will eventually, after compaction, give back their space), and
1826        // metadata_reservation (which is available for use during compaction, both to persist the
1827        // new layer files and for other transient costs).
1828        let required_reservation = self.objects.required_reservation();
1829        // reclaim_size determines the frequency of compaction.  The journal will automatically
1830        // compact when reclaim_size / 2 <= J <= reclaim_size, and will block most operations once
1831        // J > reclaim_size, so any given compaction should have no more than reclaim_size bytes.
1832        let reclaim_size = self.inner.lock().reclaim_size;
1833        // metadata_reservation = required_reservation - borrowed
1834        let metadata_reservation = self.objects.metadata_reservation().amount();
1835        // max_store_reservation is the greatest amount reserved by any object store.  This
1836        // corresponds to the minimum amount which is needed to perform a major compaction of that
1837        // store (and therefore the minimum amount needed to do a major compaction overall, since
1838        // we compact stores one by one and immediately purge old layer files after compaction,
1839        // releasing their space back into metaadta_resrvation).
1840        let max_store_reservation = self.objects.max_store_reservation();
1841
1842        // Compaction will usually return borrowed space by freeing up space in the journal, but a
1843        // major or merge compaction is sometimes necessary -- consider if you have one layer file
1844        // which adds many objects, and another layer file which deletes all of these objects.  If
1845        // the layers are merged, they cancel out and become very small.  If they are not merged,
1846        // they do not cancel out and both layers remain large.  The space borrowed to persist the
1847        // deletion mutations is only given back once they cancel out with the creation mutations.
1848        //
1849        // Thus, borrowed space can only be guaranteed to be returned if we eventually
1850        // major-compact LSM trees.
1851        //
1852        // Since major compaction is expensive, we don't want to do it too often, but we also need
1853        // to make sure we don't wait too long and drain the metadata reservation too much (in
1854        // particular, it cannot go below `max_store_reservation`, because at that point major
1855        // compaction might not be possible any more).
1856        //
1857        // Thus, we force a major compaction when metadata_reservation <= max_store_reservation -
1858        // M, for some margin M.  The reason that `reclaim_size` is used as the margin is because
1859        // `reclaim_size` determines how often we compact, and in theory, the journal is going to
1860        // be compacted at some point when J < reclaim_size.  This establishes the ideal margin --
1861        // at this point we are confident that major compaction is still possible, but the next
1862        // compaction might be another reclaim_size bytes later, which is too late.
1863        if required_reservation >= reclaim_size
1864            && metadata_reservation <= max_store_reservation.saturating_sub(reclaim_size)
1865        {
1866            ForceMajor::True
1867        } else {
1868            ForceMajor::False
1869        }
1870    }
1871
1872    #[trace]
1873    async fn compact(&self) -> Result<(), Error> {
1874        assert!(
1875            self.inner.lock().image_builder_mode.is_none(),
1876            "compact called in image builder mode"
1877        );
1878        let bytes_before = self.objects.compaction_bytes_written();
1879        let _measure = crate::metrics::DurationMeasureScope::new(
1880            &crate::metrics::lsm_tree_metrics().journal_compaction_time,
1881        );
1882        crate::metrics::lsm_tree_metrics().journal_compactions_total.add(1);
1883        let trace = self.trace.load(Ordering::Relaxed);
1884        let force_major = self.should_force_major_compaction();
1885        debug!("Compaction starting, force_major={force_major:?}");
1886        if trace {
1887            info!("J: start compaction, force_major={force_major:?}");
1888        }
1889        let earliest_version = self
1890            .objects
1891            .flush(FlushReason::Journal(force_major))
1892            .await
1893            .context("Failed to flush objects")?;
1894        self.inner.lock().super_block_header.earliest_version = earliest_version;
1895        self.write_super_block().await.context("Failed to write superblock")?;
1896        if trace {
1897            info!("J: end compaction");
1898        }
1899        debug!("Compaction finished");
1900        let bytes_after = self.objects.compaction_bytes_written();
1901        crate::metrics::lsm_tree_metrics()
1902            .journal_compaction_bytes_written
1903            .add(bytes_after.saturating_sub(bytes_before));
1904        Ok(())
1905    }
1906
1907    /// This should generally NOT be called externally. It is public to allow use by FIDL service
1908    /// fxfs.Debug.
1909    pub async fn force_compact(&self) -> Result<(), Error> {
1910        self.inner.lock().forced_compaction = true;
1911        scopeguard::defer! { self.inner.lock().forced_compaction = false; }
1912        self.compact().await
1913    }
1914
1915    pub async fn stop_compactions(&self) {
1916        loop {
1917            debug_assert_not_too_long!({
1918                let mut inner = self.inner.lock();
1919                inner.disable_compactions = true;
1920                if !inner.compaction_running {
1921                    return;
1922                }
1923                self.reclaim_event.listen()
1924            });
1925        }
1926    }
1927
1928    pub async fn pause_compactions(&self) {
1929        loop {
1930            debug_assert_not_too_long!({
1931                let mut inner = self.inner.lock();
1932                inner.compactions_paused = true;
1933                if !inner.compaction_running {
1934                    return;
1935                }
1936                self.reclaim_event.listen()
1937            });
1938        }
1939    }
1940
1941    pub fn resume_compactions(&self) {
1942        let mut inner = self.inner.lock();
1943        inner.compactions_paused = false;
1944        if let Some(waker) = inner.flush_waker.take() {
1945            waker.wake();
1946        }
1947    }
1948
1949    /// Creates a lazy inspect node named `str` under `parent` which will yield statistics for the
1950    /// journal when queried.
1951    pub fn track_statistics(self: &Arc<Self>, parent: &fuchsia_inspect::Node, name: &str) {
1952        let this = Arc::downgrade(self);
1953        parent.record_lazy_child(name, move || {
1954            let this_clone = this.clone();
1955            async move {
1956                let inspector = fuchsia_inspect::Inspector::default();
1957                if let Some(this) = this_clone.upgrade() {
1958                    let (journal_min, journal_max, journal_reclaim_size) = {
1959                        // TODO(https://fxbug.dev/42069513): Push-back or rate-limit to prevent DoS.
1960                        let inner = this.inner.lock();
1961                        (
1962                            BLOCK_SIZE.align_down(
1963                                inner.super_block_header.journal_checkpoint.file_offset,
1964                            ),
1965                            inner.flushed_offset,
1966                            inner.reclaim_size,
1967                        )
1968                    };
1969                    let root = inspector.root();
1970                    root.record_uint("journal_min_offset", journal_min);
1971                    root.record_uint("journal_max_offset", journal_max);
1972                    root.record_uint("journal_size", journal_max - journal_min);
1973                    root.record_uint("journal_reclaim_size", journal_reclaim_size);
1974
1975                    // TODO(https://fxbug.dev/42068224): Post-compute rather than manually computing metrics.
1976                    if let Some(x) = round_div(
1977                        100 * (journal_max - journal_min),
1978                        this.objects.allocator().get_disk_bytes(),
1979                    ) {
1980                        root.record_uint("journal_size_to_disk_size_percent", x);
1981                    }
1982                }
1983                Ok(inspector)
1984            }
1985            .boxed()
1986        });
1987    }
1988
1989    /// Terminate all journal activity.
1990    pub fn terminate(&self) {
1991        self.inner.lock().terminate(/*reason*/ None);
1992        self.reclaim_event.notify(usize::MAX);
1993    }
1994}
1995
1996/// Wrapper to allow records to be written to the journal.
1997pub struct Writer<'a>(u64, &'a mut JournalWriter);
1998
1999impl Writer<'_> {
2000    pub fn write(&mut self, mutation: Mutation) {
2001        self.1.write_record(&JournalRecord::Mutation { object_id: self.0, mutation }).unwrap();
2002    }
2003}
2004
2005#[cfg(target_os = "fuchsia")]
2006mod yielder {
2007    use super::Journal;
2008    use crate::lsm_tree::Yielder;
2009    use fuchsia_async as fasync;
2010
2011    /// CompactionYielder uses fuchsia-async to yield if other tasks are being polled, which should
2012    /// be a proxy for how busy the system is.  We can afford to delay compactions for a small
2013    /// amount of time but not so long that we end up blocking new transactions.
2014    pub struct CompactionYielder<'a> {
2015        journal: &'a Journal,
2016        low_priority_task: Option<fasync::LowPriorityTask>,
2017    }
2018
2019    impl<'a> CompactionYielder<'a> {
2020        pub fn new(journal: &'a Journal) -> Self {
2021            Self { journal, low_priority_task: None }
2022        }
2023    }
2024
2025    impl Yielder for CompactionYielder<'_> {
2026        async fn yield_now(&mut self) {
2027            // We will wait for the executor to be idle for 4ms, but no longer than 16ms.  We need
2028            // to cap the maximum amount of time we wait in case we've reached a point where
2029            // compaction is now urgent or else we could block new transactions.
2030            const IDLE_PERIOD: zx::MonotonicDuration = zx::MonotonicDuration::from_millis(4);
2031            const MAX_YIELD_DURATION: zx::MonotonicDuration =
2032                zx::MonotonicDuration::from_millis(16);
2033
2034            {
2035                let inner = self.journal.inner.lock();
2036                if inner.forced_compaction {
2037                    return;
2038                }
2039                let outstanding = self.journal.objects.last_end_offset()
2040                    - inner.super_block_header.journal_checkpoint.file_offset;
2041                let half_reclaim_size = inner.reclaim_size / 2;
2042                if outstanding
2043                    .checked_sub(half_reclaim_size)
2044                    .is_some_and(|x| x >= half_reclaim_size / 2)
2045                {
2046                    // If we have got to the point where we have used up 3/4 of reclaim size in the
2047                    // journal, do not delay any further.  If we continue to yield we will get to
2048                    // the point where we block new transactions.
2049                    self.low_priority_task = None;
2050                    return;
2051                }
2052            }
2053
2054            self.low_priority_task
2055                .get_or_insert_with(|| fasync::LowPriorityTask::new())
2056                .wait_until_idle_for(
2057                    IDLE_PERIOD,
2058                    fasync::MonotonicInstant::after(MAX_YIELD_DURATION),
2059                )
2060                .await;
2061        }
2062    }
2063}
2064
2065#[cfg(not(target_os = "fuchsia"))]
2066mod yielder {
2067    use super::Journal;
2068    use crate::lsm_tree::Yielder;
2069
2070    #[expect(dead_code)]
2071    pub struct CompactionYielder<'a>(&'a Journal);
2072
2073    impl<'a> CompactionYielder<'a> {
2074        pub fn new(journal: &'a Journal) -> Self {
2075            Self(journal)
2076        }
2077    }
2078
2079    impl Yielder for CompactionYielder<'_> {
2080        async fn yield_now(&mut self) {}
2081    }
2082}
2083
2084pub use yielder::*;
2085
2086#[cfg(test)]
2087mod tests {
2088    use super::SuperBlockInstance;
2089    use crate::filesystem::{FxFilesystem, FxFilesystemBuilder, SyncOptions};
2090    use crate::fsck::fsck;
2091    use crate::object_handle::{ObjectHandle, WriteObjectHandle};
2092    use crate::object_store::directory::Directory;
2093    use crate::object_store::transaction::Options;
2094    use crate::object_store::volume::root_volume;
2095    use crate::object_store::{
2096        HandleOptions, LockKey, NewChildStoreOptions, ObjectStore, StoreOptions, lock_keys,
2097    };
2098    #[cfg(target_os = "fuchsia")]
2099    use fuchsia_async::TestExecutor;
2100    use fuchsia_async::{self as fasync, MonotonicDuration};
2101    use storage_device::DeviceHolder;
2102    use storage_device::fake_device::FakeDevice;
2103
2104    const TEST_DEVICE_BLOCK_SIZE: u32 = 512;
2105
2106    #[fuchsia::test]
2107    async fn test_replay() {
2108        const TEST_DATA: &[u8] = b"hello";
2109
2110        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
2111
2112        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2113
2114        let object_id = {
2115            let root_store = fs.root_store();
2116            let root_directory =
2117                Directory::open(&root_store, root_store.root_directory_object_id())
2118                    .await
2119                    .expect("open failed");
2120            let mut transaction = fs
2121                .root_store()
2122                .new_transaction(
2123                    lock_keys![LockKey::object(
2124                        root_store.store_object_id(),
2125                        root_store.root_directory_object_id(),
2126                    )],
2127                    Options::default(),
2128                )
2129                .await
2130                .expect("new_transaction failed");
2131            let handle = root_directory
2132                .create_child_file(&mut transaction, "test")
2133                .await
2134                .expect("create_child_file failed");
2135
2136            transaction.commit().await.expect("commit failed");
2137            let mut buf = handle.allocate_buffer(TEST_DATA.len()).await;
2138            buf.copy_from_slice(TEST_DATA);
2139            handle.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
2140            // As this is the first sync, this will actually trigger a new super-block, but normally
2141            // this would not be the case.
2142            fs.sync(SyncOptions::default()).await.expect("sync failed");
2143            handle.object_id()
2144        };
2145
2146        {
2147            fs.close().await.expect("Close failed");
2148            let device = fs.take_device().await;
2149            device.reopen(false);
2150            let fs = FxFilesystem::open(device).await.expect("open failed");
2151            let handle = ObjectStore::open_object(
2152                &fs.root_store(),
2153                object_id,
2154                HandleOptions::default(),
2155                None,
2156            )
2157            .await
2158            .expect("open_object failed");
2159            let mut buf = handle.allocate_buffer(handle.block_size().get() as usize).await;
2160            assert_eq!(
2161                handle.read_aligned(0, buf.as_mut()).await.expect("read failed"),
2162                TEST_DATA.len()
2163            );
2164            assert_eq!(buf.subslice(..TEST_DATA.len()).to_vec(), TEST_DATA);
2165            fsck(fs.clone()).await.expect("fsck failed");
2166            fs.close().await.expect("Close failed");
2167        }
2168    }
2169
2170    #[fuchsia::test]
2171    async fn test_reset() {
2172        const TEST_DATA: &[u8] = b"hello";
2173
2174        let device = DeviceHolder::new(FakeDevice::new(32768, TEST_DEVICE_BLOCK_SIZE));
2175
2176        let mut object_ids = Vec::new();
2177
2178        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2179        {
2180            let root_store = fs.root_store();
2181            let root_directory =
2182                Directory::open(&root_store, root_store.root_directory_object_id())
2183                    .await
2184                    .expect("open failed");
2185            let mut transaction = fs
2186                .root_store()
2187                .new_transaction(
2188                    lock_keys![LockKey::object(
2189                        root_store.store_object_id(),
2190                        root_store.root_directory_object_id(),
2191                    )],
2192                    Options::default(),
2193                )
2194                .await
2195                .expect("new_transaction failed");
2196            let handle = root_directory
2197                .create_child_file(&mut transaction, "test")
2198                .await
2199                .expect("create_child_file failed");
2200            transaction.commit().await.expect("commit failed");
2201            let mut buf = handle.allocate_buffer(TEST_DATA.len()).await;
2202            buf.copy_from_slice(TEST_DATA);
2203            handle.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
2204            fs.sync(SyncOptions::default()).await.expect("sync failed");
2205            object_ids.push(handle.object_id());
2206
2207            // Create a lot of objects but don't sync at the end. This should leave the filesystem
2208            // with a half finished transaction that cannot be replayed.
2209            for i in 0..1000 {
2210                let mut transaction = fs
2211                    .root_store()
2212                    .new_transaction(
2213                        lock_keys![LockKey::object(
2214                            root_store.store_object_id(),
2215                            root_store.root_directory_object_id(),
2216                        )],
2217                        Options::default(),
2218                    )
2219                    .await
2220                    .expect("new_transaction failed");
2221                let handle = root_directory
2222                    .create_child_file(&mut transaction, &format!("{}", i))
2223                    .await
2224                    .expect("create_child_file failed");
2225                transaction.commit().await.expect("commit failed");
2226                let mut buf = handle.allocate_buffer(TEST_DATA.len()).await;
2227                buf.copy_from_slice(TEST_DATA);
2228                handle.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
2229                object_ids.push(handle.object_id());
2230            }
2231        }
2232        fs.close().await.expect("fs close failed");
2233        let device = fs.take_device().await;
2234        device.reopen(false);
2235        let fs = FxFilesystem::open(device).await.expect("open failed");
2236        fsck(fs.clone()).await.expect("fsck failed");
2237        {
2238            let root_store = fs.root_store();
2239            // Check the first two objects which should exist.
2240            for &object_id in &object_ids[0..1] {
2241                let handle = ObjectStore::open_object(
2242                    &root_store,
2243                    object_id,
2244                    HandleOptions::default(),
2245                    None,
2246                )
2247                .await
2248                .expect("open_object failed");
2249                let mut buf = handle.allocate_buffer(handle.block_size().get() as usize).await;
2250                assert_eq!(
2251                    handle.read_aligned(0, buf.as_mut()).await.expect("read failed"),
2252                    TEST_DATA.len()
2253                );
2254                let data = buf.subslice(..TEST_DATA.len()).to_vec();
2255                assert_eq!(data, TEST_DATA);
2256            }
2257
2258            // Write one more object and sync.
2259            let root_directory =
2260                Directory::open(&root_store, root_store.root_directory_object_id())
2261                    .await
2262                    .expect("open failed");
2263            let mut transaction = fs
2264                .root_store()
2265                .new_transaction(
2266                    lock_keys![LockKey::object(
2267                        root_store.store_object_id(),
2268                        root_store.root_directory_object_id(),
2269                    )],
2270                    Options::default(),
2271                )
2272                .await
2273                .expect("new_transaction failed");
2274            let handle = root_directory
2275                .create_child_file(&mut transaction, "test2")
2276                .await
2277                .expect("create_child_file failed");
2278            transaction.commit().await.expect("commit failed");
2279            let mut buf = handle.allocate_buffer(TEST_DATA.len()).await;
2280            buf.copy_from_slice(TEST_DATA);
2281            handle.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
2282            fs.sync(SyncOptions::default()).await.expect("sync failed");
2283            object_ids.push(handle.object_id());
2284        }
2285
2286        fs.close().await.expect("close failed");
2287        let device = fs.take_device().await;
2288        device.reopen(false);
2289        let fs = FxFilesystem::open(device).await.expect("open failed");
2290        {
2291            fsck(fs.clone()).await.expect("fsck failed");
2292
2293            // Check the first two and the last objects.
2294            for &object_id in object_ids[0..1].iter().chain(object_ids.last().cloned().iter()) {
2295                let handle = ObjectStore::open_object(
2296                    &fs.root_store(),
2297                    object_id,
2298                    HandleOptions::default(),
2299                    None,
2300                )
2301                .await
2302                .unwrap_or_else(|e| {
2303                    panic!("open_object failed (object_id: {}): {:?}", object_id, e)
2304                });
2305                let mut buf = handle.allocate_buffer(handle.block_size().get() as usize).await;
2306                assert_eq!(
2307                    handle.read_aligned(0, buf.as_mut()).await.expect("read failed"),
2308                    TEST_DATA.len()
2309                );
2310                let data = buf.subslice(..TEST_DATA.len()).to_vec();
2311                assert_eq!(data, TEST_DATA);
2312            }
2313        }
2314        fs.close().await.expect("close failed");
2315    }
2316
2317    #[fuchsia::test]
2318    async fn test_discard() {
2319        let device = {
2320            let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
2321            let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2322            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
2323
2324            let store = root_volume
2325                .new_volume("test", NewChildStoreOptions::default())
2326                .await
2327                .expect("new_volume failed");
2328            let root_directory = Directory::open(&store, store.root_directory_object_id())
2329                .await
2330                .expect("open failed");
2331
2332            // Create enough data so that another journal extent is used.
2333            let mut i = 0;
2334            loop {
2335                let mut transaction = fs
2336                    .root_store()
2337                    .new_transaction(
2338                        lock_keys![LockKey::object(
2339                            store.store_object_id(),
2340                            store.root_directory_object_id()
2341                        )],
2342                        Options::default(),
2343                    )
2344                    .await
2345                    .expect("new_transaction failed");
2346                root_directory
2347                    .create_child_file(&mut transaction, &format!("a {i}"))
2348                    .await
2349                    .expect("create_child_file failed");
2350                if transaction.commit().await.expect("commit failed") > super::CHUNK_SIZE {
2351                    break;
2352                }
2353                i += 1;
2354            }
2355
2356            // Compact and then disable compactions.
2357            fs.journal().force_compact().await.expect("compact failed");
2358            fs.journal().stop_compactions().await;
2359
2360            // Keep going until we need another journal extent.
2361            let mut i = 0;
2362            loop {
2363                let mut transaction = fs
2364                    .root_store()
2365                    .new_transaction(
2366                        lock_keys![LockKey::object(
2367                            store.store_object_id(),
2368                            store.root_directory_object_id()
2369                        )],
2370                        Options::default(),
2371                    )
2372                    .await
2373                    .expect("new_transaction failed");
2374                root_directory
2375                    .create_child_file(&mut transaction, &format!("b {i}"))
2376                    .await
2377                    .expect("create_child_file failed");
2378                if transaction.commit().await.expect("commit failed") > 2 * super::CHUNK_SIZE {
2379                    break;
2380                }
2381                i += 1;
2382            }
2383
2384            // Allow the journal to flush, but we don't want to sync.
2385            fasync::Timer::new(MonotonicDuration::from_millis(10)).await;
2386            // Because we're not gracefully closing the filesystem, a Discard record will be
2387            // emitted.
2388            fs.device().snapshot().expect("snapshot failed")
2389        };
2390
2391        let fs = FxFilesystem::open(device).await.expect("open failed");
2392
2393        {
2394            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
2395
2396            let store =
2397                root_volume.volume("test", StoreOptions::default()).await.expect("volume failed");
2398
2399            let root_directory = Directory::open(&store, store.root_directory_object_id())
2400                .await
2401                .expect("open failed");
2402
2403            // Write one more transaction.
2404            let mut transaction = fs
2405                .root_store()
2406                .new_transaction(
2407                    lock_keys![LockKey::object(
2408                        store.store_object_id(),
2409                        store.root_directory_object_id()
2410                    )],
2411                    Options::default(),
2412                )
2413                .await
2414                .expect("new_transaction failed");
2415            root_directory
2416                .create_child_file(&mut transaction, &format!("d"))
2417                .await
2418                .expect("create_child_file failed");
2419            transaction.commit().await.expect("commit failed");
2420        }
2421
2422        fs.close().await.expect("close failed");
2423        let device = fs.take_device().await;
2424        device.reopen(false);
2425
2426        let fs = FxFilesystem::open(device).await.expect("open failed");
2427        fsck(fs.clone()).await.expect("fsck failed");
2428        fs.close().await.expect("close failed");
2429    }
2430
2431    #[fuchsia::test]
2432    async fn test_use_existing_generation() {
2433        let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
2434
2435        // First format should be generation 1.
2436        let fs = FxFilesystemBuilder::new()
2437            .format(true)
2438            .image_builder_mode(Some(SuperBlockInstance::A))
2439            .open(device)
2440            .await
2441            .expect("open failed");
2442        fs.enable_allocations();
2443        let generation0 = fs.super_block_header().generation;
2444        assert_eq!(generation0, 1);
2445        fs.close().await.expect("close failed");
2446        let device = fs.take_device().await;
2447        device.reopen(false);
2448
2449        // Format the device normally (again, generation 1).
2450        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2451        let generation1 = fs.super_block_header().generation;
2452        {
2453            let root_volume = crate::object_store::volume::root_volume(fs.clone())
2454                .await
2455                .expect("root_volume failed");
2456            root_volume
2457                .new_volume("test", crate::object_store::NewChildStoreOptions::default())
2458                .await
2459                .expect("new_volume failed");
2460        }
2461        fs.close().await.expect("close failed");
2462        let device = fs.take_device().await;
2463        device.reopen(false);
2464
2465        // Format again with image_builder_mode.
2466        let fs = FxFilesystemBuilder::new()
2467            .format(true)
2468            .image_builder_mode(Some(SuperBlockInstance::A))
2469            .open(device)
2470            .await
2471            .expect("open failed");
2472        fs.enable_allocations();
2473        let generation2 = fs.super_block_header().generation;
2474        assert!(
2475            generation2 > generation1,
2476            "generation2 ({}) should be greater than generation1 ({})",
2477            generation2,
2478            generation1
2479        );
2480        fs.close().await.expect("close failed");
2481    }
2482
2483    #[fuchsia::test]
2484    async fn test_image_builder_mode_generation_bump_512_byte_block() {
2485        let device = DeviceHolder::new(FakeDevice::new(16384, 512));
2486
2487        // Format initial filesystem (generation 1)
2488        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2489        let generation1 = fs.super_block_header().generation;
2490        fs.close().await.expect("close failed");
2491        let device = fs.take_device().await;
2492        device.reopen(false);
2493
2494        // Format again with image_builder_mode (should bump generation)
2495        let fs = FxFilesystemBuilder::new()
2496            .format(true)
2497            .image_builder_mode(Some(SuperBlockInstance::A))
2498            .open(device)
2499            .await
2500            .expect("open failed");
2501
2502        fs.enable_allocations();
2503        let generation2 = fs.super_block_header().generation;
2504        assert!(
2505            generation2 > generation1,
2506            "Expected generation bump, got {} vs {}",
2507            generation2,
2508            generation1
2509        );
2510        fs.close().await.expect("close failed");
2511    }
2512
2513    #[fuchsia::test]
2514    #[cfg(target_os = "fuchsia")]
2515    fn test_low_priority_compaction() {
2516        let mut executor = TestExecutor::new_with_fake_time();
2517        let mut fut = std::pin::pin!(async {
2518            use std::sync::Arc;
2519            use std::sync::atomic::{AtomicBool, Ordering};
2520
2521            let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
2522            let fs = FxFilesystemBuilder::new()
2523                .journal_options(super::JournalOptions {
2524                    reclaim_size: 65536,
2525                    ..Default::default()
2526                })
2527                .format(true)
2528                .open(device)
2529                .await
2530                .expect("open failed");
2531
2532            let _low = fasync::LowPriorityTask::new();
2533
2534            // Add some data to the tree.
2535            {
2536                let root_store = fs.root_store();
2537                let root_directory =
2538                    Directory::open(&root_store, root_store.root_directory_object_id())
2539                        .await
2540                        .expect("open failed");
2541                for i in 0..100 {
2542                    let mut transaction = fs
2543                        .root_store()
2544                        .new_transaction(
2545                            lock_keys![LockKey::object(
2546                                root_store.store_object_id(),
2547                                root_store.root_directory_object_id(),
2548                            )],
2549                            Options::default(),
2550                        )
2551                        .await
2552                        .expect("new_transaction failed");
2553                    root_directory
2554                        .create_child_file(&mut transaction, &format!("test{}", i))
2555                        .await
2556                        .expect("create_child_file failed");
2557                    transaction.commit().await.expect("commit failed");
2558                }
2559            }
2560
2561            // Spawn a task that polls every 1ms.
2562            let stop = Arc::new(AtomicBool::new(false));
2563            let stop_clone = stop.clone();
2564            let _normal_task = fasync::Task::spawn(async move {
2565                while !stop_clone.load(Ordering::Relaxed) {
2566                    fasync::Timer::new(fasync::MonotonicInstant::after(
2567                        MonotonicDuration::from_millis(1),
2568                    ))
2569                    .await;
2570                }
2571            });
2572
2573            // Trigger journal compaction.
2574            // We can do this by writing more data until outstanding > reclaim_size / 2.
2575            {
2576                let root_store = fs.root_store();
2577                let root_directory =
2578                    Directory::open(&root_store, root_store.root_directory_object_id())
2579                        .await
2580                        .expect("open failed");
2581                let mut i = 0;
2582                loop {
2583                    let mut transaction = fs
2584                        .root_store()
2585                        .new_transaction(
2586                            lock_keys![LockKey::object(
2587                                root_store.store_object_id(),
2588                                root_store.root_directory_object_id(),
2589                            )],
2590                            Options::default(),
2591                        )
2592                        .await
2593                        .expect("new_transaction failed");
2594                    root_directory
2595                        .create_child_file(&mut transaction, &format!("trigger{i}"))
2596                        .await
2597                        .expect("create_child_file failed");
2598                    transaction.commit().await.expect("commit failed");
2599
2600                    if fs.journal().inner.lock().compaction_running {
2601                        break;
2602                    }
2603                    TestExecutor::advance_to(fasync::MonotonicInstant::after(
2604                        MonotonicDuration::from_millis(1),
2605                    ))
2606                    .await;
2607                    i += 1;
2608                }
2609            }
2610
2611            // Compaction should now be running. Because of our 1ms poller, it should be yielding.
2612            // It will yield for 16ms at each point.
2613            for _ in 0..10 {
2614                TestExecutor::advance_to(fasync::MonotonicInstant::after(
2615                    MonotonicDuration::from_millis(1),
2616                ))
2617                .await;
2618                assert!(fs.journal().inner.lock().compaction_running);
2619            }
2620
2621            // Stop the normal task.
2622            stop.store(true, Ordering::Relaxed);
2623            TestExecutor::advance_to(fasync::MonotonicInstant::after(
2624                MonotonicDuration::from_millis(1),
2625            ))
2626            .await;
2627
2628            // Compaction should still be running because it hasn't been 4 ms since the normal task
2629            // finished.
2630            assert!(fs.journal().inner.lock().compaction_running);
2631
2632            // For the next 3ms, compaction should still be running.
2633            for _ in 0..3 {
2634                TestExecutor::advance_to(fasync::MonotonicInstant::after(
2635                    MonotonicDuration::from_millis(1),
2636                ))
2637                .await;
2638                assert!(fs.journal().inner.lock().compaction_running);
2639            }
2640
2641            // 1 more ms and compaction should be unblocked.
2642            TestExecutor::advance_to(fasync::MonotonicInstant::after(
2643                MonotonicDuration::from_millis(1),
2644            ))
2645            .await;
2646
2647            // When the executor next stalls, compaction should be done.
2648            let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2649            assert!(!fs.journal().inner.lock().compaction_running);
2650
2651            fs.close().await.expect("Close failed");
2652        });
2653        assert!(executor.run_until_stalled(&mut fut).is_ready());
2654    }
2655
2656    #[fuchsia::test]
2657    #[cfg(target_os = "fuchsia")]
2658    fn test_low_priority_compaction_deadline() {
2659        let mut executor = TestExecutor::new_with_fake_time();
2660        let mut fut = std::pin::pin!(async {
2661            use std::sync::Arc;
2662            use std::sync::atomic::{AtomicBool, Ordering};
2663
2664            let device = DeviceHolder::new(FakeDevice::new(216384, TEST_DEVICE_BLOCK_SIZE));
2665            let fs = FxFilesystemBuilder::new()
2666                .journal_options(super::JournalOptions {
2667                    reclaim_size: 65536,
2668                    ..Default::default()
2669                })
2670                .format(true)
2671                .open(device)
2672                .await
2673                .expect("open failed");
2674
2675            let _low = fasync::LowPriorityTask::new();
2676
2677            // Add some data to the tree.
2678            {
2679                let root_store = fs.root_store();
2680                let root_directory =
2681                    Directory::open(&root_store, root_store.root_directory_object_id())
2682                        .await
2683                        .expect("open failed");
2684                for i in 0..10 {
2685                    let mut transaction = fs
2686                        .root_store()
2687                        .new_transaction(
2688                            lock_keys![LockKey::object(
2689                                root_store.store_object_id(),
2690                                root_store.root_directory_object_id(),
2691                            )],
2692                            Options::default(),
2693                        )
2694                        .await
2695                        .expect("new_transaction failed");
2696                    root_directory
2697                        .create_child_file(&mut transaction, &format!("test{}", i))
2698                        .await
2699                        .expect("create_child_file failed");
2700                    transaction.commit().await.expect("commit failed");
2701                }
2702            }
2703
2704            // Spawn a task that polls every 3ms.
2705            let stop = Arc::new(AtomicBool::new(false));
2706            let stop_clone = stop.clone();
2707            let _normal_task = fasync::Task::spawn(async move {
2708                while !stop_clone.load(Ordering::Relaxed) {
2709                    fasync::Timer::new(fasync::MonotonicInstant::after(
2710                        MonotonicDuration::from_millis(3),
2711                    ))
2712                    .await;
2713                }
2714            });
2715
2716            // Trigger journal compaction.
2717            {
2718                let root_store = fs.root_store();
2719                let root_directory =
2720                    Directory::open(&root_store, root_store.root_directory_object_id())
2721                        .await
2722                        .expect("open failed");
2723                let mut i = 0;
2724                loop {
2725                    let mut transaction = fs
2726                        .root_store()
2727                        .new_transaction(
2728                            lock_keys![LockKey::object(
2729                                root_store.store_object_id(),
2730                                root_store.root_directory_object_id(),
2731                            )],
2732                            Options::default(),
2733                        )
2734                        .await
2735                        .expect("new_transaction failed");
2736                    root_directory
2737                        .create_child_file(&mut transaction, &format!("trigger{i}"))
2738                        .await
2739                        .expect("create_child_file failed");
2740                    transaction.commit().await.expect("commit failed");
2741
2742                    if fs.journal().inner.lock().compaction_running {
2743                        break;
2744                    }
2745                    TestExecutor::advance_to(fasync::MonotonicInstant::after(
2746                        MonotonicDuration::from_millis(1),
2747                    ))
2748                    .await;
2749                    i += 1;
2750                }
2751            }
2752
2753            // Advance time in 20ms increments. Each increment should allow compaction to make progress
2754            // on one item (since MAX_YIELD_DURATION is 16ms).
2755            let mut count = 0;
2756            for _ in 0..1000 {
2757                TestExecutor::advance_to(fasync::MonotonicInstant::after(
2758                    MonotonicDuration::from_millis(20),
2759                ))
2760                .await;
2761                if !fs.journal().inner.lock().compaction_running {
2762                    break;
2763                }
2764                count += 1;
2765            }
2766            assert!(!fs.journal().inner.lock().compaction_running);
2767
2768            // Make sure it took a few iterations to complete.  It's difficult to know what the exact
2769            // number should be.
2770            assert!(count > 200);
2771
2772            stop.store(true, Ordering::Relaxed);
2773            fs.close().await.expect("Close failed");
2774        });
2775        assert!(executor.run_until_stalled(&mut fut).is_ready());
2776    }
2777
2778    #[fuchsia::test]
2779    #[cfg(target_os = "fuchsia")]
2780    fn test_low_priority_compaction_no_yielding_when_full() {
2781        let mut executor = TestExecutor::new_with_fake_time();
2782        let mut fut = std::pin::pin!(async {
2783            use std::sync::Arc;
2784            use std::sync::atomic::{AtomicBool, Ordering};
2785
2786            let reclaim_size = 65536;
2787            let device = DeviceHolder::new(FakeDevice::new(216384, TEST_DEVICE_BLOCK_SIZE));
2788            let fs = FxFilesystemBuilder::new()
2789                .journal_options(super::JournalOptions { reclaim_size, ..Default::default() })
2790                .format(true)
2791                .open(device)
2792                .await
2793                .expect("open failed");
2794
2795            let _low = fasync::LowPriorityTask::new();
2796
2797            // Add some data to the tree.
2798            {
2799                let root_store = fs.root_store();
2800                let root_directory =
2801                    Directory::open(&root_store, root_store.root_directory_object_id())
2802                        .await
2803                        .expect("open failed");
2804                for i in 0..10 {
2805                    let mut transaction = fs
2806                        .root_store()
2807                        .new_transaction(
2808                            lock_keys![LockKey::object(
2809                                root_store.store_object_id(),
2810                                root_store.root_directory_object_id(),
2811                            )],
2812                            Options::default(),
2813                        )
2814                        .await
2815                        .expect("new_transaction failed");
2816                    root_directory
2817                        .create_child_file(&mut transaction, &format!("test{}", i))
2818                        .await
2819                        .expect("create_child_file failed");
2820                    transaction.commit().await.expect("commit failed");
2821                }
2822            }
2823
2824            // Spawn a task that polls every 3ms.
2825            let stop = Arc::new(AtomicBool::new(false));
2826            let stop_clone = stop.clone();
2827            let _normal_task = fasync::Task::spawn(async move {
2828                while !stop_clone.load(Ordering::Relaxed) {
2829                    fasync::Timer::new(fasync::MonotonicInstant::after(
2830                        MonotonicDuration::from_millis(3),
2831                    ))
2832                    .await;
2833                }
2834            });
2835
2836            // Trigger journal compaction, but this time we fill it up to 3/4 full.
2837            {
2838                let root_store = fs.root_store();
2839                let root_directory =
2840                    Directory::open(&root_store, root_store.root_directory_object_id())
2841                        .await
2842                        .expect("open failed");
2843                let mut i = 0;
2844                loop {
2845                    let mut transaction = fs
2846                        .root_store()
2847                        .new_transaction(
2848                            lock_keys![LockKey::object(
2849                                root_store.store_object_id(),
2850                                root_store.root_directory_object_id(),
2851                            )],
2852                            Options::default(),
2853                        )
2854                        .await
2855                        .expect("new_transaction failed");
2856                    root_directory
2857                        .create_child_file(&mut transaction, &format!("trigger{i}"))
2858                        .await
2859                        .expect("create_child_file failed");
2860                    transaction.commit().await.expect("commit failed");
2861
2862                    let outstanding = {
2863                        let inner = fs.journal().inner.lock();
2864                        fs.journal().objects.last_end_offset()
2865                            - inner.super_block_header.journal_checkpoint.file_offset
2866                    };
2867                    if outstanding >= reclaim_size * 7 / 8 {
2868                        break;
2869                    }
2870                    // We don't advance time here to try and reach 3/4 before compaction can yield.
2871                    i += 1;
2872                }
2873            }
2874
2875            // Advancing by 4ms should be enough to wake compaction up.
2876            TestExecutor::advance_to(fasync::MonotonicInstant::after(
2877                MonotonicDuration::from_millis(4),
2878            ))
2879            .await;
2880
2881            // When the executor next stalls, compaction should be done.
2882            let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2883            assert!(!fs.journal().inner.lock().compaction_running);
2884
2885            stop.store(true, Ordering::Relaxed);
2886            fs.close().await.expect("Close failed");
2887        });
2888        assert!(executor.run_until_stalled(&mut fut).is_ready());
2889    }
2890}
2891
2892#[cfg(fuzz)]
2893mod fuzz {
2894    use fuzz::fuzz;
2895
2896    #[fuzz]
2897    fn fuzz_journal_bytes(input: Vec<u8>) {
2898        use crate::filesystem::FxFilesystem;
2899        use fuchsia_async as fasync;
2900        use std::io::Write;
2901        use storage_device::DeviceHolder;
2902        use storage_device::fake_device::FakeDevice;
2903
2904        fasync::SendExecutorBuilder::new().num_threads(4).build().run(async move {
2905            let device = DeviceHolder::new(FakeDevice::new(32768, 512));
2906            let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2907            fs.journal().inner.lock().writer.write_all(&input).expect("write failed");
2908            fs.close().await.expect("close failed");
2909            let device = fs.take_device().await;
2910            device.reopen(false);
2911            if let Ok(fs) = FxFilesystem::open(device).await {
2912                // `close()` can fail if there were objects to be tombstoned. If the said object is
2913                // corrupted, there will be an error when we compact the journal.
2914                let _ = fs.close().await;
2915            }
2916        });
2917    }
2918
2919    #[fuzz]
2920    fn fuzz_journal(input: Vec<super::JournalRecord>) {
2921        use crate::filesystem::FxFilesystem;
2922        use fuchsia_async as fasync;
2923        use storage_device::DeviceHolder;
2924        use storage_device::fake_device::FakeDevice;
2925
2926        fasync::SendExecutorBuilder::new().num_threads(4).build().run(async move {
2927            let device = DeviceHolder::new(FakeDevice::new(32768, 512));
2928            let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2929            {
2930                let mut inner = fs.journal().inner.lock();
2931                for record in &input {
2932                    let _ = inner.writer.write_record(record);
2933                }
2934            }
2935            fs.close().await.expect("close failed");
2936            let device = fs.take_device().await;
2937            device.reopen(false);
2938            if let Ok(fs) = FxFilesystem::open(device).await {
2939                // `close()` can fail if there were objects to be tombstoned. If the said object is
2940                // corrupted, there will be an error when we compact the journal.
2941                let _ = fs.close().await;
2942            }
2943        });
2944    }
2945}