Skip to main content

fxfs/
object_store.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
5pub mod allocator;
6pub mod caching_object_handle;
7pub mod data_object_handle;
8pub mod directory;
9pub mod extent;
10mod extent_mapping_iterator;
11mod extent_record;
12pub mod flush;
13use crate::filesystem::FlushReason;
14pub mod graveyard;
15mod install;
16pub mod journal;
17mod key_manager;
18pub mod merge;
19pub mod object_manager;
20pub mod object_record;
21pub mod project_id;
22mod store_object_handle;
23pub mod transaction;
24mod tree;
25mod tree_cache;
26pub mod volume;
27
28pub use data_object_handle::{
29    DataObjectHandle, DataObjectState, DirectWriter, FileExtent, FsverityStateInner, RangeType,
30};
31pub use directory::Directory;
32pub use object_record::{
33    ChildValue, DirType, FscryptDirInfo, FscryptPolicyFlags, LEGACY_FSCRYPT_FLAGS,
34    ObjectDescriptor, PosixAttributes, Timestamp,
35};
36pub use store_object_handle::{MAX_INLINE_XATTR_SIZE, SetExtendedAttributeMode, StoreObjectHandle};
37
38use crate::errors::FxfsError;
39use crate::filesystem::{
40    ApplyContext, ApplyMode, FxFilesystem, JournalingObject, MAX_FILE_SIZE, SyncOptions,
41    TruncateGuard,
42};
43use crate::log::*;
44use crate::lsm_tree::cache::ObjectCache;
45use crate::lsm_tree::types::{Existence, Item, ItemRef, LayerIterator};
46use crate::lsm_tree::{LSMTree, Query, open_layers};
47use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, ObjectProperties, ReadObjectHandle};
48use crate::object_store::allocator::Allocator;
49use crate::object_store::graveyard::Graveyard;
50use crate::object_store::journal::{JournalCheckpoint, JournalCheckpointV32, JournaledTransaction};
51use crate::object_store::key_manager::KeyManager;
52use crate::object_store::transaction::{
53    AssocObj, AssociatedObject, LockKey, LockKeys, MutationV59, ObjectMutationIterator,
54    ObjectStoreMutation, Operation, Options, ReservationOptions, Transaction, WriteGuard,
55    lock_keys,
56};
57use crate::serialized_types::{
58    AES_JOURNAL_ENCRYPTION_VERSION, DEFAULT_MAX_SERIALIZED_RECORD_SIZE, Version, Versioned,
59    VersionedLatest,
60};
61use anyhow::{Context, Error, anyhow, bail, ensure};
62use async_trait::async_trait;
63use fidl_fuchsia_io as fio;
64use fprint::TypeFingerprint;
65use fuchsia_sync::Mutex;
66use fxfs_crypto::ff1::Ff1;
67use fxfs_crypto::{
68    CipherHolder, Crypt, JournalCipher, JournalXtsCipher, KeyPurpose, ObjectType, StreamCipher,
69    UnwrappedKey, key_to_cipher,
70};
71use rand::Rng as _;
72use scopeguard::ScopeGuard;
73use serde::{Deserialize, Serialize};
74use std::collections::HashSet;
75use std::fmt;
76use std::num::NonZero;
77use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
78use std::sync::{Arc, OnceLock, Weak};
79use storage_device::Device;
80use storage_units::BlockSize;
81use uuid::Uuid;
82
83pub use extent::Extent;
84pub use extent_record::{ExtentMode, ExtentValue};
85pub use object_record::{
86    AttributeId, AttributeKey, BytesAndNodes, EncryptionKey, EncryptionKeys,
87    ExtendedAttributeValue, FsverityMetadata, FxfsKey, FxfsKeyV49, ObjectAttributes, ObjectKey,
88    ObjectKeyData, ObjectKind, ObjectValue, ProjectProperty, RootDigest,
89};
90pub use project_id::{ProjectId, ProjectIdExt};
91pub use transaction::Mutation;
92
93// For encrypted stores, the lower 32 bits of the object ID are encrypted to make side-channel
94// attacks more difficult. This mask can be used to extract the hi part of the object ID.
95const OBJECT_ID_HI_MASK: u64 = 0xffffffff00000000;
96
97// At time of writing, this threshold limits transactions that delete extents to about 10,000 bytes.
98const TRANSACTION_MUTATION_THRESHOLD: usize = 200;
99
100// The number of keys we want to pre-cache for flushes.  We cache up to 2 keys.  One is for the
101// current flush, and one is pre-cached for the next flush so that we don't need to call the crypt
102// service during a flush.  To understand why, consider two threads T1 and T2.  Just prior to
103// committing a transaction, T1 ensures there are two keys.  Then, before T1 has committed the
104// transaction, T2 flushes and consumes one of the keys.  T1 then commits the transaction. The next
105// time a flush occurs, there's a key ready. If T1 had only ensured there was one key, there'd be no
106// key.
107const CACHED_KEYS_LIMIT: usize = 2;
108
109// Encrypted files and directories use the fscrypt key (identified by `FSCRYPT_KEY_ID`) to encrypt
110// file contents and filenames respectively. All non-fscrypt encrypted files otherwise default to
111// using the `VOLUME_DATA_KEY_ID` key. Note, the filesystem always uses the `VOLUME_DATA_KEY_ID`
112// key to encrypt large extended attributes. Thus, encrypted files and directories with large
113// xattrs will have both an fscrypt and volume data key.
114pub const VOLUME_DATA_KEY_ID: u64 = 0;
115pub const FSCRYPT_KEY_ID: u64 = 1;
116
117/// DataObjectHandle stores an owner that must implement this trait, which allows the handle to get
118/// back to an ObjectStore.
119pub trait HandleOwner: AsRef<ObjectStore> + Send + Sync + 'static {}
120
121/// StoreInfo stores information about the object store.  This is stored within the parent object
122/// store, and is used, for example, to get the persistent layer objects.
123pub type StoreInfo = StoreInfoV52;
124
125#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
126pub struct StoreInfoV52 {
127    /// The globally unique identifier for the associated object store. If unset, will be all zero.
128    guid: [u8; 16],
129
130    /// The last used object ID.  Note that this field is not accurate in memory; ObjectStore's
131    /// last_object_id field is the one to use in that case.  Technically, this might not be the
132    /// last object ID used for the latest transaction that created an object because we use this at
133    /// the point of creating the object but before we commit the transaction.  Transactions can
134    /// then get committed in an arbitrary order (or not at all).
135    last_object_id: LastObjectIdInfo,
136
137    /// Object ids for layers.  TODO(https://fxbug.dev/42178036): need a layer of indirection here
138    /// so we can support snapshots.
139    pub layers: Vec<u64>,
140
141    /// The object ID for the root directory.
142    root_directory_object_id: u64,
143
144    /// The object ID for the graveyard.
145    graveyard_directory_object_id: u64,
146
147    /// The number of live objects in the store.  This should *not* be trusted; it can be invalid
148    /// due to filesystem inconsistencies.
149    object_count: u64,
150
151    /// The (wrapped) key that encrypted mutations should use.
152    mutations_key: Option<FxfsKeyV49>,
153
154    /// Mutations for the store are encrypted using a stream cipher.  To decrypt the mutations, we
155    /// need to know the offset in the cipher stream to start it.
156    mutations_cipher_offset: u64,
157
158    /// If we have to flush the store whilst we do not have the key, we need to write the encrypted
159    /// mutations to an object. This is the object ID of that file if it exists.
160    pub encrypted_mutations_object_id: u64,
161
162    /// A directory for storing internal files in a directory structure. Holds INVALID_OBJECT_ID
163    /// when the directory doesn't yet exist.
164    internal_directory_object_id: u64,
165}
166
167#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint)]
168enum LastObjectIdInfo {
169    Unencrypted {
170        id: u64,
171    },
172    Encrypted {
173        /// The *unencrypted* value of the last object ID.
174        id: u64,
175
176        /// Object IDs are encrypted to reduce the amount of information that sequential object IDs
177        /// reveal (such as the number of files in the system and the ordering of their creation in
178        /// time).  Only the bottom 32 bits of the object ID are encrypted whilst the top 32 bits
179        /// will increment after 2^32 object IDs have been used and this allows us to roll the key.
180        key: FxfsKeyV49,
181    },
182    Low32Bit,
183}
184
185impl Default for LastObjectIdInfo {
186    fn default() -> Self {
187        LastObjectIdInfo::Unencrypted { id: 0 }
188    }
189}
190
191impl StoreInfo {
192    /// Returns the parent objects for this store.
193    pub fn parent_objects(&self) -> Vec<u64> {
194        // We should not include the ID of the store itself, since that should be referred to in the
195        // volume directory.
196        let mut objects = self.layers.to_vec();
197        if self.encrypted_mutations_object_id != INVALID_OBJECT_ID {
198            objects.push(self.encrypted_mutations_object_id);
199        }
200        objects
201    }
202}
203
204// TODO(https://fxbug.dev/42178037): We should test or put checks in place to ensure this limit isn't exceeded.
205// It will likely involve placing limits on the maximum number of layers.
206pub const MAX_STORE_INFO_SERIALIZED_SIZE: usize = 131072;
207
208// This needs to be large enough to accommodate the maximum amount of unflushed data (data that is
209// in the journal but hasn't yet been written to layer files) for a store.  We set a limit because
210// we want to limit the amount of memory use in the case the filesystem is corrupt or under attack.
211pub const MAX_ENCRYPTED_MUTATIONS_SIZE: usize = 8 * journal::DEFAULT_RECLAIM_SIZE as usize;
212
213#[derive(Default)]
214pub struct HandleOptions {
215    /// If true, transactions used by this handle will skip journal space checks.
216    pub skip_journal_checks: bool,
217    /// If true, data written to any attribute of this handle will not have per-block checksums
218    /// computed.
219    pub skip_checksums: bool,
220    /// If true, any files using fsverity will not attempt to perform any verification. This is
221    /// useful to open an object without the correct encryption keys to look at the metadata.
222    pub skip_fsverity: bool,
223}
224
225/// Parameters for encrypting a newly created object.
226pub struct ObjectEncryptionOptions {
227    /// If set, the keys are treated as permanent and never evicted from the KeyManager cache.
228    /// This is necessary when keys are managed by another store; for example, the layer files
229    /// of a child store are objects in the root store, but they are encrypted with keys from the
230    /// child store.  Generally, most objects should have this set to `false`.
231    pub permanent: bool,
232    pub key_id: u64,
233    pub key: EncryptionKey,
234    pub unwrapped_key: UnwrappedKey,
235}
236
237pub struct StoreOptions {
238    /// The store is unencrypted if store is none.
239    pub crypt: Option<Arc<dyn Crypt>>,
240}
241
242impl Default for StoreOptions {
243    fn default() -> Self {
244        Self { crypt: None }
245    }
246}
247
248#[derive(Default)]
249pub struct NewChildStoreOptions {
250    pub options: StoreOptions,
251
252    /// Specifies the object ID in the root store to be used for the store.  If set to
253    /// INVALID_OBJECT_ID (the default and typical case), a suitable ID will be chosen.
254    pub object_id: u64,
255
256    /// If true, reserve all 32 bit object_ids.  All new objects will start with IDs exceeding
257    /// 0x1_0000_0000.
258    pub reserve_32bit_object_ids: bool,
259
260    /// Object IDs will be restricted to 32 bits.  This involves a less performant algorithm and so
261    /// should not be used unless necessary.
262    pub low_32_bit_object_ids: bool,
263
264    /// If set, use this GUID for the new store.
265    pub guid: Option<[u8; 16]>,
266}
267
268pub type EncryptedMutations = EncryptedMutationsV49;
269
270#[derive(Clone, Default, Deserialize, Serialize, TypeFingerprint)]
271pub struct EncryptedMutationsV49 {
272    // Information about the mutations are held here, but the actual encrypted data is held within
273    // data.  For each transaction, we record the checkpoint and the count of mutations within the
274    // transaction.  The checkpoint is required for the log file offset (which we need to apply the
275    // mutations), and the version so that we can correctly decode the mutation after it has been
276    // decrypted. The count specifies the number of serialized mutations encoded in |data|.
277    transactions: Vec<(JournalCheckpointV32, u64)>,
278
279    // The encrypted mutations.
280    #[serde(with = "crate::zerocopy_serialization")]
281    data: Vec<u8>,
282
283    // If the mutations key was rolled, this holds the offset in `data` where the new key should
284    // apply.
285    mutations_key_roll: Vec<(usize, FxfsKeyV49)>,
286}
287
288impl std::fmt::Debug for EncryptedMutations {
289    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> Result<(), fmt::Error> {
290        f.debug_struct("EncryptedMutations")
291            .field("transactions", &self.transactions)
292            .field("len", &self.data.len())
293            .field(
294                "mutations_key_roll",
295                &self.mutations_key_roll.iter().map(|k| k.0).collect::<Vec<usize>>(),
296            )
297            .finish()
298    }
299}
300
301impl Versioned for EncryptedMutations {
302    fn max_serialized_size() -> Option<u64> {
303        Some(MAX_ENCRYPTED_MUTATIONS_SIZE as u64)
304    }
305}
306
307impl EncryptedMutations {
308    fn from_replayed_mutations(
309        store_object_id: u64,
310        transactions: Vec<JournaledTransaction>,
311    ) -> Self {
312        let mut this = Self::default();
313        for JournaledTransaction { checkpoint, non_root_mutations, .. } in transactions {
314            for (object_id, mutation) in non_root_mutations {
315                if store_object_id == object_id {
316                    if let Mutation::EncryptedObjectStore(data) = mutation {
317                        this.push(&checkpoint, data);
318                    } else if let Mutation::UpdateMutationsKey(key) = mutation {
319                        this.mutations_key_roll.push((this.data.len(), key.into()));
320                    }
321                }
322            }
323        }
324        this
325    }
326
327    fn extend(&mut self, other: &EncryptedMutations) {
328        self.transactions.extend_from_slice(&other.transactions[..]);
329        self.mutations_key_roll.extend(
330            other
331                .mutations_key_roll
332                .iter()
333                .map(|(offset, key)| (offset + self.data.len(), key.clone())),
334        );
335        self.data.extend_from_slice(&other.data[..]);
336    }
337
338    fn push(&mut self, checkpoint: &JournalCheckpoint, data: Box<[u8]>) {
339        let len = data.len();
340        self.data.append(&mut data.into());
341        if checkpoint.version >= AES_JOURNAL_ENCRYPTION_VERSION {
342            // Chunks of the same transaction share the same checkpoint, so coalesce their lengths.
343            if let Some((last_checkpoint, total_len)) = self.transactions.last_mut() {
344                if last_checkpoint.file_offset == checkpoint.file_offset {
345                    *total_len += len as u64;
346                    return;
347                }
348            }
349            self.transactions.push((checkpoint.clone(), len as u64));
350        } else {
351            // If the checkpoint is the same as the last mutation we pushed, increment the count.
352            if let Some((last_checkpoint, count)) = self.transactions.last_mut() {
353                if last_checkpoint.file_offset == checkpoint.file_offset {
354                    *count += 1;
355                    return;
356                }
357            }
358            self.transactions.push((checkpoint.clone(), 1));
359        }
360    }
361}
362
363pub enum LockState {
364    Locked,
365    Unencrypted,
366    Unlocked {
367        crypt: Arc<dyn Crypt>,
368        cached_keys: Vec<(NonZero<u64>, EncryptionKey, UnwrappedKey)>,
369        mutations_cipher: JournalCipher,
370    },
371
372    // The store is unlocked, but in a read-only state, and no flushes or other operations will be
373    // performed on the store.
374    UnlockedReadOnly(Arc<dyn Crypt>),
375
376    // The store is encrypted but is now in an unusable state (due to a failure to sync the journal
377    // after locking the store).  The store cannot be unlocked.
378    Invalid,
379
380    // Before we've read the StoreInfo we might not know whether the store is Locked or Unencrypted.
381    // This can happen when lazily opening stores (ObjectManager::lazy_open_store).
382    Unknown,
383
384    // The store is in the process of being locked.  Whilst the store is being locked, the store
385    // isn't usable; assertions will trip if any mutations are applied.
386    Locking,
387
388    // Whilst we're unlocking, we will replay encrypted mutations.  The store isn't usable until
389    // it's in the Unlocked state.
390    Unlocking,
391
392    // The store has been deleted.
393    Deleted,
394}
395
396impl fmt::Debug for LockState {
397    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
398        formatter.write_str(match self {
399            LockState::Locked => "Locked",
400            LockState::Unencrypted => "Unencrypted",
401            LockState::Unlocked { .. } => "Unlocked",
402            LockState::UnlockedReadOnly(..) => "UnlockedReadOnly",
403            LockState::Invalid => "Invalid",
404            LockState::Unknown => "Unknown",
405            LockState::Locking => "Locking",
406            LockState::Unlocking => "Unlocking",
407            LockState::Deleted => "Deleted",
408        })
409    }
410}
411
412enum LastObjectId {
413    // This is used when the store is encrypted, but the key and ID isn't yet available.
414    Pending,
415
416    Unencrypted {
417        id: u64,
418    },
419
420    Encrypted {
421        // The *unencrypted* value of the last object ID.
422        id: u64,
423
424        // Encrypted stores will use a cipher to obfuscate the object ID.
425        cipher: Box<Ff1>,
426    },
427
428    Low32Bit {
429        reserved: HashSet<u32>,
430        unreserved: Vec<u32>,
431    },
432}
433
434impl LastObjectId {
435    /// Returns true if object IDs require reservations.
436    fn uses_reserved_ids(&self) -> bool {
437        matches!(self, LastObjectId::Low32Bit { .. })
438    }
439
440    /// Tries to get the next object ID.  Returns None if a new cipher is required because all
441    /// object IDs that can be generated with the current cipher have been exhausted, or if only
442    /// using the lower 32 bits which requires an async algorithm.
443    fn try_get_next(&mut self) -> Option<NonZero<u64>> {
444        match self {
445            LastObjectId::Unencrypted { id } => {
446                NonZero::new(id.wrapping_add(1)).inspect(|next| *id = next.get())
447            }
448            LastObjectId::Encrypted { id, cipher } => {
449                let mut next = *id;
450                let hi = next & OBJECT_ID_HI_MASK;
451                loop {
452                    if next as u32 == u32::MAX {
453                        return None;
454                    }
455                    next += 1;
456                    let candidate = hi | cipher.encrypt(next as u32) as u64;
457                    if let Some(candidate) = NonZero::new(candidate) {
458                        *id = next;
459                        return Some(candidate);
460                    }
461                }
462            }
463            _ => None,
464        }
465    }
466
467    /// Returns INVALID_OBJECT_ID if it's not possible to peek at the next object ID.
468    fn peek_next(&self) -> u64 {
469        match self {
470            LastObjectId::Unencrypted { id } => id.wrapping_add(1),
471            LastObjectId::Encrypted { id, cipher } => {
472                let mut next = *id;
473                let hi = next & OBJECT_ID_HI_MASK;
474                loop {
475                    if next as u32 == u32::MAX {
476                        return INVALID_OBJECT_ID;
477                    }
478                    next += 1;
479                    let candidate = hi | cipher.encrypt(next as u32) as u64;
480                    if candidate != INVALID_OBJECT_ID {
481                        return candidate;
482                    }
483                }
484            }
485            _ => INVALID_OBJECT_ID,
486        }
487    }
488
489    /// Returns INVALID_OBJECT_ID for algorithms that don't use the last ID.
490    fn id(&self) -> u64 {
491        match self {
492            LastObjectId::Unencrypted { id } | LastObjectId::Encrypted { id, .. } => *id,
493            _ => INVALID_OBJECT_ID,
494        }
495    }
496
497    /// Returns true if `id` is reserved (it must be 32 bits).
498    fn is_reserved(&self, id: u64) -> bool {
499        match self {
500            LastObjectId::Low32Bit { reserved, .. } => {
501                if let Ok(id) = id.try_into() {
502                    reserved.contains(&id)
503                } else {
504                    false
505                }
506            }
507            _ => false,
508        }
509    }
510
511    /// Reserves `id`.
512    fn reserve(&mut self, id: u64) {
513        match self {
514            LastObjectId::Low32Bit { reserved, .. } => {
515                assert!(reserved.insert(id.try_into().unwrap()))
516            }
517            _ => unreachable!(),
518        }
519    }
520
521    /// Unreserves `id`.
522    fn unreserve(&mut self, id: u64) {
523        match self {
524            LastObjectId::Low32Bit { unreserved, .. } => {
525                // To avoid races, where a reserved ID transitions from being reserved to being
526                // actually used in a committed transaction, we delay updating `reserved` until a
527                // suitable point.
528                //
529                // On thread A, we might have:
530                //
531                //   A1. Commit transaction (insert a record into the LSM tree that uses ID)
532                //   A2. `unreserve`
533                //
534                // And on another thread B, we might have:
535                //
536                //   B1. Drain `unreserved`.
537                //   B2. Check tree and `reserved` to see if ID is used.
538                //
539                // B2 will involve calling `LsmTree::layer_set` which should be thought of as a
540                // snapshot, so the change A1 might not be visible to thread B, but it won't matter
541                // because `reserved` will still include the ID.  So long as each thread does the
542                // operations in this order, it should be safe.
543                unreserved.push(id.try_into().unwrap())
544            }
545            _ => {}
546        }
547    }
548
549    /// Removes `unreserved` IDs from the `reserved` list.
550    fn drain_unreserved(&mut self) {
551        match self {
552            LastObjectId::Low32Bit { reserved, unreserved } => {
553                for u in unreserved.drain(..) {
554                    assert!(reserved.remove(&u));
555                }
556            }
557            _ => {}
558        }
559    }
560}
561
562pub struct ReservedId<'a>(&'a ObjectStore, NonZero<u64>);
563
564impl<'a> ReservedId<'a> {
565    pub fn new(store: &'a ObjectStore, id: NonZero<u64>) -> Self {
566        Self(store, id)
567    }
568
569    pub fn get(&self) -> u64 {
570        self.1.get()
571    }
572
573    /// The caller takes responsibility for this id.
574    #[must_use]
575    pub fn release(self) -> NonZero<u64> {
576        let id = self.1;
577        std::mem::forget(self);
578        id
579    }
580}
581
582impl Drop for ReservedId<'_> {
583    fn drop(&mut self) {
584        self.0.last_object_id.lock().unreserve(self.1.get());
585    }
586}
587
588/// An object store supports a file like interface for objects.  Objects are keyed by a 64 bit
589/// identifier.  And object store has to be backed by a parent object store (which stores metadata
590/// for the object store).  The top-level object store (a.k.a. the root parent object store) is
591/// in-memory only.
592pub struct ObjectStore {
593    parent_store: Option<Arc<ObjectStore>>,
594    store_object_id: u64,
595    device: Arc<dyn Device>,
596    block_size: BlockSize,
597    filesystem: Weak<FxFilesystem>,
598    // Lock ordering: This must be taken before `lock_state`.
599    store_info: Mutex<Option<StoreInfo>>,
600    tree: LSMTree<ObjectKey, ObjectValue>,
601
602    // When replaying the journal, the store cannot read StoreInfo until the whole journal
603    // has been replayed, so during that time, store_info_handle will be None and records
604    // just get sent to the tree. Once the journal has been replayed, we can open the store
605    // and load all the other layer information.
606    store_info_handle: OnceLock<DataObjectHandle<ObjectStore>>,
607
608    // Current lock state of the store.
609    // Lock ordering: This must be taken after `store_info`.
610    lock_state: Mutex<LockState>,
611    pub key_manager: KeyManager,
612
613    // Enable/disable tracing.
614    trace: AtomicBool,
615
616    // Informational counters for events occurring within the store.
617    counters: Mutex<ObjectStoreCounters>,
618
619    // These are updated in performance-sensitive code paths so we use atomics instead of counters.
620    device_read_ops: AtomicU64,
621    device_write_ops: AtomicU64,
622    logical_read_ops: AtomicU64,
623    logical_write_ops: AtomicU64,
624    graveyard_entries: AtomicU64,
625
626    // Contains the last object ID and, optionally, a cipher to be used when generating new object
627    // IDs.
628    last_object_id: Mutex<LastObjectId>,
629
630    // An optional callback to be invoked each time the ObjectStore flushes.  The callback is
631    // invoked at the end of flush, while the write lock is still held.
632    flush_callback: Mutex<Option<Box<dyn Fn(&ObjectStore) + Send + Sync + 'static>>>,
633}
634
635#[derive(Clone, Default)]
636struct ObjectStoreCounters {
637    mutations_applied: u64,
638    mutations_dropped: u64,
639    num_flushes: u64,
640    last_flush_time: Option<std::time::SystemTime>,
641}
642
643impl ObjectStore {
644    fn new(
645        parent_store: Option<Arc<ObjectStore>>,
646        store_object_id: u64,
647        filesystem: Arc<FxFilesystem>,
648        store_info: Option<StoreInfo>,
649        object_cache: Option<Box<dyn ObjectCache<ObjectKey, ObjectValue>>>,
650        lock_state: LockState,
651        last_object_id: LastObjectId,
652    ) -> Arc<ObjectStore> {
653        let device = filesystem.device();
654        let block_size = filesystem.block_size();
655        Arc::new(ObjectStore {
656            parent_store,
657            store_object_id,
658            device,
659            block_size,
660            filesystem: Arc::downgrade(&filesystem),
661            store_info: Mutex::new(store_info),
662            tree: LSMTree::new(merge::merge, object_cache),
663            store_info_handle: OnceLock::new(),
664            lock_state: Mutex::new(lock_state),
665            key_manager: KeyManager::new(),
666            trace: AtomicBool::new(false),
667            counters: Mutex::new(ObjectStoreCounters::default()),
668            device_read_ops: AtomicU64::new(0),
669            device_write_ops: AtomicU64::new(0),
670            logical_read_ops: AtomicU64::new(0),
671            logical_write_ops: AtomicU64::new(0),
672            graveyard_entries: AtomicU64::new(0),
673            last_object_id: Mutex::new(last_object_id),
674            flush_callback: Mutex::new(None),
675        })
676    }
677
678    fn new_empty(
679        parent_store: Option<Arc<ObjectStore>>,
680        store_object_id: u64,
681        filesystem: Arc<FxFilesystem>,
682        object_cache: Option<Box<dyn ObjectCache<ObjectKey, ObjectValue>>>,
683    ) -> Arc<Self> {
684        Self::new(
685            parent_store,
686            store_object_id,
687            filesystem,
688            Some(StoreInfo::default()),
689            object_cache,
690            LockState::Unencrypted,
691            LastObjectId::Unencrypted { id: 0 },
692        )
693    }
694
695    /// Cycle breaker constructor that returns an ObjectStore without a filesystem.
696    /// This should only be used from super block code.
697    pub fn new_root_parent(
698        device: Arc<dyn Device>,
699        block_size: BlockSize,
700        store_object_id: u64,
701    ) -> Self {
702        ObjectStore {
703            parent_store: None,
704            store_object_id,
705            device,
706            block_size,
707            filesystem: Weak::<FxFilesystem>::new(),
708            store_info: Mutex::new(Some(StoreInfo::default())),
709            tree: LSMTree::new(merge::merge, None),
710            store_info_handle: OnceLock::new(),
711            lock_state: Mutex::new(LockState::Unencrypted),
712            key_manager: KeyManager::new(),
713            trace: AtomicBool::new(false),
714            counters: Mutex::new(ObjectStoreCounters::default()),
715            device_read_ops: AtomicU64::new(0),
716            device_write_ops: AtomicU64::new(0),
717            logical_read_ops: AtomicU64::new(0),
718            logical_write_ops: AtomicU64::new(0),
719            graveyard_entries: AtomicU64::new(0),
720            last_object_id: Mutex::new(LastObjectId::Unencrypted { id: 0 }),
721            flush_callback: Mutex::new(None),
722        }
723    }
724
725    /// Used to set filesystem on root_parent stores at bootstrap time after the filesystem has
726    /// been created.
727    pub fn attach_filesystem(mut this: ObjectStore, filesystem: Arc<FxFilesystem>) -> ObjectStore {
728        this.filesystem = Arc::downgrade(&filesystem);
729        this
730    }
731
732    /// Acquires appropriate locks and starts a new transaction.  A transaction should be associated
733    /// with a store, even though it may have mutations for other objects and parent stores.
734    pub async fn new_transaction<'a>(
735        &self,
736        locks: LockKeys,
737        options: Options<'a>,
738    ) -> Result<Transaction<'a>, Error> {
739        let fs = self.filesystem();
740        Transaction::new(fs, options, locks).await
741    }
742
743    /// Ensures that `cached_keys` is filled up to `CACHED_KEYS_LIMIT`.
744    async fn pre_cache_keys(&self) -> Result<(), Error> {
745        let crypt = match &*self.lock_state.lock() {
746            LockState::Unlocked { cached_keys, crypt, .. }
747                if cached_keys.len() < CACHED_KEYS_LIMIT =>
748            {
749                crypt.clone()
750            }
751            _ => return Ok(()),
752        };
753        loop {
754            let parent_store = self.parent_store.as_ref().unwrap();
755
756            // We assert that the parent store is not using object IDs that need reservation, since
757            // if it did, release() would leak the reservation below.
758            assert!(!parent_store.last_object_id.lock().uses_reserved_ids());
759
760            // Allocate a raw ID from the parent store.  Since the parent store is unencrypted, this
761            // is fast and won't block.
762            let raw_id = {
763                let reserved_id = parent_store
764                    .maybe_get_next_object_id()
765                    .expect("maybe_get_next_object_id failed on parent store");
766                reserved_id.release()
767            };
768
769            let (wrapped, unwrapped) = match crypt.create_key(raw_id.get(), KeyPurpose::Data).await
770            {
771                Ok(v) => v,
772                Err(error) => {
773                    log::warn!(
774                        error:?,
775                        store_id = self.store_object_id();
776                        "Failed to pre-cache key"
777                    );
778                    return Err(error.into());
779                }
780            };
781
782            let mut lock_state = self.lock_state.lock();
783            if let LockState::Unlocked { cached_keys, .. } = &mut *lock_state {
784                cached_keys.push((raw_id, EncryptionKey::Fxfs(wrapped), unwrapped));
785                if cached_keys.len() >= CACHED_KEYS_LIMIT {
786                    break;
787                }
788            } else {
789                // Store was locked while we were awaiting; discard the key.
790                break;
791            }
792        }
793        Ok(())
794    }
795
796    /// Create a child store. It is a multi-step process:
797    ///
798    ///   1. Call `ObjectStore::new_child_store`.
799    ///   2. Register the store with the object-manager.
800    ///   3. Call `ObjectStore::create` to write the store-info.
801    ///
802    /// If the procedure fails, care must be taken to unregister store with the object-manager.
803    ///
804    /// The steps have to be separate because of lifetime issues when working with a transaction.
805    async fn new_child_store(
806        self: &Arc<Self>,
807        transaction: &mut Transaction<'_>,
808        options: NewChildStoreOptions,
809        object_cache: Option<Box<dyn ObjectCache<ObjectKey, ObjectValue>>>,
810    ) -> Result<Arc<Self>, Error> {
811        ensure!(
812            !options.reserve_32bit_object_ids || !options.low_32_bit_object_ids,
813            FxfsError::InvalidArgs
814        );
815        let handle = if let Some(object_id) = NonZero::new(options.object_id) {
816            self.update_last_object_id(object_id.get());
817            let handle = ObjectStore::create_object_with_id(
818                self,
819                transaction,
820                ReservedId::new(self, object_id),
821                HandleOptions::default(),
822                None,
823            )?;
824            handle
825        } else {
826            ObjectStore::create_object(self, transaction, HandleOptions::default(), None).await?
827        };
828        let filesystem = self.filesystem();
829        let id = if options.reserve_32bit_object_ids { 0x1_0000_0000 } else { 0 };
830        let (last_object_id, last_object_id_in_memory) = if options.low_32_bit_object_ids {
831            (
832                LastObjectIdInfo::Low32Bit,
833                LastObjectId::Low32Bit { reserved: HashSet::new(), unreserved: Vec::new() },
834            )
835        } else if let Some(crypt) = &options.options.crypt {
836            let (object_id_wrapped, object_id_unwrapped) =
837                crypt.create_key(handle.object_id(), KeyPurpose::Metadata).await?;
838            (
839                LastObjectIdInfo::Encrypted { id, key: object_id_wrapped },
840                LastObjectId::Encrypted { id, cipher: Box::new(Ff1::new(&object_id_unwrapped)) },
841            )
842        } else {
843            (LastObjectIdInfo::Unencrypted { id }, LastObjectId::Unencrypted { id })
844        };
845        let store = if let Some(crypt) = options.options.crypt {
846            let (wrapped_key, unwrapped_key) =
847                crypt.create_key(handle.object_id(), KeyPurpose::Metadata).await?;
848            Self::new(
849                Some(self.clone()),
850                handle.object_id(),
851                filesystem.clone(),
852                Some(StoreInfo {
853                    mutations_key: Some(wrapped_key),
854                    last_object_id,
855                    guid: options.guid.unwrap_or_else(|| *Uuid::new_v4().as_bytes()),
856                    ..Default::default()
857                }),
858                object_cache,
859                LockState::Unlocked {
860                    crypt,
861                    cached_keys: Vec::new(),
862                    mutations_cipher: JournalCipher::new_aes256_xts(&unwrapped_key, 0),
863                },
864                last_object_id_in_memory,
865            )
866        } else {
867            Self::new(
868                Some(self.clone()),
869                handle.object_id(),
870                filesystem.clone(),
871                Some(StoreInfo {
872                    last_object_id,
873                    guid: options.guid.unwrap_or_else(|| *Uuid::new_v4().as_bytes()),
874                    ..Default::default()
875                }),
876                object_cache,
877                LockState::Unencrypted,
878                last_object_id_in_memory,
879            )
880        };
881        assert!(store.store_info_handle.set(handle).is_ok());
882        Ok(store)
883    }
884
885    /// Actually creates the store in a transaction.  This will also create a root directory and
886    /// graveyard directory for the store.  See `new_child_store` above.
887    async fn create<'a>(
888        self: &'a Arc<Self>,
889        transaction: &mut Transaction<'a>,
890    ) -> Result<(), Error> {
891        let buf = {
892            // Create a root directory and graveyard directory.
893            let graveyard_directory_object_id = Graveyard::create(transaction, &self).await?;
894            let root_directory = Directory::create(transaction, &self, None).await?;
895
896            let serialized_info = {
897                let mut store_info = self.store_info.lock();
898                let store_info = store_info.as_mut().unwrap();
899
900                store_info.graveyard_directory_object_id = graveyard_directory_object_id;
901                store_info.root_directory_object_id = root_directory.object_id();
902
903                let mut serialized_info = Vec::new();
904                store_info.serialize_with_version(&mut serialized_info)?;
905                serialized_info
906            };
907            let mut buf = self.device.allocate_buffer(serialized_info.len()).await;
908            buf.copy_from_slice(&serialized_info[..]);
909            buf
910        };
911
912        if self.filesystem().options().image_builder_mode.is_some() {
913            // If we're in image builder mode, we want to avoid writing to disk unless explicitly
914            // asked to. New object stores will have their StoreInfo written when we compact in
915            // FxFilesystem::finalize().
916            Ok(())
917        } else {
918            self.store_info_handle.get().unwrap().txn_write(transaction, 0u64, buf.as_ref()).await
919        }
920    }
921
922    pub fn set_trace(&self, trace: bool) {
923        let old_value = self.trace.swap(trace, Ordering::Relaxed);
924        if trace != old_value {
925            info!(store_id = self.store_object_id(), trace; "OS: trace",);
926        }
927    }
928
929    /// Sets a callback to be invoked each time the ObjectStore flushes.  The callback is invoked at
930    /// the end of flush, while the write lock is still held.
931    pub fn set_flush_callback<F: Fn(&ObjectStore) + Send + Sync + 'static>(&self, callback: F) {
932        let mut flush_callback = self.flush_callback.lock();
933        *flush_callback = Some(Box::new(callback));
934    }
935
936    pub fn is_root(&self) -> bool {
937        if let Some(parent) = &self.parent_store {
938            parent.parent_store.is_none()
939        } else {
940            // The root parent store isn't the root store.
941            false
942        }
943    }
944
945    /// Populates an inspect node with store statistics.
946    pub fn record_data(self: &Arc<Self>, root: &fuchsia_inspect::Node) {
947        // TODO(https://fxbug.dev/42069513): Push-back or rate-limit to prevent DoS.
948        let counters = self.counters.lock();
949        if let Some(store_info) = self.store_info() {
950            root.record_string("guid", Uuid::from_bytes(store_info.guid).to_string());
951        };
952        root.record_uint("store_object_id", self.store_object_id);
953        root.record_uint("mutations_applied", counters.mutations_applied);
954        root.record_uint("mutations_dropped", counters.mutations_dropped);
955        root.record_uint("num_flushes", counters.num_flushes);
956        if let Some(last_flush_time) = counters.last_flush_time.as_ref() {
957            root.record_uint(
958                "last_flush_time_ms",
959                last_flush_time
960                    .duration_since(std::time::UNIX_EPOCH)
961                    .unwrap_or(std::time::Duration::ZERO)
962                    .as_millis()
963                    .try_into()
964                    .unwrap_or(0u64),
965            );
966        }
967        root.record_uint("device_read_ops", self.device_read_ops.load(Ordering::Relaxed));
968        root.record_uint("device_write_ops", self.device_write_ops.load(Ordering::Relaxed));
969        root.record_uint("logical_read_ops", self.logical_read_ops.load(Ordering::Relaxed));
970        root.record_uint("logical_write_ops", self.logical_write_ops.load(Ordering::Relaxed));
971        root.record_uint("graveyard_entries", self.graveyard_entries.load(Ordering::Relaxed));
972        {
973            let last_object_id = self.last_object_id.lock();
974            root.record_uint("object_id_hi", last_object_id.id() >> 32);
975            root.record_bool(
976                "low_32_bit_object_ids",
977                matches!(&*last_object_id, LastObjectId::Low32Bit { .. }),
978            );
979        }
980
981        let this = self.clone();
982        root.record_child("lsm_tree", move |node| this.tree().record_inspect_data(node));
983    }
984
985    pub fn device(&self) -> &Arc<dyn Device> {
986        &self.device
987    }
988
989    pub fn block_size(&self) -> BlockSize {
990        self.block_size
991    }
992
993    pub fn filesystem(&self) -> Arc<FxFilesystem> {
994        self.filesystem.upgrade().unwrap()
995    }
996
997    pub fn store_object_id(&self) -> u64 {
998        self.store_object_id
999    }
1000
1001    pub fn tree(&self) -> &LSMTree<ObjectKey, ObjectValue> {
1002        &self.tree
1003    }
1004
1005    pub fn root_directory_object_id(&self) -> u64 {
1006        self.store_info.lock().as_ref().unwrap().root_directory_object_id
1007    }
1008
1009    pub fn guid(&self) -> [u8; 16] {
1010        self.store_info.lock().as_ref().unwrap().guid
1011    }
1012
1013    pub fn graveyard_directory_object_id(&self) -> u64 {
1014        self.store_info.lock().as_ref().unwrap().graveyard_directory_object_id
1015    }
1016
1017    fn set_graveyard_directory_object_id(&self, oid: u64) {
1018        assert_eq!(
1019            std::mem::replace(
1020                &mut self.store_info.lock().as_mut().unwrap().graveyard_directory_object_id,
1021                oid
1022            ),
1023            INVALID_OBJECT_ID
1024        );
1025    }
1026
1027    pub fn object_count(&self) -> u64 {
1028        self.store_info.lock().as_ref().unwrap().object_count
1029    }
1030
1031    /// Returns INVALID_OBJECT_ID for algorithms that don't use the last ID.
1032    pub(crate) fn unencrypted_last_object_id(&self) -> u64 {
1033        self.last_object_id.lock().id()
1034    }
1035
1036    pub fn key_manager(&self) -> &KeyManager {
1037        &self.key_manager
1038    }
1039
1040    /// Clears all in-memory caches (LSM tree object cache, persistent layer chunk caches, and
1041    /// non-permanent unwrapped keys) for this object store.
1042    pub fn clear_caches(&self) {
1043        self.tree.clear_cache();
1044        self.key_manager.clear_cached_keys();
1045    }
1046
1047    pub fn parent_store(&self) -> Option<&Arc<ObjectStore>> {
1048        self.parent_store.as_ref()
1049    }
1050
1051    /// Returns the crypt object for the store.  Returns None if the store is unencrypted.
1052    pub fn crypt(&self) -> Option<Arc<dyn Crypt>> {
1053        match &*self.lock_state.lock() {
1054            LockState::Locked => panic!("Store is locked"),
1055            LockState::Invalid
1056            | LockState::Unencrypted
1057            | LockState::Locking
1058            | LockState::Unlocking
1059            | LockState::Deleted => None,
1060            LockState::Unlocked { crypt, .. } => Some(crypt.clone()),
1061            LockState::UnlockedReadOnly(crypt) => Some(crypt.clone()),
1062            LockState::Unknown => {
1063                panic!("Store is of unknown lock state; has the journal been replayed yet?")
1064            }
1065        }
1066    }
1067
1068    /// Returns the id of the internal directory. Returns a NotFound error if this has not been
1069    /// initialized.
1070    pub fn get_internal_directory_id(self: &Arc<Self>) -> Result<u64, Error> {
1071        if let Some(store_info) = self.store_info.lock().as_ref() {
1072            if store_info.internal_directory_object_id == INVALID_OBJECT_ID {
1073                Err(FxfsError::NotFound.into())
1074            } else {
1075                Ok(store_info.internal_directory_object_id)
1076            }
1077        } else {
1078            Err(FxfsError::Unavailable.into())
1079        }
1080    }
1081
1082    pub async fn get_or_create_internal_directory_id(self: &Arc<Self>) -> Result<u64, Error> {
1083        // Create the transaction first to use the object store lock.
1084        let mut transaction = self
1085            .new_transaction(
1086                lock_keys![LockKey::InternalDirectory { store_object_id: self.store_object_id }],
1087                Options::default(),
1088            )
1089            .await?;
1090        let obj_id = self.store_info.lock().as_ref().unwrap().internal_directory_object_id;
1091        if obj_id != INVALID_OBJECT_ID {
1092            return Ok(obj_id);
1093        }
1094
1095        // Need to create an internal directory.
1096        let directory = Directory::create(&mut transaction, self, None).await?;
1097
1098        transaction.add(self.store_object_id, Mutation::CreateInternalDir(directory.object_id()));
1099        transaction.commit().await?;
1100        Ok(directory.object_id())
1101    }
1102
1103    /// Returns the file size for the object without opening the object.
1104    async fn get_file_size(&self, object_id: u64) -> Result<u64, Error> {
1105        self.tree
1106            .find_map(
1107                &ObjectKey::attribute(object_id, AttributeId::DATA, AttributeKey::Attribute),
1108                |item| match item.value {
1109                    ObjectValue::Attribute { size, .. } => Ok(*size),
1110                    _ => Err(anyhow!(FxfsError::NotFile)),
1111                },
1112            )
1113            .await?
1114            .ok_or(FxfsError::NotFound)?
1115    }
1116
1117    #[cfg(feature = "migration")]
1118    pub fn last_object_id(&self) -> u64 {
1119        self.last_object_id.lock().id()
1120    }
1121
1122    /// Provides access to the allocator to mark a specific region of the device as allocated.
1123    #[cfg(feature = "migration")]
1124    pub fn mark_allocated(
1125        &self,
1126        transaction: &mut Transaction<'_>,
1127        store_object_id: u64,
1128        device_range: std::ops::Range<u64>,
1129    ) -> Result<(), Error> {
1130        self.allocator().mark_allocated(transaction, store_object_id, device_range)
1131    }
1132
1133    /// `crypt` can be provided if the crypt service should be different to the default; see the
1134    /// comment on create_object.  Users should avoid having more than one handle open for the same
1135    /// object at the same time because they might get out-of-sync; there is no code that will
1136    /// prevent this.  One example where this can cause an issue is if the object ends up using a
1137    /// permanent key (which is the case if a value is passed for `crypt`), the permanent key is
1138    /// dropped when a handle is dropped, which will impact any other handles for the same object.
1139    pub async fn open_object<S: HandleOwner>(
1140        owner: &Arc<S>,
1141        obj_id: u64,
1142        options: HandleOptions,
1143        crypt: Option<Arc<dyn Crypt>>,
1144    ) -> Result<DataObjectHandle<S>, Error> {
1145        let store = owner.as_ref().as_ref();
1146        let mut fsverity_descriptor = None;
1147        let mut overwrite_ranges = Vec::new();
1148        let value = store
1149            .tree
1150            .find_value(&ObjectKey::attribute(obj_id, AttributeId::DATA, AttributeKey::Attribute))
1151            .await?
1152            .ok_or(FxfsError::NotFound)?;
1153
1154        let (size, track_overwrite_extents) = match value {
1155            ObjectValue::Attribute { size, has_overwrite_extents } => (size, has_overwrite_extents),
1156            ObjectValue::VerifiedAttribute { size, fsverity_metadata } => {
1157                if !options.skip_fsverity {
1158                    fsverity_descriptor = Some(fsverity_metadata);
1159                }
1160                // We only track the overwrite extents in memory for writes, reads handle them
1161                // implicitly, which means verified files (where the data won't change anymore)
1162                // don't need to track them.
1163                (size, false)
1164            }
1165            _ => bail!(anyhow!(FxfsError::Inconsistent).context("open_object: Expected attibute")),
1166        };
1167
1168        ensure!(size <= MAX_FILE_SIZE, FxfsError::Inconsistent);
1169
1170        if track_overwrite_extents {
1171            let layer_set = store.tree.layer_set();
1172            let mut merger = layer_set.merger();
1173            let mut iter = merger
1174                .query(Query::FullRange(&ObjectKey::attribute(
1175                    obj_id,
1176                    AttributeId::DATA,
1177                    AttributeKey::Extent(Extent::search_key_from_offset(0)),
1178                )))
1179                .await?;
1180            loop {
1181                match iter.get() {
1182                    Some(ItemRef {
1183                        key:
1184                            ObjectKey {
1185                                object_id,
1186                                data:
1187                                    ObjectKeyData::Attribute(
1188                                        AttributeId::DATA,
1189                                        AttributeKey::Extent(extent),
1190                                    ),
1191                            },
1192                        value,
1193                        ..
1194                    }) if *object_id == obj_id => {
1195                        match value {
1196                            ObjectValue::Extent(ExtentValue::None)
1197                            | ObjectValue::Extent(ExtentValue::Some {
1198                                mode: ExtentMode::Raw,
1199                                ..
1200                            })
1201                            | ObjectValue::Extent(ExtentValue::Some {
1202                                mode: ExtentMode::Cow(_),
1203                                ..
1204                            }) => (),
1205                            ObjectValue::Extent(ExtentValue::Some {
1206                                mode: ExtentMode::OverwritePartial(_),
1207                                ..
1208                            })
1209                            | ObjectValue::Extent(ExtentValue::Some {
1210                                mode: ExtentMode::Overwrite,
1211                                ..
1212                            }) => overwrite_ranges.push(extent.clone().into()),
1213                            _ => bail!(
1214                                anyhow!(FxfsError::Inconsistent)
1215                                    .context("open_object: Expected extent")
1216                            ),
1217                        }
1218                        iter.advance().await?;
1219                    }
1220                    _ => break,
1221                }
1222            }
1223        }
1224
1225        // If a crypt service has been specified, it needs to be a permanent key because cached
1226        // keys can only use the store's crypt service.
1227        let permanent = if let Some(crypt) = crypt {
1228            store
1229                .key_manager
1230                .get_keys(
1231                    obj_id,
1232                    crypt.as_ref(),
1233                    &mut Some(async || store.get_keys(obj_id).await),
1234                    /* permanent= */ true,
1235                    /* force= */ false,
1236                )
1237                .await?;
1238            true
1239        } else {
1240            false
1241        };
1242        let data_object_handle = DataObjectHandle::new(
1243            owner.clone(),
1244            obj_id,
1245            permanent,
1246            AttributeId::DATA,
1247            size,
1248            options,
1249            false,
1250            &overwrite_ranges,
1251        );
1252        if let Some(descriptor) = fsverity_descriptor {
1253            data_object_handle
1254                .set_fsverity_state_some(descriptor)
1255                .await
1256                .context("Invalid or mismatched merkle tree")?;
1257        }
1258        Ok(data_object_handle)
1259    }
1260
1261    pub fn create_object_with_id<S: HandleOwner>(
1262        owner: &Arc<S>,
1263        transaction: &mut Transaction<'_>,
1264        reserved_object_id: ReservedId<'_>,
1265        options: HandleOptions,
1266        encryption_options: Option<ObjectEncryptionOptions>,
1267    ) -> Result<DataObjectHandle<S>, Error> {
1268        let store = owner.as_ref().as_ref();
1269        // Don't permit creating unencrypted objects in an encrypted store.  The converse is OK.
1270        debug_assert!(store.crypt().is_none() || encryption_options.is_some());
1271        let now = Timestamp::now();
1272        let object_id = reserved_object_id.get();
1273        assert!(
1274            transaction
1275                .add(
1276                    store.store_object_id(),
1277                    Mutation::insert_object(
1278                        ObjectKey::object(reserved_object_id.release().get()),
1279                        ObjectValue::file(
1280                            1,
1281                            0,
1282                            now.clone(),
1283                            now.clone(),
1284                            now.clone(),
1285                            now,
1286                            None,
1287                            None
1288                        ),
1289                    ),
1290                )
1291                .is_none()
1292        );
1293        let mut permanent_keys = false;
1294        if let Some(ObjectEncryptionOptions { permanent, key_id, key, unwrapped_key }) =
1295            encryption_options
1296        {
1297            permanent_keys = permanent;
1298            let cipher = key_to_cipher(&key, &unwrapped_key)?;
1299            store.key_manager.insert_fscrypt_file_cipher(&key, cipher.clone());
1300            transaction.add(
1301                store.store_object_id(),
1302                Mutation::insert_object(
1303                    ObjectKey::keys(object_id),
1304                    ObjectValue::keys(vec![(key_id, key)].into()),
1305                ),
1306            );
1307            store.key_manager.insert(
1308                object_id,
1309                Arc::new(vec![(key_id, CipherHolder::Cipher(cipher))].into()),
1310                permanent,
1311            );
1312        }
1313        transaction.add(
1314            store.store_object_id(),
1315            Mutation::insert_object(
1316                ObjectKey::attribute(object_id, AttributeId::DATA, AttributeKey::Attribute),
1317                // This is a new object so nothing has pre-allocated overwrite extents yet.
1318                ObjectValue::attribute(0, false),
1319            ),
1320        );
1321        Ok(DataObjectHandle::new(
1322            owner.clone(),
1323            object_id,
1324            permanent_keys,
1325            AttributeId::DATA,
1326            0,
1327            options,
1328            false,
1329            &[],
1330        ))
1331    }
1332
1333    /// Creates an object in the store.
1334    ///
1335    /// If the store is encrypted, the object will be automatically encrypted as well.
1336    /// If `wrapping_key_id` is set, the new keys will be wrapped with that specific key, and
1337    /// otherwise the default data key is used.
1338    pub async fn create_object<S: HandleOwner>(
1339        owner: &Arc<S>,
1340        mut transaction: &mut Transaction<'_>,
1341        options: HandleOptions,
1342        fscrypt_info: Option<FscryptDirInfo>,
1343    ) -> Result<DataObjectHandle<S>, Error> {
1344        let store = owner.as_ref().as_ref();
1345        let object_id = store.get_next_object_id(transaction).await?;
1346        let crypt = store.crypt();
1347        let encryption_options = if let Some(crypt) = crypt {
1348            let key_id = if fscrypt_info.is_some() { FSCRYPT_KEY_ID } else { VOLUME_DATA_KEY_ID };
1349            let (key, unwrapped_key) = if let Some(info) = fscrypt_info {
1350                crypt
1351                    .create_key_with_id(
1352                        object_id.get(),
1353                        info.wrapping_key_id,
1354                        ObjectType::File,
1355                        info.flags.into(),
1356                    )
1357                    .await?
1358            } else {
1359                let (fxfs_key, unwrapped_key) =
1360                    crypt.create_key(object_id.get(), KeyPurpose::Data).await?;
1361                (EncryptionKey::Fxfs(fxfs_key), unwrapped_key)
1362            };
1363            Some(ObjectEncryptionOptions { permanent: false, key_id, key, unwrapped_key })
1364        } else {
1365            None
1366        };
1367        ObjectStore::create_object_with_id(
1368            owner,
1369            &mut transaction,
1370            object_id,
1371            options,
1372            encryption_options,
1373        )
1374    }
1375
1376    /// Creates an object using explicitly provided keys.
1377    ///
1378    /// There are some cases where an encrypted object needs to be created in an unencrypted store.
1379    /// For example, when layer files for a child store are created in the root store, but they must
1380    /// be encrypted using the child store's keys.  This method exists for that purpose.
1381    pub(crate) async fn create_object_with_key<S: HandleOwner>(
1382        owner: &Arc<S>,
1383        mut transaction: &mut Transaction<'_>,
1384        object_id: ReservedId<'_>,
1385        options: HandleOptions,
1386        key: EncryptionKey,
1387        unwrapped_key: UnwrappedKey,
1388    ) -> Result<DataObjectHandle<S>, Error> {
1389        ObjectStore::create_object_with_id(
1390            owner,
1391            &mut transaction,
1392            object_id,
1393            options,
1394            Some(ObjectEncryptionOptions {
1395                permanent: true,
1396                key_id: VOLUME_DATA_KEY_ID,
1397                key,
1398                unwrapped_key,
1399            }),
1400        )
1401    }
1402
1403    /// Adjusts the reference count for a given object.  If the reference count reaches zero, the
1404    /// object is moved into the graveyard and true is returned.
1405    pub async fn adjust_refs(
1406        &self,
1407        transaction: &mut Transaction<'_>,
1408        object_id: u64,
1409        delta: i64,
1410    ) -> Result<bool, Error> {
1411        self.adjust_refs_impl(transaction, object_id, delta, true).await
1412    }
1413
1414    /// Adjusts the reference count for a given object.  If the reference count reaches zero and
1415    /// `add_to_graveyard` is true, the object is moved into the graveyard.  Returns whether the
1416    /// reference count reached zero.
1417    pub async fn adjust_refs_impl(
1418        &self,
1419        transaction: &mut Transaction<'_>,
1420        object_id: u64,
1421        delta: i64,
1422        add_to_graveyard: bool,
1423    ) -> Result<bool, Error> {
1424        let mut mutation = self.txn_get_object_mutation(transaction, object_id).await?;
1425        let refs = if let ObjectValue::Object {
1426            kind:
1427                ObjectKind::File { refs, .. }
1428                | ObjectKind::Symlink { refs, .. }
1429                | ObjectKind::EncryptedSymlink { refs, .. },
1430            ..
1431        } = &mut mutation.item.value
1432        {
1433            *refs =
1434                refs.checked_add_signed(delta).ok_or_else(|| anyhow!("refs underflow/overflow"))?;
1435            refs
1436        } else {
1437            bail!(FxfsError::NotFile);
1438        };
1439        if *refs == 0 {
1440            if add_to_graveyard {
1441                self.add_to_graveyard(transaction, object_id);
1442
1443                // We might still need to adjust the reference count if delta was something other
1444                // than -1.
1445                if delta != -1 {
1446                    *refs = 1;
1447                    transaction.add(self.store_object_id, Mutation::ObjectStore(mutation));
1448                }
1449
1450                // Otherwise, we don't commit the mutation as we want to keep reference count as 1
1451                // for objects in graveyard.
1452            }
1453            Ok(true)
1454        } else {
1455            transaction.add(self.store_object_id, Mutation::ObjectStore(mutation));
1456            Ok(false)
1457        }
1458    }
1459
1460    /// Tombstones (purges) an object that is in the graveyard by deallocating its extents,
1461    /// removing its graveyard entry, and inserting an LSM-tree tombstone record
1462    /// (`ObjectValue::None`). If a truncate guard is not provided, one will be acquired.
1463    ///
1464    /// Callers outside of [`ObjectStore`] should generally prefer
1465    /// [`FxFilesystem::tombstone_object`] or [`Graveyard::queue_tombstone_object`] so that
1466    /// transaction reservations are configured appropriately for the store.
1467    pub async fn tombstone_object(
1468        &self,
1469        object_id: u64,
1470        txn_options: Options<'_>,
1471        truncate_guard: Option<&TruncateGuard<'_>>,
1472    ) -> Result<(), Error> {
1473        debug_assert!(
1474            self.tree.exists(&ObjectKey::object(object_id)).await?,
1475            "Tombstoning missing object"
1476        );
1477        debug_assert!(
1478            self.tree
1479                .exists(&ObjectKey::graveyard_entry(
1480                    self.graveyard_directory_object_id(),
1481                    object_id
1482                ))
1483                .await?,
1484            "Tombstoning object not in graveyard"
1485        );
1486        self.key_manager.remove(object_id).await;
1487        let _fs;
1488        let _guard;
1489        let truncate_guard = match truncate_guard {
1490            Some(guard) => guard,
1491            None => {
1492                _fs = self.filesystem();
1493                _guard = _fs.truncate_guard(self.store_object_id, object_id).await;
1494                &_guard
1495            }
1496        };
1497        self.trim_or_tombstone(object_id, true, txn_options, truncate_guard).await
1498    }
1499
1500    /// Trim extents beyond the end of a file for all attributes.  This will remove the entry from
1501    /// the graveyard when done.
1502    pub async fn trim(
1503        &self,
1504        object_id: u64,
1505        truncate_guard: &TruncateGuard<'_>,
1506    ) -> Result<(), Error> {
1507        // For the root and root parent store, we would need to use the metadata reservation which
1508        // we don't currently support, so assert that we're not those stores.
1509        assert!(self.parent_store.as_ref().unwrap().parent_store.is_some());
1510
1511        self.trim_or_tombstone(
1512            object_id,
1513            false,
1514            Options { reservation: ReservationOptions::BorrowedMetadata, ..Default::default() },
1515            truncate_guard,
1516        )
1517        .await
1518    }
1519
1520    /// Trims or tombstones an object.
1521    async fn trim_or_tombstone(
1522        &self,
1523        object_id: u64,
1524        for_tombstone: bool,
1525        txn_options: Options<'_>,
1526        _truncate_guard: &TruncateGuard<'_>,
1527    ) -> Result<(), Error> {
1528        let mut next_attribute = Some(AttributeId::SORTED_START);
1529        while let Some(attribute_id) = next_attribute.take() {
1530            let mut transaction = self
1531                .new_transaction(
1532                    lock_keys![
1533                        LockKey::object_attribute(self.store_object_id, object_id, attribute_id),
1534                        LockKey::object(self.store_object_id, object_id),
1535                    ],
1536                    txn_options,
1537                )
1538                .await?;
1539
1540            match self
1541                .trim_some(
1542                    &mut transaction,
1543                    object_id,
1544                    attribute_id,
1545                    if for_tombstone {
1546                        TrimMode::Tombstone(TombstoneMode::Object)
1547                    } else {
1548                        TrimMode::UseSize
1549                    },
1550                )
1551                .await?
1552            {
1553                TrimResult::Incomplete => next_attribute = Some(attribute_id),
1554                TrimResult::Done(None) => {
1555                    if for_tombstone
1556                        || matches!(
1557                            self.tree
1558                                .find_value(&ObjectKey::graveyard_entry(
1559                                    self.graveyard_directory_object_id(),
1560                                    object_id,
1561                                ))
1562                                .await?,
1563                            Some(ObjectValue::Trim)
1564                        )
1565                    {
1566                        self.remove_from_graveyard(&mut transaction, object_id);
1567                    }
1568                    // The last attribute was not the default attribute, it may have been added to
1569                    // the graveyard alongside the object.
1570                    if for_tombstone && attribute_id != AttributeId::DATA {
1571                        self.remove_attribute_from_graveyard(
1572                            &mut transaction,
1573                            object_id,
1574                            attribute_id,
1575                        );
1576                    }
1577                }
1578                TrimResult::Done(id) => {
1579                    // Moved to the next attribute. This one is finished and it may have been
1580                    // added to the graveyard alongside the object.
1581                    if for_tombstone && attribute_id != AttributeId::DATA {
1582                        self.remove_attribute_from_graveyard(
1583                            &mut transaction,
1584                            object_id,
1585                            attribute_id,
1586                        );
1587                    }
1588                    next_attribute = id;
1589                }
1590            }
1591
1592            if !transaction.mutations().is_empty() {
1593                transaction.commit().await?;
1594            }
1595        }
1596        Ok(())
1597    }
1598
1599    /// Tombstones (purges) an object's attribute that is in the graveyard by deallocating its
1600    /// extents, removing its graveyard entry, and inserting an LSM-tree tombstone record
1601    /// (`ObjectValue::None`).
1602    ///
1603    /// Callers outside of [`ObjectStore`] should generally prefer
1604    /// [`FxFilesystem::tombstone_attribute`] or [`Graveyard::queue_tombstone_attribute`] so that
1605    /// transaction reservations are configured appropriately for the store.
1606    pub async fn tombstone_attribute(
1607        &self,
1608        object_id: u64,
1609        attribute_id: AttributeId,
1610        txn_options: Options<'_>,
1611    ) -> Result<(), Error> {
1612        // Ensure that we don't double-delete things, it should still exist and be in the graveyard.
1613        debug_assert!(
1614            self.tree
1615                .exists(&ObjectKey::attribute(object_id, attribute_id, AttributeKey::Attribute))
1616                .await?,
1617            "Tombstoning missing attribute"
1618        );
1619        debug_assert!(
1620            self.tree
1621                .exists(&ObjectKey::graveyard_attribute_entry(
1622                    self.graveyard_directory_object_id(),
1623                    object_id,
1624                    attribute_id
1625                ))
1626                .await?,
1627            "Tombstoning attribute not in graveyard"
1628        );
1629        let mut trim_result = TrimResult::Incomplete;
1630        while matches!(trim_result, TrimResult::Incomplete) {
1631            let mut transaction = self
1632                .new_transaction(
1633                    lock_keys![
1634                        LockKey::object_attribute(self.store_object_id, object_id, attribute_id),
1635                        LockKey::object(self.store_object_id, object_id),
1636                    ],
1637                    txn_options,
1638                )
1639                .await?;
1640            trim_result = self
1641                .trim_some(
1642                    &mut transaction,
1643                    object_id,
1644                    attribute_id,
1645                    TrimMode::Tombstone(TombstoneMode::Attribute),
1646                )
1647                .await?;
1648            if let TrimResult::Done(..) = trim_result {
1649                self.remove_attribute_from_graveyard(&mut transaction, object_id, attribute_id)
1650            }
1651            if !transaction.mutations().is_empty() {
1652                transaction.commit().await?;
1653            }
1654        }
1655        Ok(())
1656    }
1657
1658    /// Deletes extents for attribute `attribute_id` in object `object_id`.  Also see the comments
1659    /// for TrimMode and TrimResult. Should hold a lock on the attribute, and the object as it
1660    /// performs a read-modify-write on the sizes.
1661    pub async fn trim_some(
1662        &self,
1663        transaction: &mut Transaction<'_>,
1664        object_id: u64,
1665        attribute_id: AttributeId,
1666        mode: TrimMode,
1667    ) -> Result<TrimResult, Error> {
1668        let layer_set = self.tree.layer_set();
1669        let mut merger = layer_set.merger();
1670
1671        let aligned_offset = match mode {
1672            TrimMode::FromOffset(offset) => {
1673                self.block_size.align_up(offset).ok_or(FxfsError::Inconsistent)?
1674            }
1675            TrimMode::Tombstone(..) => 0,
1676            TrimMode::UseSize => {
1677                let iter = merger
1678                    .query(Query::FullRange(&ObjectKey::attribute(
1679                        object_id,
1680                        attribute_id,
1681                        AttributeKey::Attribute,
1682                    )))
1683                    .await?;
1684                if let Some(item_ref) = iter.get() {
1685                    if item_ref.key.object_id != object_id {
1686                        return Ok(TrimResult::Done(None));
1687                    }
1688
1689                    if let ItemRef {
1690                        key:
1691                            ObjectKey {
1692                                data:
1693                                    ObjectKeyData::Attribute(size_attribute_id, AttributeKey::Attribute),
1694                                ..
1695                            },
1696                        value: ObjectValue::Attribute { size, .. },
1697                        ..
1698                    } = item_ref
1699                    {
1700                        // If we found a different attribute_id, return so we can get the
1701                        // right lock.
1702                        if *size_attribute_id != attribute_id {
1703                            return Ok(TrimResult::Done(Some(*size_attribute_id)));
1704                        }
1705                        self.block_size.align_up(*size).ok_or(FxfsError::Inconsistent)?
1706                    } else {
1707                        // At time of writing, we should always see a size record or None here, but
1708                        // asserting here would be brittle so just skip to the the next attribute
1709                        // instead.
1710                        return Ok(TrimResult::Done(Some(attribute_id.next())));
1711                    }
1712                } else {
1713                    // End of the tree.
1714                    return Ok(TrimResult::Done(None));
1715                }
1716            }
1717        };
1718
1719        // Loop over the extents and deallocate them.
1720        let mut iter = merger
1721            .query(Query::FullRange(&ObjectKey::from_extent(
1722                object_id,
1723                attribute_id,
1724                Extent::search_key_from_offset(aligned_offset),
1725            )))
1726            .await?;
1727        let mut end = 0;
1728        let allocator = self.allocator();
1729        let mut result = TrimResult::Done(None);
1730        let mut deallocated = 0;
1731        let block_size = self.block_size;
1732
1733        while let Some(item_ref) = iter.get() {
1734            if item_ref.key.object_id != object_id {
1735                break;
1736            }
1737            if let ObjectKey {
1738                data: ObjectKeyData::Attribute(extent_attribute_id, attribute_key),
1739                ..
1740            } = item_ref.key
1741            {
1742                if *extent_attribute_id != attribute_id {
1743                    result = TrimResult::Done(Some(*extent_attribute_id));
1744                    break;
1745                }
1746                if let (
1747                    AttributeKey::Extent(extent),
1748                    ObjectValue::Extent(ExtentValue::Some { device_offset, .. }),
1749                ) = (attribute_key, item_ref.value)
1750                {
1751                    let start = std::cmp::max(extent.start, aligned_offset);
1752                    ensure!(start < extent.end, FxfsError::Inconsistent);
1753                    let device_offset = device_offset
1754                        .checked_add(start - extent.start)
1755                        .ok_or(FxfsError::Inconsistent)?;
1756                    end = extent.end;
1757                    let len = end - start;
1758                    let device_range = device_offset..device_offset + len;
1759                    ensure!(block_size.is_aligned(&device_range), FxfsError::Inconsistent);
1760                    allocator.deallocate(transaction, self.store_object_id, device_range).await?;
1761                    deallocated += len;
1762                    // Stop if the transaction is getting too big.
1763                    if transaction.mutations().len() >= TRANSACTION_MUTATION_THRESHOLD {
1764                        result = TrimResult::Incomplete;
1765                        break;
1766                    }
1767                }
1768            }
1769            iter.advance().await?;
1770        }
1771
1772        let finished_tombstone_object = matches!(mode, TrimMode::Tombstone(TombstoneMode::Object))
1773            && matches!(result, TrimResult::Done(None));
1774        let finished_tombstone_attribute =
1775            matches!(mode, TrimMode::Tombstone(TombstoneMode::Attribute))
1776                && !matches!(result, TrimResult::Incomplete);
1777        let mut object_mutation = None;
1778        let nodes = if finished_tombstone_object { -1 } else { 0 };
1779        if nodes != 0 || deallocated != 0 {
1780            let mutation = self.txn_get_object_mutation(transaction, object_id).await?;
1781            if let ObjectValue::Object { attributes: ObjectAttributes { project_id, .. }, .. } =
1782                mutation.item.value
1783            {
1784                if let Some(project_id) = project_id {
1785                    transaction.merge_bytes_and_nodes(
1786                        self.store_object_id,
1787                        ObjectKey::project_usage(self.root_directory_object_id(), project_id),
1788                        BytesAndNodes { bytes: -i64::try_from(deallocated).unwrap(), nodes },
1789                    );
1790                }
1791                object_mutation = Some(mutation);
1792            } else {
1793                panic!("Inconsistent object type.");
1794            }
1795        }
1796
1797        // Deletion marker records *must* be merged so as to consume all other records for the
1798        // object.
1799        if finished_tombstone_object {
1800            transaction.add(
1801                self.store_object_id,
1802                Mutation::merge_object(ObjectKey::object(object_id), ObjectValue::None),
1803            );
1804        } else {
1805            if finished_tombstone_attribute {
1806                transaction.add(
1807                    self.store_object_id,
1808                    Mutation::merge_object(
1809                        ObjectKey::attribute(object_id, attribute_id, AttributeKey::Attribute),
1810                        ObjectValue::None,
1811                    ),
1812                );
1813            }
1814            if deallocated > 0 {
1815                let mut mutation = match object_mutation {
1816                    Some(mutation) => mutation,
1817                    None => self.txn_get_object_mutation(transaction, object_id).await?,
1818                };
1819                transaction.add(
1820                    self.store_object_id,
1821                    Mutation::merge_object(
1822                        ObjectKey::extent(object_id, attribute_id, aligned_offset..end),
1823                        ObjectValue::deleted_extent(),
1824                    ),
1825                );
1826                // Update allocated size.
1827                if let ObjectValue::Object {
1828                    attributes: ObjectAttributes { allocated_size, .. },
1829                    ..
1830                } = &mut mutation.item.value
1831                {
1832                    // The only way for these to fail are if the volume is inconsistent.
1833                    *allocated_size = allocated_size.checked_sub(deallocated).ok_or_else(|| {
1834                        anyhow!(FxfsError::Inconsistent).context("Allocated size overflow")
1835                    })?;
1836                } else {
1837                    panic!("Unexpected object value");
1838                }
1839                transaction.add(self.store_object_id, Mutation::ObjectStore(mutation));
1840            }
1841        }
1842        Ok(result)
1843    }
1844
1845    /// Returns all objects that exist in the parent store that pertain to this object store.
1846    /// Note that this doesn't include the object_id of the store itself which is generally
1847    /// referenced externally.
1848    pub fn parent_objects(&self) -> Vec<u64> {
1849        assert!(self.store_info_handle.get().is_some());
1850        self.store_info.lock().as_ref().unwrap().parent_objects()
1851    }
1852
1853    /// Returns root objects for this store.
1854    pub fn root_objects(&self) -> Vec<u64> {
1855        let mut objects = Vec::new();
1856        let store_info = self.store_info.lock();
1857        let info = store_info.as_ref().unwrap();
1858        if info.root_directory_object_id != INVALID_OBJECT_ID {
1859            objects.push(info.root_directory_object_id);
1860        }
1861        if info.graveyard_directory_object_id != INVALID_OBJECT_ID {
1862            objects.push(info.graveyard_directory_object_id);
1863        }
1864        if info.internal_directory_object_id != INVALID_OBJECT_ID {
1865            objects.push(info.internal_directory_object_id);
1866        }
1867        objects
1868    }
1869
1870    pub fn store_info(&self) -> Option<StoreInfo> {
1871        self.store_info.lock().as_ref().cloned()
1872    }
1873
1874    /// Returns None if called during journal replay.
1875    pub fn store_info_handle_object_id(&self) -> Option<u64> {
1876        self.store_info_handle.get().map(|h| h.object_id())
1877    }
1878
1879    pub fn graveyard_count(&self) -> u64 {
1880        self.graveyard_entries.load(Ordering::Relaxed)
1881    }
1882
1883    /// Called to open a store, before replay of this store's mutations.
1884    async fn open(
1885        parent_store: &Arc<ObjectStore>,
1886        store_object_id: u64,
1887        object_cache: Option<Box<dyn ObjectCache<ObjectKey, ObjectValue>>>,
1888    ) -> Result<Arc<ObjectStore>, Error> {
1889        let handle =
1890            ObjectStore::open_object(parent_store, store_object_id, HandleOptions::default(), None)
1891                .await?;
1892
1893        let info = load_store_info(parent_store, store_object_id).await?;
1894        let is_encrypted = info.mutations_key.is_some();
1895
1896        let mut total_layer_size = 0;
1897        let last_object_id;
1898
1899        // TODO(https://fxbug.dev/42178043): the layer size here could be bad and cause overflow.
1900
1901        // If the store is encrypted, we can't open the object tree layers now, but we need to
1902        // compute the size of the layers.
1903        if is_encrypted {
1904            for &oid in &info.layers {
1905                total_layer_size += parent_store.get_file_size(oid).await?;
1906            }
1907            if info.encrypted_mutations_object_id != INVALID_OBJECT_ID {
1908                total_layer_size += layer_size_from_encrypted_mutations_size(
1909                    parent_store.get_file_size(info.encrypted_mutations_object_id).await?,
1910                );
1911            }
1912            last_object_id = LastObjectId::Pending;
1913            ensure!(
1914                matches!(
1915                    info.last_object_id,
1916                    LastObjectIdInfo::Encrypted { .. } | LastObjectIdInfo::Low32Bit { .. }
1917                ),
1918                FxfsError::Inconsistent
1919            );
1920        } else {
1921            last_object_id = match info.last_object_id {
1922                LastObjectIdInfo::Unencrypted { id } => LastObjectId::Unencrypted { id },
1923                LastObjectIdInfo::Low32Bit => {
1924                    LastObjectId::Low32Bit { reserved: HashSet::new(), unreserved: Vec::new() }
1925                }
1926                _ => bail!(FxfsError::Inconsistent),
1927            };
1928        }
1929
1930        let fs = parent_store.filesystem();
1931
1932        let store = ObjectStore::new(
1933            Some(parent_store.clone()),
1934            store_object_id,
1935            fs.clone(),
1936            if is_encrypted { None } else { Some(info) },
1937            object_cache,
1938            if is_encrypted { LockState::Locked } else { LockState::Unencrypted },
1939            last_object_id,
1940        );
1941
1942        assert!(store.store_info_handle.set(handle).is_ok(), "Failed to set store_info_handle!");
1943
1944        if !is_encrypted {
1945            let object_tree_layer_object_ids =
1946                store.store_info.lock().as_ref().unwrap().layers.clone();
1947            total_layer_size = store
1948                .tree
1949                .append_layers(&parent_store, object_tree_layer_object_ids, None)
1950                .await
1951                .context("Failed to read object store layers")?;
1952        }
1953
1954        fs.object_manager().update_reservation(
1955            store_object_id,
1956            tree::reservation_amount_from_layer_size(total_layer_size),
1957        );
1958
1959        Ok(store)
1960    }
1961
1962    async fn load_store_info(&self) -> Result<StoreInfo, Error> {
1963        load_store_info_from_handle(self.store_info_handle.get().unwrap()).await
1964    }
1965
1966    /// Unlocks a store so that it is ready to be used.
1967    /// This is not thread-safe.
1968    pub async fn unlock(self: &Arc<Self>, crypt: Arc<dyn Crypt>) -> Result<(), Error> {
1969        self.unlock_inner(crypt, /*read_only=*/ false).await
1970    }
1971
1972    /// Unlocks a store so that it is ready to be read from.
1973    /// The store will generally behave like it is still locked: when flushed, the store will
1974    /// write out its mutations into the encrypted mutations file, rather than directly updating
1975    /// the layer files of the object store.
1976    /// Re-locking the store (which *must* be done with `Self::lock_read_only` will not trigger a
1977    /// flush, although the store might still be flushed during other operations.
1978    /// This is not thread-safe.
1979    pub async fn unlock_read_only(self: &Arc<Self>, crypt: Arc<dyn Crypt>) -> Result<(), Error> {
1980        self.unlock_inner(crypt, /*read_only=*/ true).await
1981    }
1982
1983    // This function is *not* thread-safe.  Callers must ensure mutual exclusion when unlocking.
1984    async fn unlock_inner(
1985        self: &Arc<Self>,
1986        crypt: Arc<dyn Crypt>,
1987        read_only: bool,
1988    ) -> Result<(), Error> {
1989        // To avoid blocking compactions while waiting for a stalled crypt service during store
1990        // unlock, we perform all crypt operations outside of the flush lock. `CryptFields` tracks
1991        // the `StoreInfo` fields we depend on for these crypt operations, so we can verify they
1992        // haven't changed once we acquire the flush lock.
1993        #[derive(Debug, PartialEq)]
1994        struct CryptFields {
1995            mutations_key: Option<FxfsKey>,
1996            mutations_cipher_offset: u64,
1997            last_object_id: LastObjectIdInfo,
1998            layers: Vec<u64>,
1999        }
2000        impl From<&StoreInfo> for CryptFields {
2001            fn from(info: &StoreInfo) -> Self {
2002                Self {
2003                    mutations_key: info.mutations_key.clone(),
2004                    mutations_cipher_offset: info.mutations_cipher_offset,
2005                    last_object_id: info.last_object_id.clone(),
2006                    layers: info.layers.clone(),
2007                }
2008            }
2009        }
2010
2011        // Unless we are unlocking the store as read-only, the filesystem must not be read-only.
2012        assert!(read_only || !self.filesystem().options().read_only);
2013        match &*self.lock_state.lock() {
2014            LockState::Locked => {}
2015            LockState::Unencrypted => bail!(FxfsError::InvalidArgs),
2016            LockState::Invalid | LockState::Deleted => bail!(FxfsError::Internal),
2017            LockState::Unlocked { .. } | LockState::UnlockedReadOnly(..) => {
2018                bail!(FxfsError::AlreadyBound)
2019            }
2020            LockState::Unknown => panic!("Store was unlocked before replay"),
2021            LockState::Locking => panic!("Store is being locked"),
2022            LockState::Unlocking => panic!("Store is being unlocked"),
2023        }
2024
2025        // --- PHASE 1: Do everything that uses the crypt service outside of the flush lock ---
2026
2027        let mut last_num_flushes = self.counters.lock().num_flushes;
2028        let mut store_info = self.load_store_info().await?;
2029        let crypt_fields = CryptFields::from(&store_info);
2030
2031        let mut update_store_info = async |store_info: &mut StoreInfo| -> Result<(), Error> {
2032            let new_num_flushes = self.counters.lock().num_flushes;
2033            if last_num_flushes == new_num_flushes {
2034                return Ok(());
2035            }
2036            last_num_flushes = new_num_flushes;
2037            *store_info = self.load_store_info().await?;
2038            if crypt_fields != CryptFields::from(&*store_info) {
2039                return Err(
2040                    anyhow!(FxfsError::Inconsistent).context("Crypt fields changed during unlock")
2041                );
2042            }
2043            Ok(())
2044        };
2045
2046        // Open layers (uses crypt).
2047        let (layers, _) = open_layers(
2048            self.parent_store.as_ref().unwrap(),
2049            store_info.layers.iter().cloned(),
2050            Some(crypt.clone()),
2051        )
2052        .await
2053        .context("Failed to read object tree layer file contents")?;
2054
2055        // Unwrap mutations key.
2056        let wrapped_key =
2057            fxfs_crypto::WrappedKey::Fxfs(store_info.mutations_key.clone().unwrap().into());
2058        let unwrapped_key = crypt
2059            .unwrap_key(&wrapped_key, self.store_object_id)
2060            .await
2061            .context("Failed to unwrap mutations keys")?;
2062
2063        // Unwrap last object ID key.
2064        let last_object_id_cipher = match &store_info.last_object_id {
2065            LastObjectIdInfo::Encrypted { id: _, key } => {
2066                let wrapped_key = fxfs_crypto::WrappedKey::Fxfs(key.clone().into());
2067                let unwrapped = crypt
2068                    .unwrap_key(&wrapped_key, self.store_object_id)
2069                    .await
2070                    .context("Failed to unwrap last object ID key")?;
2071                Some(Box::new(Ff1::new(&unwrapped)))
2072            }
2073            _ => None,
2074        };
2075
2076        // Roll mutations key (create new key outside of lock).
2077        let (new_wrapped_mutations_key, new_unwrapped_mutations_key) =
2078            crypt.create_key(self.store_object_id, KeyPurpose::Metadata).await?;
2079
2080        // Pre-cache keys outside of lock.
2081        let mut keys_to_cache = Vec::new();
2082        if !read_only && !self.filesystem().options().read_only {
2083            let parent_store = self.parent_store.as_ref().unwrap();
2084            for _ in 0..CACHED_KEYS_LIMIT {
2085                let raw_id = {
2086                    let reserved_id = parent_store
2087                        .maybe_get_next_object_id()
2088                        .expect("maybe_get_next_object_id failed on parent store");
2089                    reserved_id.release()
2090                };
2091                let (wrapped, unwrapped) = crypt
2092                    .create_key(raw_id.get(), KeyPurpose::Data)
2093                    .await
2094                    .context("Failed to pre-cache key during unlock")?;
2095                keys_to_cache.push((raw_id, EncryptionKey::Fxfs(wrapped), unwrapped));
2096            }
2097        }
2098
2099        // --- PHASE 2: Take the flush lock to read mutations ---
2100
2101        let do_not_use_crypt = crypt;
2102
2103        let fs = self.filesystem();
2104        let guard =
2105            fs.lock_manager().write_lock(lock_keys![LockKey::flush(self.store_object_id())]).await;
2106
2107        update_store_info(&mut store_info).await?;
2108
2109        // Read mutations.
2110        let mut mutations = {
2111            if store_info.encrypted_mutations_object_id == INVALID_OBJECT_ID {
2112                EncryptedMutations::default()
2113            } else {
2114                let parent_store = self.parent_store.as_ref().unwrap();
2115                let handle = ObjectStore::open_object(
2116                    &parent_store,
2117                    store_info.encrypted_mutations_object_id,
2118                    HandleOptions::default(),
2119                    None,
2120                )
2121                .await?;
2122                let mut cursor = std::io::Cursor::new(
2123                    handle
2124                        .contents(MAX_ENCRYPTED_MUTATIONS_SIZE)
2125                        .await
2126                        .context(FxfsError::Inconsistent)?,
2127                );
2128                let mut mutations = EncryptedMutations::deserialize_with_version(&mut cursor)
2129                    .context("Failed to deserialize EncryptedMutations")?
2130                    .0;
2131                let len = cursor.get_ref().len() as u64;
2132                while cursor.position() < len {
2133                    mutations.extend(
2134                        &EncryptedMutations::deserialize_with_version(&mut cursor)
2135                            .context("Failed to deserialize EncryptedMutations")?
2136                            .0,
2137                    );
2138                }
2139                mutations
2140            }
2141        };
2142
2143        // This assumes that the journal has no buffered mutations for this store (see Self::lock).
2144        let journaled = EncryptedMutations::from_replayed_mutations(
2145            self.store_object_id,
2146            fs.journal()
2147                .read_transactions_for_object(self.store_object_id)
2148                .await
2149                .context("Failed to read encrypted mutations from journal")?,
2150        );
2151        mutations.extend(&journaled);
2152
2153        // Drop the lock before we do any crypt operations (decryption/unwrapping).
2154        std::mem::drop(guard);
2155
2156        // It's safe to use the crypt service now because we're not holding the flush lock.
2157        let crypt = do_not_use_crypt;
2158
2159        self.filesystem().hooks().on_unlock_resources_acquired();
2160
2161        // --- PHASE 3: Decrypt mutations (outside of lock) ---
2162
2163        let EncryptedMutations { transactions, mut data, mutations_key_roll } = mutations;
2164
2165        let mut unwrapped_key_rolls = Vec::with_capacity(mutations_key_roll.len());
2166        for (offset, key) in mutations_key_roll {
2167            let unwrapped = crypt
2168                .unwrap_key(&fxfs_crypto::WrappedKey::Fxfs(key.into()), self.store_object_id)
2169                .await
2170                .context("Failed to unwrap mutations keys")?;
2171            unwrapped_key_rolls.push((offset, unwrapped));
2172        }
2173
2174        let mut mutations_to_apply = Vec::new();
2175
2176        // All the data, transactions and keyrolls will be split with the first half being in the
2177        // older version and the rest in the new. Find that split so that they can work separately
2178        // from each other. Though it is very likely that they're all in one version or the other.
2179        let split_idx = transactions
2180            .partition_point(|(checkpoint, _)| checkpoint.version < AES_JOURNAL_ENCRYPTION_VERSION);
2181        let (chacha_transactions, aes_transactions) = transactions.split_at(split_idx);
2182        let data_offset = if chacha_transactions.is_empty() {
2183            0
2184        } else {
2185            let aes_len = aes_transactions
2186                .iter()
2187                .try_fold(0usize, |total, (_, next)| total.checked_add(*next as usize))
2188                .ok_or(FxfsError::Inconsistent)?;
2189            ensure!(aes_len <= data.len(), FxfsError::Inconsistent);
2190            data.len() - aes_len
2191        };
2192        let (chacha_data, aes_data) = data.split_at_mut(data_offset);
2193
2194        // Split the keyrolls based on the data offset found above, remap AES rolls based on the
2195        // partial data set it receives.
2196        let split_key_idx =
2197            unwrapped_key_rolls.partition_point(|(offset, _)| *offset < data_offset);
2198        let aes_key_rolls: Vec<(usize, UnwrappedKey)> = unwrapped_key_rolls
2199            .split_off(split_key_idx)
2200            .into_iter()
2201            .map(|(offset, key)| (offset - data_offset, key))
2202            .collect();
2203        let chacha_key_rolls = unwrapped_key_rolls;
2204        let mut cipher_sequence_number = store_info.mutations_cipher_offset;
2205
2206        if !chacha_transactions.is_empty() {
2207            ensure!(store_info.mutations_cipher_offset <= u32::MAX as u64, FxfsError::Inconsistent);
2208            let mut cipher = StreamCipher::new(&unwrapped_key, cipher_sequence_number);
2209            mutations_to_apply.append(&mut Self::decrypt_mutations_chacha20(
2210                &mut cipher,
2211                chacha_transactions,
2212                chacha_data,
2213                &chacha_key_rolls,
2214            )?);
2215            cipher_sequence_number = cipher.offset();
2216        }
2217
2218        if !aes_transactions.is_empty() {
2219            let mut cipher = JournalXtsCipher::new(&unwrapped_key, cipher_sequence_number);
2220            mutations_to_apply.append(&mut Self::decrypt_mutations_aes256_xts(
2221                &mut cipher,
2222                aes_transactions,
2223                aes_data,
2224                &aes_key_rolls,
2225            )?);
2226            cipher_sequence_number = cipher.current_tweak();
2227        }
2228
2229        // --- PHASE 4: Re-acquire the flush lock and apply changes ---
2230
2231        let do_not_use_crypt = crypt;
2232
2233        let guard =
2234            fs.lock_manager().write_lock(lock_keys![LockKey::flush(self.store_object_id())]).await;
2235
2236        update_store_info(&mut store_info).await?;
2237
2238        let _ = std::mem::replace(&mut *self.lock_state.lock(), LockState::Unlocking);
2239
2240        store_info.mutations_key = Some(new_wrapped_mutations_key);
2241        *self.store_info.lock() = Some(store_info.clone());
2242
2243        let clean_up = scopeguard::guard((), |_| {
2244            *self.lock_state.lock() = LockState::Locked;
2245            *self.store_info.lock() = None;
2246            *self.last_object_id.lock() = LastObjectId::Pending;
2247            // Make sure we don't leave unencrypted data lying around in memory.
2248            self.tree.reset();
2249        });
2250
2251        // Apply layers.
2252        self.tree.append_open_layers(layers);
2253
2254        // Set last object ID.
2255        match &store_info.last_object_id {
2256            LastObjectIdInfo::Encrypted { id, .. } => {
2257                *self.last_object_id.lock() =
2258                    LastObjectId::Encrypted { id: *id, cipher: last_object_id_cipher.unwrap() };
2259            }
2260            LastObjectIdInfo::Low32Bit => {
2261                *self.last_object_id.lock() = LastObjectId::Low32Bit {
2262                    reserved: Default::default(),
2263                    unreserved: Default::default(),
2264                }
2265            }
2266            _ => unreachable!(),
2267        }
2268
2269        // Apply mutations.
2270        for (checkpoint, mutation) in mutations_to_apply {
2271            let context = ApplyContext { mode: ApplyMode::Replay, checkpoint };
2272            self.apply_mutation(mutation, &context, AssocObj::None)
2273                .context("failed to apply encrypted mutation")?;
2274        }
2275
2276        // Transition to Unlocked.
2277        *self.lock_state.lock() = if read_only {
2278            LockState::UnlockedReadOnly(do_not_use_crypt)
2279        } else {
2280            LockState::Unlocked {
2281                crypt: do_not_use_crypt,
2282                cached_keys: keys_to_cache,
2283                mutations_cipher: JournalCipher::new_aes256_xts(
2284                    &new_unwrapped_mutations_key,
2285                    cipher_sequence_number,
2286                ),
2287            }
2288        };
2289
2290        // To avoid unbounded memory growth, we should flush the encrypted mutations now. Otherwise
2291        // it's possible for more writes to be queued and for the store to be locked before we can
2292        // flush anything and that can repeat.
2293        std::mem::drop(guard);
2294
2295        if !read_only && !self.filesystem().options().read_only {
2296            self.flush_with_reason(FlushReason::EncryptedMutations).await?;
2297
2298            // Reap purged files within this store.
2299            let _ = self.filesystem().graveyard().initial_reap(&self).await?;
2300        }
2301
2302        // Return and cancel the clean up.
2303        Ok(ScopeGuard::into_inner(clean_up))
2304    }
2305
2306    fn decrypt_mutations_aes256_xts(
2307        cipher: &mut JournalXtsCipher,
2308        transactions: &[(JournalCheckpoint, u64)],
2309        data: &mut [u8],
2310        key_rolls: &[(usize, UnwrappedKey)],
2311    ) -> Result<Vec<(JournalCheckpoint, Mutation)>, Error> {
2312        let mut mutations_to_apply = Vec::new();
2313        let mut key_roll_iter = key_rolls.iter().peekable();
2314        let mut offset = 0;
2315
2316        for (checkpoint, chunk_len) in transactions {
2317            ensure!(checkpoint.version >= AES_JOURNAL_ENCRYPTION_VERSION, FxfsError::Inconsistent);
2318            while let Some((roll_offset, unwrapped_key)) = key_roll_iter.peek() {
2319                if offset >= *roll_offset {
2320                    key_roll_iter.next();
2321                    *cipher = JournalXtsCipher::new(unwrapped_key, cipher.current_tweak());
2322                } else {
2323                    break;
2324                }
2325            }
2326            let chunk_len = *chunk_len as usize;
2327            ensure!(offset + chunk_len <= data.len(), FxfsError::Inconsistent);
2328            let chunk = &mut data[offset..offset + chunk_len];
2329            let transaction =
2330                EncryptedTransaction::decrypt_and_deserialize(cipher, chunk, checkpoint.version)?;
2331            for mutation in transaction.0 {
2332                mutations_to_apply.push((checkpoint.clone(), mutation));
2333            }
2334            offset += chunk_len;
2335        }
2336        Ok(mutations_to_apply)
2337    }
2338
2339    fn decrypt_mutations_chacha20(
2340        cipher: &mut StreamCipher,
2341        transactions: &[(JournalCheckpoint, u64)],
2342        data: &mut [u8],
2343        key_rolls: &[(usize, UnwrappedKey)],
2344    ) -> Result<Vec<(JournalCheckpoint, Mutation)>, Error> {
2345        let mut mutations_to_apply = Vec::new();
2346        let mut slice = &mut data[..];
2347        let mut last_offset = 0;
2348        for (offset, unwrapped_key) in key_rolls {
2349            let split_offset = offset.checked_sub(last_offset).ok_or(FxfsError::Inconsistent)?;
2350            ensure!(split_offset <= slice.len(), FxfsError::Inconsistent);
2351            last_offset = *offset;
2352            let (old, new) = slice.split_at_mut(split_offset);
2353            cipher.decrypt(old);
2354            *cipher = StreamCipher::new(unwrapped_key, 0);
2355            slice = new;
2356        }
2357        cipher.decrypt(slice);
2358
2359        let mut cursor = std::io::Cursor::new(&*data);
2360        for (checkpoint, count) in transactions {
2361            for _ in 0..*count {
2362                let mutation = Mutation::deserialize_from_version(&mut cursor, checkpoint.version)
2363                    .context("failed to deserialize encrypted mutation")?;
2364                mutations_to_apply.push((checkpoint.clone(), mutation));
2365            }
2366        }
2367        Ok(mutations_to_apply)
2368    }
2369
2370    pub fn is_locked(&self) -> bool {
2371        matches!(
2372            *self.lock_state.lock(),
2373            LockState::Locked | LockState::Locking | LockState::Unknown
2374        )
2375    }
2376
2377    /// NB: This is not the converse of `is_locked`, as there are lock states where neither are
2378    /// true.
2379    pub fn is_unlocked(&self) -> bool {
2380        matches!(
2381            *self.lock_state.lock(),
2382            LockState::Unlocked { .. } | LockState::UnlockedReadOnly(..) | LockState::Unlocking
2383        )
2384    }
2385
2386    pub fn is_unknown(&self) -> bool {
2387        matches!(*self.lock_state.lock(), LockState::Unknown)
2388    }
2389
2390    pub fn is_encrypted(&self) -> bool {
2391        self.store_info.lock().as_ref().unwrap().mutations_key.is_some()
2392    }
2393
2394    // Locks a store.
2395    // This operation will take a flush lock on the store, in case any flushes are ongoing.  Any
2396    // ongoing store accesses might be interrupted by this.  See `Self::crypt`.
2397    // Whilst this can return an error, the store will be placed into an unusable but safe state
2398    // (i.e. no lingering unencrypted data) if an error is encountered.
2399    pub async fn lock(&self) -> Result<(), Error> {
2400        // We must lock flushing since it is not safe for that to be happening whilst we are locking
2401        // the store.
2402        let keys = lock_keys![LockKey::flush(self.store_object_id())];
2403        let fs = self.filesystem();
2404        let _guard = fs.lock_manager().write_lock(keys).await;
2405
2406        {
2407            let mut lock_state = self.lock_state.lock();
2408            if let LockState::Unlocked { .. } = &*lock_state {
2409                *lock_state = LockState::Locking;
2410            } else {
2411                panic!("Unexpected lock state: {:?}", *lock_state);
2412            }
2413        }
2414
2415        // Sync the journal now to ensure that any buffered mutations for this store make it out to
2416        // disk.  This is necessary to be able to unlock the store again.
2417        // We need to establish a barrier at this point (so that the journaled writes are observable
2418        // by any future attempts to unlock the store), hence the flush_device.
2419        let sync_result =
2420            self.filesystem().sync(SyncOptions { flush_device: true, ..Default::default() }).await;
2421
2422        *self.lock_state.lock() = if let Err(error) = &sync_result {
2423            error!(error:?; "Failed to sync journal; store will no longer be usable");
2424            LockState::Invalid
2425        } else {
2426            LockState::Locked
2427        };
2428        self.key_manager.clear();
2429        *self.store_info.lock() = None;
2430        *self.last_object_id.lock() = LastObjectId::Pending;
2431        self.tree.reset();
2432
2433        sync_result
2434    }
2435
2436    // Locks a store which was previously unlocked read-only (see `Self::unlock_read_only`).  Data
2437    // is not flushed, and instead any journaled mutations are buffered back into the ObjectStore
2438    // and will be replayed next time the store is unlocked.
2439    pub fn lock_read_only(&self) {
2440        *self.lock_state.lock() = LockState::Locked;
2441        *self.store_info.lock() = None;
2442        *self.last_object_id.lock() = LastObjectId::Pending;
2443        self.tree.reset();
2444    }
2445
2446    // Returns None if the object ID cipher needs to be created or rolled, or a more expensive
2447    // algorithm needs to be used.
2448    fn maybe_get_next_object_id(&self) -> Option<ReservedId<'_>> {
2449        self.last_object_id.lock().try_get_next().map(|id| ReservedId::new(self, id))
2450    }
2451
2452    /// Returns a new object ID that can be used.  This will create an object ID cipher if needed.
2453    ///
2454    /// If the object ID key needs to be rolled, a new transaction will be created and committed.
2455    pub(super) async fn get_next_object_id(
2456        &self,
2457        parent_transaction: &Transaction<'_>,
2458    ) -> Result<ReservedId<'_>, Error> {
2459        let low_32_bit = {
2460            let mut last_object_id = self.last_object_id.lock();
2461            if let Some(id) = last_object_id.try_get_next() {
2462                return Ok(ReservedId::new(self, id));
2463            }
2464            ensure!(
2465                !matches!(&*last_object_id, LastObjectId::Unencrypted { .. }),
2466                FxfsError::Inconsistent
2467            );
2468            matches!(&*last_object_id, LastObjectId::Low32Bit { .. })
2469        };
2470
2471        let parent_store = self.parent_store().unwrap();
2472        let lock_keys =
2473            lock_keys![LockKey::object(parent_store.store_object_id, self.store_object_id)];
2474        if low_32_bit {
2475            // NOTE: Since `parent_transaction` may already hold locks, we must take care to avoid
2476            // deadlocks; no more locks should be taken whilst we hold this lock.
2477            let fs = self.filesystem();
2478            let _guard = fs.lock_manager().txn_lock(lock_keys).await;
2479
2480            // Keep picking an object ID at random until we find one free.
2481
2482            // To avoid races, this must be before we capture the layer set.
2483            self.last_object_id.lock().drain_unreserved();
2484
2485            let layer_set = self.tree.layer_set();
2486            let mut key = ObjectKey::object(0);
2487            loop {
2488                let next_id = rand::rng().next_u32() as u64;
2489                let Some(next_id) = NonZero::new(next_id) else { continue };
2490                if self.last_object_id.lock().is_reserved(next_id.get()) {
2491                    continue;
2492                }
2493                key.object_id = next_id.get();
2494                if layer_set.key_exists(&key).await? == Existence::Missing {
2495                    self.last_object_id.lock().reserve(next_id.get());
2496                    return Ok(ReservedId::new(self, next_id));
2497                }
2498            }
2499        } else {
2500            // Create a transaction (which has a lock) and then check again.
2501            //
2502            // NOTE: Since this is a nested transaction, we must take care to avoid deadlocks; no
2503            // more locks should be taken whilst we hold this lock.
2504            let mut transaction = parent_store
2505                .new_transaction(
2506                    lock_keys,
2507                    Options { parent_transaction: Some(parent_transaction), ..Default::default() },
2508                )
2509                .await?;
2510
2511            let next_id_hi = {
2512                let mut last_object_id = self.last_object_id.lock();
2513                if let Some(id) = last_object_id.try_get_next() {
2514                    // Something else raced and created/rolled the cipher.
2515                    return Ok(ReservedId::new(self, id));
2516                }
2517
2518                match &*last_object_id {
2519                    LastObjectId::Encrypted { id, .. } => {
2520                        // It shouldn't be possible for last_object_id to wrap within our lifetime,
2521                        // so if this happens, it's most likely due to corruption.
2522                        info!(store_id = self.store_object_id; "Rolling object ID key");
2523
2524                        id.checked_add(1 << 32).ok_or(FxfsError::Inconsistent)? & OBJECT_ID_HI_MASK
2525                    }
2526                    _ => unreachable!(),
2527                }
2528            };
2529
2530            // Create a key.
2531            let (object_id_wrapped, object_id_unwrapped) = self
2532                .crypt()
2533                .unwrap()
2534                .create_key(self.store_object_id, KeyPurpose::Metadata)
2535                .await?;
2536
2537            // Normally we would use a mutation to note the updated key, but that would complicate
2538            // replay.  During replay, we need to keep track of the highest used object ID and this
2539            // is done by watching mutations to see when we create objects, and then decrypting
2540            // the object ID.  This relies on the unwrapped key being available, so as soon as
2541            // we detect the key has changed, we would need to immediately unwrap the key via the
2542            // crypt service.  Currently, this isn't easy to do during replay.  An option we could
2543            // consider would be to include the unencrypted object ID when we create objects, which
2544            // would avoid us having to decrypt the object ID during replay.
2545            //
2546            // For now and for historical reasons, the approach we take is to just write a new
2547            // version of StoreInfo here.  We must take care that we only update the key and not any
2548            // other information contained within StoreInfo because other information should only be
2549            // updated when we flush.  We are holding the lock on the StoreInfo file, so this will
2550            // prevent potential races with flushing.  To make sure we only change the key, we read
2551            // StoreInfo from storage rather than using our in-memory copy.  This won't be
2552            // performant, but rolling the object ID key will be extremely rare.
2553            let new_store_info = StoreInfo {
2554                last_object_id: LastObjectIdInfo::Encrypted {
2555                    id: next_id_hi,
2556                    key: object_id_wrapped.clone(),
2557                },
2558                ..self.load_store_info().await?
2559            };
2560
2561            self.write_store_info(&mut transaction, &new_store_info).await?;
2562
2563            transaction
2564                .commit_with_callback(|_| {
2565                    self.store_info.lock().as_mut().unwrap().last_object_id =
2566                        new_store_info.last_object_id;
2567                    match &mut *self.last_object_id.lock() {
2568                        LastObjectId::Encrypted { id, cipher } => {
2569                            **cipher = Ff1::new(&object_id_unwrapped);
2570                            *id = next_id_hi;
2571                            ReservedId::new(
2572                                self,
2573                                NonZero::new(next_id_hi | cipher.encrypt(0) as u64).unwrap(),
2574                            )
2575                        }
2576                        _ => unreachable!(),
2577                    }
2578                })
2579                .await
2580        }
2581    }
2582
2583    /// Query the next object ID that will be used. Intended for use when checking filesystem
2584    /// consistency. Prefer [`Self::get_next_object_id()`] for general use.
2585    pub(crate) fn query_next_object_id(&self) -> u64 {
2586        self.last_object_id.lock().peek_next()
2587    }
2588
2589    fn allocator(&self) -> Arc<Allocator> {
2590        self.filesystem().allocator()
2591    }
2592
2593    // If |transaction| has an impending mutation for the underlying object, returns that.
2594    // Otherwise, looks up the object from the tree and returns a suitable mutation for it.  The
2595    // mutation is returned here rather than the item because the mutation includes the operation
2596    // which has significance: inserting an object implies it's the first of its kind unlike
2597    // replacing an object.
2598    async fn txn_get_object_mutation(
2599        &self,
2600        transaction: &Transaction<'_>,
2601        object_id: u64,
2602    ) -> Result<ObjectStoreMutation, Error> {
2603        if let Some(mutation) =
2604            transaction.get_object_mutation(self.store_object_id, ObjectKey::object(object_id))
2605        {
2606            Ok(mutation.clone())
2607        } else {
2608            Ok(ObjectStoreMutation {
2609                item: self
2610                    .tree
2611                    .find(&ObjectKey::object(object_id))
2612                    .await?
2613                    .ok_or(FxfsError::Inconsistent)
2614                    .context("Object id missing")?,
2615                op: Operation::ReplaceOrInsert,
2616            })
2617        }
2618    }
2619
2620    /// Like txn_get_object_mutation but with expanded visibility.
2621    /// Only available in migration code.
2622    #[cfg(feature = "migration")]
2623    pub async fn get_object_mutation(
2624        &self,
2625        transaction: &Transaction<'_>,
2626        object_id: u64,
2627    ) -> Result<ObjectStoreMutation, Error> {
2628        self.txn_get_object_mutation(transaction, object_id).await
2629    }
2630
2631    fn update_last_object_id(&self, object_id: u64) {
2632        let mut last_object_id = self.last_object_id.lock();
2633        match &mut *last_object_id {
2634            LastObjectId::Pending => unreachable!(),
2635            LastObjectId::Unencrypted { id } => {
2636                if object_id > *id {
2637                    *id = object_id
2638                }
2639            }
2640            LastObjectId::Encrypted { id, cipher } => {
2641                // For encrypted stores, object_id will be encrypted here, so we must decrypt first.
2642
2643                // If the object ID cipher has been rolled, then it's possible we might see object
2644                // IDs that were generated using a different cipher so the decrypt here will return
2645                // the wrong value, but that won't matter because the hi part of the object ID
2646                // should still discriminate.
2647                let object_id =
2648                    object_id & OBJECT_ID_HI_MASK | cipher.decrypt(object_id as u32) as u64;
2649                if object_id > *id {
2650                    *id = object_id;
2651                }
2652            }
2653            LastObjectId::Low32Bit { .. } => {}
2654        }
2655    }
2656
2657    /// If possible, converts the given object ID to its unencrypted value.  Returns None if it is
2658    /// not possible to convert to its unencrypted value because the key is unavailable.
2659    pub fn to_unencrypted_object_id(&self, object_id: u64) -> Option<u64> {
2660        let last_object_id = self.last_object_id.lock();
2661        match &*last_object_id {
2662            LastObjectId::Pending => None,
2663            LastObjectId::Unencrypted { .. } | LastObjectId::Low32Bit { .. } => Some(object_id),
2664            LastObjectId::Encrypted { id, cipher } => {
2665                if id & OBJECT_ID_HI_MASK != object_id & OBJECT_ID_HI_MASK {
2666                    None
2667                } else {
2668                    Some(object_id & OBJECT_ID_HI_MASK | cipher.decrypt(object_id as u32) as u64)
2669                }
2670            }
2671        }
2672    }
2673
2674    /// Adds the specified object to the graveyard.
2675    pub fn add_to_graveyard(&self, transaction: &mut Transaction<'_>, object_id: u64) {
2676        let graveyard_id = self.graveyard_directory_object_id();
2677        assert_ne!(graveyard_id, INVALID_OBJECT_ID);
2678        transaction.add(
2679            self.store_object_id,
2680            Mutation::replace_or_insert_object(
2681                ObjectKey::graveyard_entry(graveyard_id, object_id),
2682                ObjectValue::Some,
2683            ),
2684        );
2685    }
2686
2687    /// Removes the specified object from the graveyard.  NB: Care should be taken when calling
2688    /// this because graveyard entries are used for purging deleted files *and* for trimming
2689    /// extents.  For example, consider the following sequence:
2690    ///
2691    ///     1. Add Trim graveyard entry.
2692    ///     2. Replace with Some graveyard entry (see above).
2693    ///     3. Remove graveyard entry.
2694    ///
2695    /// If the desire in #3 is just to cancel the effect of the Some entry, then #3 should
2696    /// actually be:
2697    ///
2698    ///     3. Replace with Trim graveyard entry.
2699    pub fn remove_from_graveyard(&self, transaction: &mut Transaction<'_>, object_id: u64) {
2700        transaction.add(
2701            self.store_object_id,
2702            Mutation::replace_or_insert_object(
2703                ObjectKey::graveyard_entry(self.graveyard_directory_object_id(), object_id),
2704                ObjectValue::None,
2705            ),
2706        );
2707    }
2708
2709    /// Removes the specified attribute from the graveyard. Unlike object graveyard entries,
2710    /// attribute graveyard entries only have one functionality (i.e. to purge deleted attributes)
2711    /// so the caller does not need to be concerned about replacing the graveyard attribute entry
2712    /// with its prior state when cancelling it. See comment on `remove_from_graveyard()`.
2713    pub fn remove_attribute_from_graveyard(
2714        &self,
2715        transaction: &mut Transaction<'_>,
2716        object_id: u64,
2717        attribute_id: AttributeId,
2718    ) {
2719        transaction.add(
2720            self.store_object_id,
2721            Mutation::replace_or_insert_object(
2722                ObjectKey::graveyard_attribute_entry(
2723                    self.graveyard_directory_object_id(),
2724                    object_id,
2725                    attribute_id,
2726                ),
2727                ObjectValue::None,
2728            ),
2729        );
2730    }
2731
2732    // When the symlink is unlocked, this function decrypts `link` and returns a bag of bytes that
2733    // is identical to that which was passed in as the target on `create_symlink`.
2734    // If the symlink is locked, this function hashes the encrypted `link` with Sha256 in order to
2735    // get a standard length and then base64 encodes the hash and returns that to the caller.
2736    pub async fn read_encrypted_symlink(
2737        &self,
2738        object_id: u64,
2739        link: Vec<u8>,
2740    ) -> Result<Vec<u8>, Error> {
2741        let mut link = link;
2742        let key = self
2743            .key_manager()
2744            .get_fscrypt_key(object_id, self.crypt().unwrap().as_ref(), async || {
2745                self.get_keys(object_id).await
2746            })
2747            .await?;
2748        if let Some(key) = key.into_cipher() {
2749            key.decrypt_symlink(object_id, &mut link)?;
2750            Ok(link)
2751        } else {
2752            // Locked symlinks are encoded using a hash_code of 0.
2753            let proxy_filename =
2754                fscrypt::proxy_filename::ProxyFilename::new_with_hash_code(0, &link);
2755            let proxy_filename_str: String = proxy_filename.into();
2756            Ok(proxy_filename_str.into_bytes())
2757        }
2758    }
2759
2760    /// Returns the link of a symlink object.
2761    pub async fn read_symlink(&self, object_id: u64) -> Result<Vec<u8>, Error> {
2762        match self.tree.find_value(&ObjectKey::object(object_id)).await? {
2763            None => bail!(FxfsError::NotFound),
2764            Some(ObjectValue::Object {
2765                kind: ObjectKind::EncryptedSymlink { link, .. }, ..
2766            }) => self.read_encrypted_symlink(object_id, link.into_vec()).await,
2767            Some(ObjectValue::Object { kind: ObjectKind::Symlink { link, .. }, .. }) => {
2768                Ok(link.into_vec())
2769            }
2770            Some(value) => Err(anyhow!(FxfsError::Inconsistent)
2771                .context(format!("Unexpected value in symlink lookup: {value:?}"))),
2772        }
2773    }
2774
2775    /// Fail if any keys in the root store are not basic Fxfs keys. Fscrypt keys shouldn't be
2776    /// there and use a weaker key derivation function. Prevent accidentally using it on a volume's
2777    /// root key.
2778    fn fail_on_illegal_keys(&self, keys: &EncryptionKeys) -> Result<(), Error> {
2779        if !self.is_root() {
2780            return Ok(());
2781        }
2782
2783        for key in keys.iter() {
2784            if !matches!(key.1, EncryptionKey::Fxfs(_)) {
2785                return Err(
2786                    anyhow!(FxfsError::IntegrityError).context("Illegal key type in root store")
2787                );
2788            }
2789        }
2790        Ok(())
2791    }
2792
2793    /// Retrieves the wrapped keys for the given object.  The keys *should* be known to exist and it
2794    /// will be considered an inconsistency if they don't.
2795    pub async fn get_keys(&self, object_id: u64) -> Result<EncryptionKeys, Error> {
2796        self.tree
2797            .find_map(&ObjectKey::keys(object_id), |item| match item.value {
2798                ObjectValue::Keys(keys) => {
2799                    self.fail_on_illegal_keys(keys)?;
2800                    Ok(keys.clone())
2801                }
2802                _ => Err(anyhow!(FxfsError::Inconsistent).context("open_object: Expected keys")),
2803            })
2804            .await?
2805            .ok_or(FxfsError::Inconsistent)?
2806    }
2807
2808    pub async fn update_attributes<'a>(
2809        &self,
2810        transaction: &mut Transaction<'a>,
2811        object_id: u64,
2812        node_attributes: Option<&fio::MutableNodeAttributes>,
2813        change_time: Option<Timestamp>,
2814    ) -> Result<(), Error> {
2815        if change_time.is_none() {
2816            if let Some(attributes) = node_attributes {
2817                let empty_attributes = fio::MutableNodeAttributes { ..Default::default() };
2818                if *attributes == empty_attributes {
2819                    return Ok(());
2820                }
2821            } else {
2822                return Ok(());
2823            }
2824        }
2825        let mut mutation = self.txn_get_object_mutation(transaction, object_id).await?;
2826        if let ObjectValue::Object { ref mut attributes, .. } = mutation.item.value {
2827            if let Some(time) = change_time {
2828                attributes.change_time = time;
2829            }
2830            if let Some(node_attributes) = node_attributes {
2831                if let Some(time) = node_attributes.creation_time {
2832                    attributes.creation_time = Timestamp::from_nanos(time);
2833                }
2834                if let Some(time) = node_attributes.modification_time {
2835                    attributes.modification_time = Timestamp::from_nanos(time);
2836                }
2837                if let Some(time) = node_attributes.access_time {
2838                    attributes.access_time = Timestamp::from_nanos(time);
2839                }
2840                if node_attributes.mode.is_some()
2841                    || node_attributes.uid.is_some()
2842                    || node_attributes.gid.is_some()
2843                    || node_attributes.rdev.is_some()
2844                {
2845                    if let Some(a) = &mut attributes.posix_attributes {
2846                        if let Some(mode) = node_attributes.mode {
2847                            a.mode = mode;
2848                        }
2849                        if let Some(uid) = node_attributes.uid {
2850                            a.uid = uid;
2851                        }
2852                        if let Some(gid) = node_attributes.gid {
2853                            a.gid = gid;
2854                        }
2855                        if let Some(rdev) = node_attributes.rdev {
2856                            a.rdev = rdev;
2857                        }
2858                    } else {
2859                        attributes.posix_attributes = Some(PosixAttributes {
2860                            mode: node_attributes.mode.unwrap_or_default(),
2861                            uid: node_attributes.uid.unwrap_or_default(),
2862                            gid: node_attributes.gid.unwrap_or_default(),
2863                            rdev: node_attributes.rdev.unwrap_or_default(),
2864                        });
2865                    }
2866                }
2867            }
2868        } else {
2869            bail!(
2870                anyhow!(FxfsError::Inconsistent)
2871                    .context("ObjectStore.update_attributes: Expected object value")
2872            );
2873        };
2874        transaction.add(self.store_object_id(), Mutation::ObjectStore(mutation));
2875        Ok(())
2876    }
2877
2878    // Updates and commits the changes to access time in ObjectProperties. The update matches
2879    // Linux's RELATIME. That is, access time is updated to the current time if access time is less
2880    // than or equal to the last modification or status change, or if it has been more than a day
2881    // since the last access.  `precondition` is a condition to be checked *after* taking the lock
2882    // on the object.  If `precondition` returns false, no update will be performed.
2883    pub async fn update_access_time(
2884        &self,
2885        object_id: u64,
2886        props: &mut ObjectProperties,
2887        precondition: impl FnOnce() -> bool,
2888    ) -> Result<(), Error> {
2889        let access_time = props.access_time.as_nanos();
2890        let modification_time = props.modification_time.as_nanos();
2891        let change_time = props.change_time.as_nanos();
2892        let now = Timestamp::now();
2893        if access_time <= modification_time
2894            || access_time <= change_time
2895            || access_time
2896                < now.as_nanos()
2897                    - Timestamp::from(std::time::Duration::from_secs(24 * 60 * 60)).as_nanos()
2898        {
2899            let mut transaction = self
2900                .new_transaction(
2901                    lock_keys![LockKey::object(self.store_object_id, object_id,)],
2902                    Options {
2903                        reservation: ReservationOptions::BorrowedMetadata,
2904                        ..Default::default()
2905                    },
2906                )
2907                .await?;
2908            if precondition() {
2909                self.update_attributes(
2910                    &mut transaction,
2911                    object_id,
2912                    Some(&fio::MutableNodeAttributes {
2913                        access_time: Some(now.as_nanos()),
2914                        ..Default::default()
2915                    }),
2916                    None,
2917                )
2918                .await?;
2919                transaction.commit().await?;
2920                props.access_time = now;
2921            }
2922        }
2923        Ok(())
2924    }
2925
2926    async fn write_store_info<'a>(
2927        &'a self,
2928        transaction: &mut Transaction<'a>,
2929        info: &StoreInfo,
2930    ) -> Result<(), Error> {
2931        let mut serialized_info = Vec::new();
2932        info.serialize_with_version(&mut serialized_info)?;
2933        let mut buf = self.device.allocate_buffer(serialized_info.len()).await;
2934        buf.copy_from_slice(&serialized_info);
2935        self.store_info_handle.get().unwrap().txn_write(transaction, 0u64, buf.as_ref()).await
2936    }
2937
2938    pub fn mark_deleted(&self) {
2939        *self.lock_state.lock() = LockState::Deleted;
2940    }
2941
2942    #[cfg(test)]
2943    pub(crate) fn test_set_last_object_id(&self, object_id: u64) {
2944        match &mut *self.last_object_id.lock() {
2945            LastObjectId::Encrypted { id, .. } => *id = object_id,
2946            _ => unreachable!(),
2947        }
2948    }
2949
2950    /// Looks up the size of the attribute. Returns an error if either the object or attribute
2951    /// doesn't exist.
2952    pub async fn get_attribute_size(
2953        &self,
2954        object_id: u64,
2955        attribute_id: AttributeId,
2956    ) -> Result<u64, Error> {
2957        let size = self
2958            .tree
2959            .find_map(
2960                &ObjectKey::attribute(object_id, attribute_id, AttributeKey::Attribute),
2961                |item| match item.value {
2962                    ObjectValue::Attribute { size, .. }
2963                    | ObjectValue::VerifiedAttribute { size, .. } => Ok(*size),
2964                    _ => Err(anyhow!(FxfsError::Inconsistent)),
2965                },
2966            )
2967            .await?
2968            .ok_or(FxfsError::NotFound)??;
2969        Ok(size)
2970    }
2971}
2972
2973#[async_trait]
2974impl JournalingObject for ObjectStore {
2975    fn apply_mutation(
2976        &self,
2977        mutation: Mutation,
2978        context: &ApplyContext<'_, '_>,
2979        _assoc_obj: AssocObj<'_>,
2980    ) -> Result<(), Error> {
2981        match &*self.lock_state.lock() {
2982            LockState::Locked | LockState::Locking => {
2983                ensure!(
2984                    matches!(mutation, Mutation::BeginFlush | Mutation::EndFlush)
2985                        || matches!(
2986                            mutation,
2987                            Mutation::EncryptedObjectStore(_) | Mutation::UpdateMutationsKey(_)
2988                                if context.mode.is_replay()
2989                        ),
2990                    anyhow!(FxfsError::Inconsistent)
2991                        .context(format!("Unexpected mutation for encrypted store: {mutation:?}"))
2992                );
2993            }
2994            LockState::Invalid
2995            | LockState::Unlocking
2996            | LockState::Unencrypted
2997            | LockState::Unlocked { .. }
2998            | LockState::UnlockedReadOnly(..)
2999            | LockState::Deleted => {}
3000            lock_state @ _ => panic!("Unexpected lock state: {lock_state:?}"),
3001        }
3002        match mutation {
3003            Mutation::ObjectStore(ObjectStoreMutation { item, op }) => {
3004                match op {
3005                    Operation::Insert => {
3006                        let mut unreserve_id = INVALID_OBJECT_ID;
3007                        // If we are inserting an object record for the first time, it signifies the
3008                        // birth of the object so we need to adjust the object count.
3009                        if matches!(item.value, ObjectValue::Object { .. }) {
3010                            {
3011                                let info = &mut self.store_info.lock();
3012                                let object_count = &mut info.as_mut().unwrap().object_count;
3013                                *object_count = object_count.saturating_add(1);
3014                            }
3015                            if context.mode.is_replay() {
3016                                self.update_last_object_id(item.key.object_id);
3017                            } else {
3018                                unreserve_id = item.key.object_id;
3019                            }
3020                        } else if !context.mode.is_replay()
3021                            && matches!(
3022                                item.key.data,
3023                                ObjectKeyData::GraveyardEntry { .. }
3024                                    | ObjectKeyData::GraveyardAttributeEntry { .. }
3025                            )
3026                        {
3027                            if matches!(item.value, ObjectValue::Some | ObjectValue::Trim) {
3028                                self.graveyard_entries.fetch_add(1, Ordering::Relaxed);
3029                            } else if matches!(item.value, ObjectValue::None) {
3030                                self.graveyard_entries.fetch_sub(1, Ordering::Relaxed);
3031                            }
3032                        }
3033                        self.tree.insert(item)?;
3034                        if unreserve_id != INVALID_OBJECT_ID {
3035                            // To avoid races, this *must* be after the `tree.insert(..)` above.
3036                            self.last_object_id.lock().unreserve(unreserve_id);
3037                        }
3038                    }
3039                    Operation::ReplaceOrInsert => {
3040                        if !context.mode.is_replay()
3041                            && matches!(
3042                                item.key.data,
3043                                ObjectKeyData::GraveyardEntry { .. }
3044                                    | ObjectKeyData::GraveyardAttributeEntry { .. }
3045                            )
3046                        {
3047                            if matches!(item.value, ObjectValue::Some | ObjectValue::Trim) {
3048                                self.graveyard_entries.fetch_add(1, Ordering::Relaxed);
3049                            } else if matches!(item.value, ObjectValue::None) {
3050                                self.graveyard_entries.fetch_sub(1, Ordering::Relaxed);
3051                            }
3052                        }
3053                        self.tree.replace_or_insert(item);
3054                    }
3055                    Operation::Merge => {
3056                        if item.is_tombstone() {
3057                            let info = &mut self.store_info.lock();
3058                            let object_count = &mut info.as_mut().unwrap().object_count;
3059                            *object_count = object_count.saturating_sub(1);
3060                            if !context.mode.is_replay() {
3061                                // Evict the object's keys from `key_manager` when applying the
3062                                // tombstone mutation at commit time (rather than when building the
3063                                // transaction). This ensures keys are cleaned up for both
3064                                // single-transaction purges and graveyard tombstones, while
3065                                // ensuring uncommitted transactions that get rolled back have no
3066                                // side effects on `key_manager`.
3067                                let _ = self.key_manager.remove(item.key.object_id);
3068                            }
3069                        }
3070                        if !context.mode.is_replay()
3071                            && matches!(
3072                                item.key.data,
3073                                ObjectKeyData::GraveyardEntry { .. }
3074                                    | ObjectKeyData::GraveyardAttributeEntry { .. }
3075                            )
3076                        {
3077                            if matches!(item.value, ObjectValue::Some | ObjectValue::Trim) {
3078                                self.graveyard_entries.fetch_add(1, Ordering::Relaxed);
3079                            } else if matches!(item.value, ObjectValue::None) {
3080                                self.graveyard_entries.fetch_sub(1, Ordering::Relaxed);
3081                            }
3082                        }
3083                        let lower_bound = item.key.key_for_merge_into();
3084                        self.tree.merge_into(item, &lower_bound);
3085                    }
3086                }
3087            }
3088            Mutation::BeginFlush => {
3089                ensure!(self.parent_store.is_some(), FxfsError::Inconsistent);
3090                self.tree.seal();
3091            }
3092            Mutation::EndFlush => ensure!(self.parent_store.is_some(), FxfsError::Inconsistent),
3093            Mutation::EncryptedObjectStore(_) | Mutation::UpdateMutationsKey(_) => {
3094                // We will process these during Self::unlock.
3095                ensure!(
3096                    !matches!(&*self.lock_state.lock(), LockState::Unencrypted),
3097                    FxfsError::Inconsistent
3098                );
3099            }
3100            Mutation::CreateInternalDir(object_id) => {
3101                ensure!(object_id != INVALID_OBJECT_ID, FxfsError::Inconsistent);
3102                self.store_info.lock().as_mut().unwrap().internal_directory_object_id = object_id;
3103            }
3104            _ => bail!("unexpected mutation: {:?}", mutation),
3105        }
3106        self.counters.lock().mutations_applied += 1;
3107        Ok(())
3108    }
3109
3110    fn drop_mutation(&self, mutation: Mutation, _transaction: &Transaction<'_>) {
3111        self.counters.lock().mutations_dropped += 1;
3112        if let Mutation::ObjectStore(ObjectStoreMutation {
3113            item: Item { key: ObjectKey { object_id, .. }, value: ObjectValue::Object { .. }, .. },
3114            op: Operation::Insert,
3115        }) = mutation
3116        {
3117            self.last_object_id.lock().unreserve(object_id);
3118        }
3119    }
3120
3121    async fn prepare_commit<'a>(
3122        &self,
3123        filesystem: &'a FxFilesystem,
3124        _transaction: &Transaction<'_>,
3125    ) -> Result<Option<WriteGuard<'a>>, Error> {
3126        // Short circuit check to see if this is an encrypted store.
3127        if !matches!(&*self.lock_state.lock(), LockState::Unlocked { .. }) {
3128            return Ok(None);
3129        }
3130
3131        // We must acquire the keys lock before we can access or modify `cached_keys`.  This guard
3132        // is returned and held until the transaction commits, ensuring that the keys we cache (or
3133        // existing keys) remain valid and are not interfered with by other transactions.
3134        let keys = lock_keys![LockKey::pre_cache_keys(self.store_object_id())];
3135        let guard = filesystem.lock_manager().write_lock(keys).await;
3136
3137        self.pre_cache_keys().await?;
3138
3139        Ok(Some(guard))
3140    }
3141
3142    /// Push all in-memory structures to the device. This is not necessary for sync since the
3143    /// journal will take care of it.  This is supposed to be called when there is either memory or
3144    /// space pressure (flushing the store will persist in-memory data and allow the journal file to
3145    /// be trimmed).
3146    ///
3147    /// Also returns the earliest version of a struct in the filesystem (when known).
3148    async fn flush(&self, reason: FlushReason) -> Result<Version, Error> {
3149        self.flush_with_reason(reason).await
3150    }
3151
3152    fn write_mutations(
3153        &self,
3154        mutations: ObjectMutationIterator<'_, '_>,
3155        mut writer: journal::Writer<'_>,
3156    ) {
3157        let mut lock_state = self.lock_state.lock();
3158        if let LockState::Unlocked { mutations_cipher, .. } = &mut *lock_state {
3159            let mut encrypted_transaction = EncryptedTransaction::new();
3160            for mutation in mutations.cloned() {
3161                // Intentionally enumerating all variants to force a decision on any new variants.
3162                // Encrypt all mutations that could affect an encrypted object store contents or
3163                // the `StoreInfo` of the encrypted object store. During `unlock()` any mutations
3164                // which haven't been encrypted won't be replayed after reading `StoreInfo`.
3165                match mutation {
3166                    // Whilst CreateInternalDir is a mutation for `StoreInfo`, which isn't
3167                    // encrypted, we still choose to encrypt the mutation because it makes it
3168                    // easier to deal with replay. When we replay mutations for an encrypted store,
3169                    // the only thing we keep in memory are the encrypted mutations; we don't keep
3170                    // `StoreInfo` or changes to it in memory. So, by encrypting the
3171                    // CreateInternalDir mutation here, it means we don't have to track both
3172                    // encrypted mutations bound for the LSM tree and unencrypted mutations for
3173                    // `StoreInfo` to use in `unlock()`. It'll just bundle CreateInternalDir
3174                    // mutations with the other encrypted mutations and handled them all in
3175                    // sequence during `unlock()`.
3176                    Mutation::ObjectStore(_) | Mutation::CreateInternalDir(_) => {
3177                        encrypted_transaction.0.push(mutation)
3178                    }
3179                    // `EncryptedObjectStore` and `UpdateMutationsKey` are both obviously
3180                    // associated with encrypted object stores, but are either the encrypted
3181                    // mutation data itself or metadata governing how the data will be encrypted.
3182                    // They should only be produced here.
3183                    Mutation::EncryptedObjectStore(_) | Mutation::UpdateMutationsKey(_) => {
3184                        debug_assert!(
3185                            false,
3186                            "Only this method should generate encrypted mutations"
3187                        );
3188                    }
3189                    // `BeginFlush` and `EndFlush` are not needed during `unlock()` and are needed
3190                    // during the initial journal replay, so should not be encrypted. `Allocator`,
3191                    // `DeleteVolume`, `UpdateBorrowed` mutations are never associated with an
3192                    // encrypted store as we do not encrypt the allocator or root/root-parent
3193                    // stores so we can avoid the locking.
3194                    Mutation::Allocator(_)
3195                    | Mutation::BeginFlush
3196                    | Mutation::EndFlush
3197                    | Mutation::DeleteVolume
3198                    | Mutation::UpdateBorrowed(_) => {
3199                        writer.write(mutation);
3200                    }
3201                }
3202            }
3203            if !encrypted_transaction.0.is_empty() {
3204                // If this is the first time we've used this key, we must write the key out.
3205                if mutations_cipher.key_is_new() {
3206                    writer.write(Mutation::update_mutations_key(
3207                        self.store_info
3208                            .lock()
3209                            .as_ref()
3210                            .unwrap()
3211                            .mutations_key
3212                            .as_ref()
3213                            .unwrap()
3214                            .clone(),
3215                    ));
3216                }
3217                for encrypted in encrypted_transaction.serialize_and_encrypt(mutations_cipher) {
3218                    writer.write(Mutation::EncryptedObjectStore(encrypted));
3219                }
3220            }
3221        } else {
3222            for mutation in mutations {
3223                writer.write(mutation.clone());
3224            }
3225        }
3226    }
3227}
3228
3229/// The plaintext serialization of the mutations for a transaction that are encrypted then wrapped
3230/// in `Mutation::EncryptedObjectStore`.
3231pub type EncryptedTransaction = EncryptedTransactionV59;
3232
3233static_assertions::const_assert!(EncryptedTransaction::MAX_CHUNK_SIZE % 16 == 0);
3234impl EncryptedTransaction {
3235    /// The maximum size of encrypted mutation data chunked into a single
3236    /// `Mutation::EncryptedObjectStore`. This must be a multiple of 16 bytes (AES block size) and
3237    /// must fit within `DEFAULT_MAX_SERIALIZED_RECORD_SIZE` after accounting for `JournalRecord`
3238    /// and `Mutation` serialization overhead.
3239    const MAX_CHUNK_SIZE: usize = DEFAULT_MAX_SERIALIZED_RECORD_SIZE as usize - 48;
3240
3241    fn new() -> Self {
3242        Self(Vec::new())
3243    }
3244
3245    fn serialize_and_encrypt<'a>(
3246        &self,
3247        cipher: &'a mut JournalCipher,
3248    ) -> impl Iterator<Item = Box<[u8]>> + 'a {
3249        let mut buffer = Vec::new();
3250        self.serialize_into(&mut buffer).unwrap();
3251        // Need to limit the size of the mutations. Cut them into chunks.
3252        (0..buffer.len()).step_by(Self::MAX_CHUNK_SIZE).map(move |offset| {
3253            let end = std::cmp::min(offset + Self::MAX_CHUNK_SIZE, buffer.len());
3254            Box::from(cipher.encrypt(&buffer[offset..end]))
3255        })
3256    }
3257
3258    fn decrypt_and_deserialize(
3259        cipher: &mut JournalXtsCipher,
3260        data: &mut [u8],
3261        version: Version,
3262    ) -> Result<Self, Error> {
3263        // Limited by maximum mutation size. So need to rebuild the transaction.
3264        let (sub_chunks, remainder) = data.as_chunks_mut::<{ Self::MAX_CHUNK_SIZE }>();
3265        for chunk in sub_chunks {
3266            cipher.decrypt(chunk.as_mut_slice());
3267        }
3268        if remainder.len() > 0 {
3269            cipher.decrypt(remainder);
3270        }
3271        let mut chunk_cursor = std::io::Cursor::new(&*data);
3272        Self::deserialize_from_version(&mut chunk_cursor, version)
3273            .context("failed to deserialize encrypted transaction")
3274    }
3275}
3276
3277// When this type is incremented, it must increment `Mutation` since this data will actually be
3278// wrapped in Mutation and the version associated with it will come from some parent type above
3279// Mutation.
3280#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
3281pub struct EncryptedTransactionV59(pub Vec<MutationV59>);
3282
3283impl TypeFingerprint for EncryptedTransactionV59 {
3284    fn fingerprint() -> String {
3285        format!(
3286            "struct{{MAX_CHUNK_SIZE: {}, {}}}",
3287            // The MAX_CHUNK_SIZE is an important part of the format.
3288            Self::MAX_CHUNK_SIZE,
3289            Vec::<MutationV59>::fingerprint()
3290        )
3291    }
3292}
3293
3294impl Versioned for EncryptedTransactionV59 {
3295    // This can serialize much larger sizes. They will be broken up before being wrapped in
3296    // EncryptedObjectStore.
3297    fn max_serialized_size() -> Option<u64> {
3298        None
3299    }
3300}
3301
3302#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
3303pub struct EncryptedTransactionV57(pub Vec<crate::object_store::transaction::MutationV57>);
3304
3305impl TypeFingerprint for EncryptedTransactionV57 {
3306    fn fingerprint() -> String {
3307        format!(
3308            "struct{{MAX_CHUNK_SIZE: {}, {}}}",
3309            EncryptedTransaction::MAX_CHUNK_SIZE,
3310            Vec::<crate::object_store::transaction::MutationV57>::fingerprint()
3311        )
3312    }
3313}
3314
3315impl Versioned for EncryptedTransactionV57 {
3316    fn max_serialized_size() -> Option<u64> {
3317        None
3318    }
3319}
3320
3321impl From<EncryptedTransactionV57> for EncryptedTransactionV59 {
3322    fn from(old: EncryptedTransactionV57) -> Self {
3323        Self(old.0.into_iter().map(Into::into).collect())
3324    }
3325}
3326
3327impl Drop for ObjectStore {
3328    fn drop(&mut self) {
3329        let mut last_object_id = self.last_object_id.lock();
3330        last_object_id.drain_unreserved();
3331        match &*last_object_id {
3332            LastObjectId::Low32Bit { reserved, .. } => debug_assert!(reserved.is_empty()),
3333            _ => {}
3334        }
3335    }
3336}
3337
3338impl HandleOwner for ObjectStore {}
3339
3340impl AsRef<ObjectStore> for ObjectStore {
3341    fn as_ref(&self) -> &ObjectStore {
3342        self
3343    }
3344}
3345
3346fn layer_size_from_encrypted_mutations_size(size: u64) -> u64 {
3347    // This is similar to reserved_space_from_journal_usage. It needs to be a worst case estimate of
3348    // the amount of metadata space that might need to be reserved to allow the encrypted mutations
3349    // to be written to layer files.  It needs to be >= than reservation_amount_from_layer_size will
3350    // return once the data has been written to layer files and <= than
3351    // reserved_space_from_journal_usage would use.  We can't just use
3352    // reserved_space_from_journal_usage because the encrypted mutations file includes some extra
3353    // data (it includes the checkpoints) that isn't written in the same way to the journal.
3354    size * 3
3355}
3356
3357impl AssociatedObject for ObjectStore {}
3358
3359/// Argument to the trim_some method.
3360#[derive(Debug)]
3361pub enum TrimMode {
3362    /// Trim extents beyond the current size.
3363    UseSize,
3364
3365    /// Trim extents beyond the supplied offset.
3366    FromOffset(u64),
3367
3368    /// Remove the object (or attribute) from the store once it is fully trimmed.
3369    Tombstone(TombstoneMode),
3370}
3371
3372/// Sets the mode for tombstoning (either at the object or attribute level).
3373#[derive(Debug)]
3374pub enum TombstoneMode {
3375    Object,
3376    Attribute,
3377}
3378
3379/// Result of the trim_some method.
3380#[derive(Debug)]
3381pub enum TrimResult {
3382    /// We reached the limit of the transaction and more extents might follow.
3383    Incomplete,
3384
3385    /// We finished this attribute.  Returns the ID of the next attribute for the same object if
3386    /// there is one.
3387    Done(Option<AttributeId>),
3388}
3389
3390/// Loads store info.
3391pub async fn load_store_info(
3392    parent: &Arc<ObjectStore>,
3393    store_object_id: u64,
3394) -> Result<StoreInfo, Error> {
3395    load_store_info_from_handle(
3396        &ObjectStore::open_object(parent, store_object_id, HandleOptions::default(), None).await?,
3397    )
3398    .await
3399}
3400
3401async fn load_store_info_from_handle(
3402    handle: &DataObjectHandle<impl HandleOwner>,
3403) -> Result<StoreInfo, Error> {
3404    Ok(if handle.get_size() > 0 {
3405        let serialized_info = handle.contents(MAX_STORE_INFO_SERIALIZED_SIZE).await?;
3406        let mut cursor = std::io::Cursor::new(serialized_info);
3407        let (store_info, _) = StoreInfo::deserialize_with_version(&mut cursor)
3408            .context("Failed to deserialize StoreInfo")?;
3409        store_info
3410    } else {
3411        // The store_info will be absent for a newly created and empty object store.
3412        StoreInfo::default()
3413    })
3414}
3415
3416#[cfg(test)]
3417mod tests {
3418    use super::{
3419        AttributeId, DirectWriter, EncryptedMutations, FsverityMetadata, HandleOptions,
3420        LastObjectId, LastObjectIdInfo, LockKey, LockState, MAX_STORE_INFO_SERIALIZED_SIZE,
3421        Mutation, NewChildStoreOptions, OBJECT_ID_HI_MASK, ObjectEncryptionOptions, ObjectStore,
3422        RootDigest, StoreInfo, StoreOptions,
3423    };
3424    use crate::errors::FxfsError;
3425    use crate::filesystem::{
3426        FlushReason, ForceMajor, FxFilesystem, FxFilesystemBuilder, MAX_IN_FLIGHT_TRANSACTIONS,
3427        OpenFxFilesystem,
3428    };
3429    use crate::fsck::{fsck, fsck_volume};
3430    use crate::hooks::{Hooks, HooksHandle};
3431    use crate::lsm_tree::Query;
3432    use crate::lsm_tree::types::{ItemRef, LayerIterator};
3433    use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, WriteBytes, WriteObjectHandle};
3434    use crate::object_store::directory::{Directory, replace_child};
3435    use crate::object_store::journal::{JournalCheckpoint, JournalOptions};
3436    use crate::object_store::object_record::{
3437        AttributeKey, ObjectDescriptor, ObjectKey, ObjectKind, ObjectValue,
3438    };
3439    use crate::object_store::store_object_handle::MAX_INLINE_XATTR_SIZE;
3440    use crate::object_store::transaction::{Options, lock_keys};
3441    use crate::object_store::volume::root_volume;
3442    use crate::serialized_types::{
3443        DEFAULT_MAX_SERIALIZED_RECORD_SIZE, Version, Versioned, VersionedLatest,
3444    };
3445    use crate::testing;
3446    use assert_matches::assert_matches;
3447    use async_trait::async_trait;
3448    use fuchsia_async as fasync;
3449    use fuchsia_sync::Mutex;
3450    use futures::channel::oneshot;
3451    use futures::{FutureExt, join};
3452    use fxfs_crypto::ff1::Ff1;
3453    use fxfs_crypto::{
3454        Crypt, EncryptionKey, FXFS_KEY_SIZE, FXFS_WRAPPED_KEY_SIZE, FxfsKey, JournalCipher,
3455        KeyPurpose, ObjectType, StreamCipher, UnwrappedKey, WrappedKey, WrappedKeyBytes,
3456        WrappingKeyId,
3457    };
3458    use fxfs_insecure_crypto::new_insecure_crypt;
3459    use std::sync::Arc;
3460    use std::sync::atomic::{AtomicBool, AtomicIsize, Ordering};
3461    use std::time::Duration;
3462    use storage_device::DeviceHolder;
3463    use storage_device::fake_device::FakeDevice;
3464    use test_case::test_case;
3465    use zx_status as zx;
3466
3467    const TEST_DEVICE_BLOCK_SIZE: u32 = 512;
3468
3469    async fn test_filesystem_with_hooks(hooks: Arc<HooksHandle>) -> OpenFxFilesystem {
3470        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
3471        FxFilesystemBuilder::new()
3472            .hooks(hooks)
3473            .format(true)
3474            .open(device)
3475            .await
3476            .expect("open failed")
3477    }
3478
3479    async fn test_filesystem() -> OpenFxFilesystem {
3480        test_filesystem_with_hooks(Default::default()).await
3481    }
3482
3483    #[fuchsia::test]
3484    async fn test_verified_file_with_verified_attribute() {
3485        let fs: OpenFxFilesystem = test_filesystem().await;
3486        let mut transaction = fs
3487            .root_store()
3488            .new_transaction(lock_keys![], Options::default())
3489            .await
3490            .expect("new_transaction failed");
3491        let store = fs.root_store();
3492        let object = Arc::new(
3493            ObjectStore::create_object(&store, &mut transaction, HandleOptions::default(), None)
3494                .await
3495                .expect("create_object failed"),
3496        );
3497
3498        transaction.add(
3499            store.store_object_id(),
3500            Mutation::replace_or_insert_object(
3501                ObjectKey::attribute(
3502                    object.object_id(),
3503                    AttributeId::DATA,
3504                    AttributeKey::Attribute,
3505                ),
3506                ObjectValue::verified_attribute(
3507                    0,
3508                    FsverityMetadata::Internal(RootDigest::Sha256([0; 32]), vec![]),
3509                ),
3510            ),
3511        );
3512
3513        transaction.add(
3514            store.store_object_id(),
3515            Mutation::replace_or_insert_object(
3516                ObjectKey::attribute(
3517                    object.object_id(),
3518                    AttributeId::FSVERITY_MERKLE,
3519                    AttributeKey::Attribute,
3520                ),
3521                ObjectValue::attribute(0, false),
3522            ),
3523        );
3524
3525        transaction.commit().await.unwrap();
3526
3527        let handle =
3528            ObjectStore::open_object(&store, object.object_id(), HandleOptions::default(), None)
3529                .await
3530                .expect("open_object failed");
3531
3532        assert!(handle.is_verified_file());
3533
3534        fs.close().await.expect("Close failed");
3535    }
3536
3537    #[fuchsia::test]
3538    async fn test_verified_file_without_verified_attribute() {
3539        let fs: OpenFxFilesystem = test_filesystem().await;
3540        let mut transaction = fs
3541            .root_store()
3542            .new_transaction(lock_keys![], Options::default())
3543            .await
3544            .expect("new_transaction failed");
3545        let store = fs.root_store();
3546        let object = Arc::new(
3547            ObjectStore::create_object(&store, &mut transaction, HandleOptions::default(), None)
3548                .await
3549                .expect("create_object failed"),
3550        );
3551
3552        transaction.commit().await.unwrap();
3553
3554        let handle =
3555            ObjectStore::open_object(&store, object.object_id(), HandleOptions::default(), None)
3556                .await
3557                .expect("open_object failed");
3558
3559        assert!(!handle.is_verified_file());
3560
3561        fs.close().await.expect("Close failed");
3562    }
3563
3564    #[fuchsia::test]
3565    async fn test_create_and_open_store() {
3566        let fs = test_filesystem().await;
3567        let store_id = {
3568            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
3569            root_volume
3570                .new_volume(
3571                    "test",
3572                    NewChildStoreOptions {
3573                        options: StoreOptions { crypt: Some(Arc::new(new_insecure_crypt())) },
3574                        ..Default::default()
3575                    },
3576                )
3577                .await
3578                .expect("new_volume failed")
3579                .store_object_id()
3580        };
3581
3582        fs.close().await.expect("close failed");
3583        let device = fs.take_device().await;
3584        device.reopen(false);
3585        let fs = FxFilesystem::open(device).await.expect("open failed");
3586
3587        {
3588            let store = fs.object_manager().store(store_id).expect("store not found");
3589            store.unlock(Arc::new(new_insecure_crypt())).await.expect("unlock failed");
3590        }
3591        fs.close().await.expect("Close failed");
3592    }
3593
3594    #[fuchsia::test]
3595    async fn test_create_and_open_internal_dir() {
3596        let fs = test_filesystem().await;
3597        let dir_id;
3598        let store_id;
3599        {
3600            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
3601            let store = root_volume
3602                .new_volume(
3603                    "test",
3604                    NewChildStoreOptions {
3605                        options: StoreOptions { crypt: Some(Arc::new(new_insecure_crypt())) },
3606                        ..Default::default()
3607                    },
3608                )
3609                .await
3610                .expect("new_volume failed");
3611            dir_id =
3612                store.get_or_create_internal_directory_id().await.expect("Create internal dir");
3613            store_id = store.store_object_id();
3614        }
3615
3616        fs.close().await.expect("close failed");
3617        let device = fs.take_device().await;
3618        device.reopen(false);
3619        let fs = FxFilesystem::open(device).await.expect("open failed");
3620
3621        {
3622            let store = fs.object_manager().store(store_id).expect("store not found");
3623            store.unlock(Arc::new(new_insecure_crypt())).await.expect("unlock failed");
3624            assert_eq!(
3625                dir_id,
3626                store.get_or_create_internal_directory_id().await.expect("Retrieving dir")
3627            );
3628            let obj = store
3629                .tree()
3630                .find_value(&ObjectKey::object(dir_id))
3631                .await
3632                .expect("Searching tree for dir")
3633                .unwrap();
3634            assert_matches!(obj, ObjectValue::Object { kind: ObjectKind::Directory { .. }, .. });
3635        }
3636        fs.close().await.expect("Close failed");
3637    }
3638
3639    #[fuchsia::test]
3640    async fn test_create_and_open_internal_dir_unencrypted() {
3641        let fs = test_filesystem().await;
3642        let dir_id;
3643        let store_id;
3644        {
3645            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
3646            let store = root_volume
3647                .new_volume("test", NewChildStoreOptions::default())
3648                .await
3649                .expect("new_volume failed");
3650            dir_id =
3651                store.get_or_create_internal_directory_id().await.expect("Create internal dir");
3652            store_id = store.store_object_id();
3653        }
3654
3655        fs.close().await.expect("close failed");
3656        let device = fs.take_device().await;
3657        device.reopen(false);
3658        let fs = FxFilesystem::open(device).await.expect("open failed");
3659
3660        {
3661            let store = fs.object_manager().store(store_id).expect("store not found");
3662            assert_eq!(
3663                dir_id,
3664                store.get_or_create_internal_directory_id().await.expect("Retrieving dir")
3665            );
3666            let obj = store
3667                .tree()
3668                .find_value(&ObjectKey::object(dir_id))
3669                .await
3670                .expect("Searching tree for dir")
3671                .unwrap();
3672            assert_matches!(obj, ObjectValue::Object { kind: ObjectKind::Directory { .. }, .. });
3673        }
3674        fs.close().await.expect("Close failed");
3675    }
3676
3677    #[fuchsia::test(threads = 10)]
3678    async fn test_old_layers_are_purged() {
3679        let fs = test_filesystem().await;
3680
3681        let store = fs.root_store();
3682        let mut transaction = fs
3683            .root_store()
3684            .new_transaction(lock_keys![], Options::default())
3685            .await
3686            .expect("new_transaction failed");
3687        let object = Arc::new(
3688            ObjectStore::create_object(&store, &mut transaction, HandleOptions::default(), None)
3689                .await
3690                .expect("create_object failed"),
3691        );
3692        transaction.commit().await.expect("commit failed");
3693
3694        store.flush().await.expect("flush failed");
3695
3696        let mut buf = object.allocate_buffer(5).await;
3697        buf.copy_from_slice(b"hello");
3698        object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
3699
3700        // Getting the layer-set should cause the flush to stall.
3701        let layer_set = store.tree().layer_set();
3702
3703        let done = Mutex::new(false);
3704        let mut object_id = 0;
3705
3706        join!(
3707            async {
3708                store.flush().await.expect("flush failed");
3709                assert!(*done.lock());
3710            },
3711            async {
3712                // This is a halting problem so all we can do is sleep.
3713                fasync::Timer::new(Duration::from_secs(1)).await;
3714                *done.lock() = true;
3715                object_id = layer_set.layers.last().unwrap().handle().unwrap().object_id();
3716                std::mem::drop(layer_set);
3717            }
3718        );
3719
3720        if let Err(e) = ObjectStore::open_object(
3721            &store.parent_store.as_ref().unwrap(),
3722            object_id,
3723            HandleOptions::default(),
3724            store.crypt(),
3725        )
3726        .await
3727        {
3728            assert!(FxfsError::NotFound.matches(&e));
3729        } else {
3730            panic!("open_object succeeded");
3731        }
3732    }
3733
3734    #[fuchsia::test]
3735    async fn test_tombstone_deletes_data() {
3736        let fs = test_filesystem().await;
3737        let root_store = fs.root_store();
3738        let child_id = {
3739            let mut transaction = fs
3740                .root_store()
3741                .new_transaction(lock_keys![], Options::default())
3742                .await
3743                .expect("new_transaction failed");
3744            let child = ObjectStore::create_object(
3745                &root_store,
3746                &mut transaction,
3747                HandleOptions::default(),
3748                None,
3749            )
3750            .await
3751            .expect("create_object failed");
3752            root_store.add_to_graveyard(&mut transaction, child.object_id());
3753            transaction.commit().await.expect("commit failed");
3754
3755            // Allocate an extent in the file.
3756            let mut buffer = child.allocate_buffer(8192).await;
3757            buffer.fill(0xaa);
3758            child.write_or_append(Some(0), buffer.as_ref()).await.expect("write failed");
3759
3760            child.object_id()
3761        };
3762
3763        root_store
3764            .tombstone_object(child_id, Options::default(), None)
3765            .await
3766            .expect("tombstone failed");
3767
3768        // Let fsck check allocations.
3769        fsck(fs.clone()).await.expect("fsck failed");
3770    }
3771
3772    #[fuchsia::test]
3773    async fn test_tombstone_purges_keys() {
3774        let fs = test_filesystem().await;
3775        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
3776        let store = root_volume
3777            .new_volume(
3778                "test",
3779                NewChildStoreOptions {
3780                    options: StoreOptions {
3781                        crypt: Some(Arc::new(new_insecure_crypt())),
3782                        ..StoreOptions::default()
3783                    },
3784                    ..NewChildStoreOptions::default()
3785                },
3786            )
3787            .await
3788            .expect("new_volume failed");
3789        let mut transaction = fs
3790            .root_store()
3791            .new_transaction(lock_keys![], Options::default())
3792            .await
3793            .expect("new_transaction failed");
3794        let child =
3795            ObjectStore::create_object(&store, &mut transaction, HandleOptions::default(), None)
3796                .await
3797                .expect("create_object failed");
3798        store.add_to_graveyard(&mut transaction, child.object_id());
3799        transaction.commit().await.expect("commit failed");
3800        assert!(store.key_manager.get(child.object_id()).await.unwrap().is_some());
3801        store
3802            .tombstone_object(child.object_id(), Options::default(), None)
3803            .await
3804            .expect("tombstone_object failed");
3805        assert!(store.key_manager.get(child.object_id()).await.unwrap().is_none());
3806        fs.close().await.expect("close failed");
3807    }
3808
3809    #[fuchsia::test]
3810    async fn test_major_compaction_discards_unnecessary_records() {
3811        let fs = test_filesystem().await;
3812        let root_store = fs.root_store();
3813        let child_id = {
3814            let mut transaction = fs
3815                .root_store()
3816                .new_transaction(lock_keys![], Options::default())
3817                .await
3818                .expect("new_transaction failed");
3819            let child = ObjectStore::create_object(
3820                &root_store,
3821                &mut transaction,
3822                HandleOptions::default(),
3823                None,
3824            )
3825            .await
3826            .expect("create_object failed");
3827            root_store.add_to_graveyard(&mut transaction, child.object_id());
3828            transaction.commit().await.expect("commit failed");
3829
3830            // Allocate an extent in the file.
3831            let mut buffer = child.allocate_buffer(8192).await;
3832            buffer.fill(0xaa);
3833            child.write_or_append(Some(0), buffer.as_ref()).await.expect("write failed");
3834
3835            child.object_id()
3836        };
3837
3838        root_store
3839            .tombstone_object(child_id, Options::default(), None)
3840            .await
3841            .expect("tombstone failed");
3842        {
3843            let layers = root_store.tree.layer_set();
3844            let mut merger = layers.merger();
3845            let iter = merger
3846                .query(Query::FullRange(&ObjectKey::object(child_id)))
3847                .await
3848                .expect("seek failed");
3849            // Find at least one object still in the tree.
3850            match iter.get() {
3851                Some(ItemRef { key: ObjectKey { object_id, .. }, .. })
3852                    if *object_id == child_id => {}
3853                _ => panic!("Objects should still be in the tree."),
3854            }
3855        }
3856        root_store.flush().await.expect("flush failed");
3857
3858        // There should be no records for the object.
3859        let layers = root_store.tree.layer_set();
3860        let mut merger = layers.merger();
3861        let iter = merger
3862            .query(Query::FullRange(&ObjectKey::object(child_id)))
3863            .await
3864            .expect("seek failed");
3865        match iter.get() {
3866            None => {}
3867            Some(ItemRef { key: ObjectKey { object_id, .. }, .. }) => {
3868                assert_ne!(*object_id, child_id)
3869            }
3870        }
3871    }
3872
3873    #[fuchsia::test]
3874    async fn test_overlapping_extents_in_different_layers() {
3875        let fs = test_filesystem().await;
3876        let store = fs.root_store();
3877
3878        let mut transaction = store
3879            .new_transaction(
3880                lock_keys![LockKey::object(
3881                    store.store_object_id(),
3882                    store.root_directory_object_id()
3883                )],
3884                Options::default(),
3885            )
3886            .await
3887            .expect("new_transaction failed");
3888        let root_directory =
3889            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
3890        let object = root_directory
3891            .create_child_file(&mut transaction, "test")
3892            .await
3893            .expect("create_child_file failed");
3894        transaction.commit().await.expect("commit failed");
3895
3896        let buf = object.allocate_buffer(16384).await;
3897        object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
3898
3899        store.flush().await.expect("flush failed");
3900
3901        object.write_or_append(Some(0), buf.subslice(0..4096)).await.expect("write failed");
3902
3903        // At this point, we should have an extent for 0..16384 in a layer that has been flushed,
3904        // and an extent for 0..4096 that partially overwrites it.  Writing to 0..16384 should
3905        // overwrite both of those extents.
3906        object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
3907
3908        fsck(fs.clone()).await.expect("fsck failed");
3909    }
3910
3911    #[fuchsia::test(threads = 10)]
3912    async fn test_encrypted_mutations() {
3913        async fn one_iteration(
3914            fs: OpenFxFilesystem,
3915            crypt: Arc<dyn Crypt>,
3916            iteration: u64,
3917        ) -> OpenFxFilesystem {
3918            async fn reopen(fs: OpenFxFilesystem) -> OpenFxFilesystem {
3919                fs.close().await.expect("Close failed");
3920                let device = fs.take_device().await;
3921                device.reopen(false);
3922                FxFilesystem::open(device).await.expect("FS open failed")
3923            }
3924
3925            let fs = reopen(fs).await;
3926
3927            let (store_object_id, object_id) = {
3928                let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
3929                let store = root_volume
3930                    .volume(
3931                        "test",
3932                        StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
3933                    )
3934                    .await
3935                    .expect("volume failed");
3936
3937                let mut transaction = fs
3938                    .root_store()
3939                    .new_transaction(
3940                        lock_keys![LockKey::object(
3941                            store.store_object_id(),
3942                            store.root_directory_object_id(),
3943                        )],
3944                        Options::default(),
3945                    )
3946                    .await
3947                    .expect("new_transaction failed");
3948                let root_directory = Directory::open(&store, store.root_directory_object_id())
3949                    .await
3950                    .expect("open failed");
3951                let object = root_directory
3952                    .create_child_file(&mut transaction, &format!("test {}", iteration))
3953                    .await
3954                    .expect("create_child_file failed");
3955                transaction.commit().await.expect("commit failed");
3956
3957                let mut buf = object.allocate_buffer(1000).await;
3958                for (i, byte) in buf.as_mut_ptr_slice().iter_as_mut::<u8>().enumerate() {
3959                    byte.write(i as u8);
3960                }
3961                object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
3962
3963                (store.store_object_id(), object.object_id())
3964            };
3965
3966            let fs = reopen(fs).await;
3967
3968            let check_object = |fs: Arc<FxFilesystem>| {
3969                let crypt = crypt.clone();
3970                async move {
3971                    let root_volume = root_volume(fs).await.expect("root_volume failed");
3972                    let volume = root_volume
3973                        .volume(
3974                            "test",
3975                            StoreOptions { crypt: Some(crypt), ..StoreOptions::default() },
3976                        )
3977                        .await
3978                        .expect("volume failed");
3979
3980                    let object = ObjectStore::open_object(
3981                        &volume,
3982                        object_id,
3983                        HandleOptions::default(),
3984                        None,
3985                    )
3986                    .await
3987                    .expect("open_object failed");
3988                    let buf = object.read_bytes(0..1000).await.expect("read failed");
3989                    assert_eq!(buf.len(), 1000);
3990                    for (i, &byte) in buf.iter().enumerate() {
3991                        assert_eq!(byte, i as u8);
3992                    }
3993                }
3994            };
3995
3996            check_object(fs.clone()).await;
3997
3998            let fs = reopen(fs).await;
3999
4000            // At this point the "test" volume is locked.  Before checking the object, flush the
4001            // filesystem.  This should leave a file with encrypted mutations.
4002            fs.object_manager()
4003                .flush(FlushReason::Journal(ForceMajor::False))
4004                .await
4005                .expect("flush failed");
4006
4007            assert_ne!(
4008                fs.object_manager()
4009                    .store(store_object_id)
4010                    .unwrap()
4011                    .load_store_info()
4012                    .await
4013                    .expect("load_store_info failed")
4014                    .encrypted_mutations_object_id,
4015                INVALID_OBJECT_ID
4016            );
4017
4018            check_object(fs.clone()).await;
4019
4020            // Checking the object should have triggered a flush and so now there should be no
4021            // encrypted mutations object.
4022            assert_eq!(
4023                fs.object_manager()
4024                    .store(store_object_id)
4025                    .unwrap()
4026                    .load_store_info()
4027                    .await
4028                    .expect("load_store_info failed")
4029                    .encrypted_mutations_object_id,
4030                INVALID_OBJECT_ID
4031            );
4032
4033            let fs = reopen(fs).await;
4034
4035            fsck(fs.clone()).await.expect("fsck failed");
4036
4037            let fs = reopen(fs).await;
4038
4039            check_object(fs.clone()).await;
4040
4041            fs
4042        }
4043
4044        let mut fs = test_filesystem().await;
4045        let crypt = Arc::new(new_insecure_crypt());
4046
4047        {
4048            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4049            let _store = root_volume
4050                .new_volume(
4051                    "test",
4052                    NewChildStoreOptions {
4053                        options: StoreOptions {
4054                            crypt: Some(crypt.clone()),
4055                            ..StoreOptions::default()
4056                        },
4057                        ..Default::default()
4058                    },
4059                )
4060                .await
4061                .expect("new_volume failed");
4062        }
4063
4064        // Run a few iterations so that we test changes with the stream cipher offset.
4065        for i in 0..5 {
4066            fs = one_iteration(fs, crypt.clone(), i).await;
4067        }
4068    }
4069
4070    #[test_case(true; "with a flush")]
4071    #[test_case(false; "without a flush")]
4072    #[fuchsia::test(threads = 10)]
4073    async fn test_object_id_cipher_roll(with_flush: bool) {
4074        let fs = test_filesystem().await;
4075        let crypt = Arc::new(new_insecure_crypt());
4076
4077        let expected_key = {
4078            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4079            let store = root_volume
4080                .new_volume(
4081                    "test",
4082                    NewChildStoreOptions {
4083                        options: StoreOptions {
4084                            crypt: Some(crypt.clone()),
4085                            ..StoreOptions::default()
4086                        },
4087                        ..Default::default()
4088                    },
4089                )
4090                .await
4091                .expect("new_volume failed");
4092
4093            // Create some files so that our in-memory copy of StoreInfo has changes (the object
4094            // count) pending a flush.
4095            let root_dir_id = store.root_directory_object_id();
4096            let root_dir =
4097                Arc::new(Directory::open(&store, root_dir_id).await.expect("open failed"));
4098            let mut transaction = store
4099                .new_transaction(
4100                    lock_keys![LockKey::object(store.store_object_id(), root_dir_id)],
4101                    Options::default(),
4102                )
4103                .await
4104                .expect("new_transaction failed");
4105            for i in 0..10 {
4106                root_dir.create_child_file(&mut transaction, &format!("file {i}")).await.unwrap();
4107            }
4108            transaction.commit().await.expect("commit failed");
4109
4110            let orig_store_info = store.store_info().unwrap();
4111
4112            // Hack the last object ID to force a roll of the object ID cipher.
4113            {
4114                let mut last_object_id = store.last_object_id.lock();
4115                match &mut *last_object_id {
4116                    LastObjectId::Encrypted { id, .. } => {
4117                        assert_eq!(*id & OBJECT_ID_HI_MASK, 0);
4118                        *id |= 0xffffffff;
4119                    }
4120                    _ => unreachable!(),
4121                }
4122            }
4123
4124            let mut transaction = store
4125                .new_transaction(
4126                    lock_keys![LockKey::object(
4127                        store.store_object_id(),
4128                        store.root_directory_object_id()
4129                    )],
4130                    Options::default(),
4131                )
4132                .await
4133                .expect("new_transaction failed");
4134            let root_directory = Directory::open(&store, store.root_directory_object_id())
4135                .await
4136                .expect("open failed");
4137            let object = root_directory
4138                .create_child_file(&mut transaction, "test")
4139                .await
4140                .expect("create_child_file failed");
4141            transaction.commit().await.expect("commit failed");
4142
4143            assert_eq!(object.object_id() & OBJECT_ID_HI_MASK, 1u64 << 32);
4144
4145            // Check that the key has been changed.
4146            let key = match (
4147                store.store_info().unwrap().last_object_id,
4148                orig_store_info.last_object_id,
4149            ) {
4150                (
4151                    LastObjectIdInfo::Encrypted { key, id },
4152                    LastObjectIdInfo::Encrypted { key: orig_key, .. },
4153                ) => {
4154                    assert_ne!(key, orig_key);
4155                    assert_eq!(id, 1u64 << 32);
4156                    key
4157                }
4158                _ => unreachable!(),
4159            };
4160
4161            if with_flush {
4162                fs.journal().force_compact().await.unwrap();
4163            }
4164
4165            let last_object_id = store.last_object_id.lock();
4166            assert_eq!(last_object_id.id(), 1u64 << 32);
4167            key
4168        };
4169
4170        fs.close().await.expect("Close failed");
4171        let device = fs.take_device().await;
4172        device.reopen(false);
4173        let fs = FxFilesystem::open(device).await.expect("open failed");
4174        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4175        let store = root_volume
4176            .volume("test", StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() })
4177            .await
4178            .expect("volume failed");
4179
4180        assert_matches!(store.store_info().unwrap().last_object_id, LastObjectIdInfo::Encrypted { key, .. } if key == expected_key);
4181        assert_eq!(store.last_object_id.lock().id(), 1u64 << 32);
4182
4183        fsck(fs.clone()).await.expect("fsck failed");
4184        fsck_volume(&fs, store.store_object_id(), None).await.expect("fsck_volume failed");
4185    }
4186
4187    #[fuchsia::test]
4188    async fn test_object_id_cipher_roll_no_space() {
4189        let fs = test_filesystem().await;
4190        let crypt = Arc::new(new_insecure_crypt());
4191
4192        {
4193            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4194            let store = root_volume
4195                .new_volume(
4196                    "test",
4197                    NewChildStoreOptions {
4198                        options: StoreOptions {
4199                            crypt: Some(crypt.clone()),
4200                            ..StoreOptions::default()
4201                        },
4202                        ..Default::default()
4203                    },
4204                )
4205                .await
4206                .expect("new_volume failed");
4207
4208            let root_directory = Directory::open(&store, store.root_directory_object_id())
4209                .await
4210                .expect("open failed");
4211
4212            // Force the next object ID allocation to roll the object ID cipher.
4213            match &mut *store.last_object_id.lock() {
4214                LastObjectId::Encrypted { id, .. } => {
4215                    *id |= 0xffffffff;
4216                }
4217                _ => unreachable!(),
4218            }
4219
4220            let mut transaction = store
4221                .new_transaction(
4222                    lock_keys![LockKey::object(
4223                        store.store_object_id(),
4224                        store.root_directory_object_id()
4225                    )],
4226                    Options::default(),
4227                )
4228                .await
4229                .expect("new_transaction failed");
4230
4231            // Reserve all remaining free space in the allocator so that rolling the object ID
4232            // cipher fails with NoSpace rather than borrowing metadata space.
4233            let reservation = fs.allocator().reserve_with(None, |limit| limit);
4234            let err = root_directory
4235                .create_child_file(&mut transaction, "test")
4236                .await
4237                .err()
4238                .expect("create_child_file should fail with NoSpace");
4239            assert!(FxfsError::NoSpace.matches(&err), "unexpected error: {err:?}");
4240
4241            std::mem::drop(reservation);
4242
4243            let object = root_directory
4244                .create_child_file(&mut transaction, "test")
4245                .await
4246                .expect("create_child_file failed");
4247            transaction.commit().await.expect("commit failed");
4248
4249            assert_eq!(object.object_id() & OBJECT_ID_HI_MASK, 1u64 << 32);
4250        }
4251
4252        fsck(fs.clone()).await.expect("fsck failed");
4253        fs.close().await.expect("Close failed");
4254    }
4255
4256    #[fuchsia::test]
4257    async fn test_object_id_cipher_roll_checks_journal_space() {
4258        let reclaim_size = 65536;
4259        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
4260        let (mut hooks, fs_hooks) = Hooks::new();
4261        let fs = FxFilesystemBuilder::new()
4262            .hooks(fs_hooks)
4263            .journal_options(JournalOptions { reclaim_size, ..Default::default() })
4264            .format(true)
4265            .open(device)
4266            .await
4267            .expect("open failed");
4268
4269        {
4270            let crypt = Arc::new(new_insecure_crypt());
4271            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4272            let store = root_volume
4273                .new_volume(
4274                    "test",
4275                    NewChildStoreOptions {
4276                        options: StoreOptions { crypt: Some(crypt), ..StoreOptions::default() },
4277                        ..Default::default()
4278                    },
4279                )
4280                .await
4281                .expect("new_volume failed");
4282            let root_directory = Directory::open(&store, store.root_directory_object_id())
4283                .await
4284                .expect("open failed");
4285
4286            let write_xattrs = || async {
4287                for _ in 0..4 {
4288                    let mut t = store
4289                        .new_transaction(
4290                            lock_keys![LockKey::object(store.store_object_id(), 1000)],
4291                            Options::default(),
4292                        )
4293                        .await
4294                        .expect("new_transaction failed");
4295                    for i in 0..18 {
4296                        t.add(
4297                            store.store_object_id(),
4298                            Mutation::replace_or_insert_object(
4299                                ObjectKey::extended_attribute(
4300                                    1000,
4301                                    format!("attr_{i}").into_bytes(),
4302                                ),
4303                                ObjectValue::inline_extended_attribute(vec![0u8; 1024]),
4304                            ),
4305                        );
4306                    }
4307                    t.commit().await.expect("commit failed");
4308                }
4309            };
4310
4311            // Warm up the store and journal so initial borrowed metadata space is paid down.
4312            write_xattrs().await;
4313            fs.journal().force_compact().await.expect("force_compact failed");
4314            fs.journal().force_compact().await.expect("force_compact failed");
4315            fs.journal().pause_compactions().await;
4316
4317            // Start the outer transaction before the journal fills up.
4318            let mut transaction = store
4319                .new_transaction(
4320                    lock_keys![LockKey::object(
4321                        store.store_object_id(),
4322                        store.root_directory_object_id()
4323                    )],
4324                    Options::default(),
4325                )
4326                .await
4327                .expect("new_transaction failed");
4328
4329            // Fill the journal past `reclaim_size` while compactions are paused so that space is
4330            // tied up in the metadata reservation for uncompacted journal mutations.
4331            write_xattrs().await;
4332
4333            let journal_clone = fs.journal().clone();
4334            hooks.set_waiting_for_journal_space(move || {
4335                journal_clone.resume_compactions();
4336            });
4337
4338            // Reserve all remaining free space in the allocator so that rolling the object ID
4339            // cipher can only succeed if the nested transaction waits for journal compaction to
4340            // reclaim space.
4341            let _reservation = fs.allocator().reserve_with(None, |limit| limit);
4342
4343            // Force the next object ID allocation to roll the object ID cipher.
4344            match &mut *store.last_object_id.lock() {
4345                LastObjectId::Encrypted { id, .. } => *id |= 0xffffffff,
4346                _ => unreachable!(),
4347            }
4348
4349            let object = root_directory
4350                .create_child_file(&mut transaction, "test")
4351                .await
4352                .expect("create_child_file failed");
4353            transaction.commit().await.expect("commit failed");
4354
4355            assert_eq!(object.object_id() & OBJECT_ID_HI_MASK, 1u64 << 32);
4356        }
4357
4358        fsck(fs.clone()).await.expect("fsck failed");
4359        fs.close().await.expect("Close failed");
4360    }
4361
4362    #[test_case(false; "cipher_roll")]
4363    #[test_case(true; "low_32_bit")]
4364    #[fuchsia::test]
4365    async fn test_get_next_object_id_with_max_in_flight_transactions(low_32_bit_object_ids: bool) {
4366        let fs = test_filesystem().await;
4367        let crypt = Arc::new(new_insecure_crypt());
4368
4369        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4370        let store = root_volume
4371            .new_volume(
4372                "test",
4373                NewChildStoreOptions {
4374                    options: StoreOptions { crypt: Some(crypt), ..StoreOptions::default() },
4375                    low_32_bit_object_ids,
4376                    ..Default::default()
4377                },
4378            )
4379            .await
4380            .expect("new_volume failed");
4381
4382        if !low_32_bit_object_ids {
4383            // Force the next object ID allocation to roll the object ID cipher.
4384            match &mut *store.last_object_id.lock() {
4385                LastObjectId::Encrypted { id, .. } => *id |= 0xffffffff,
4386                _ => unreachable!(),
4387            }
4388        }
4389
4390        let mut transactions = Vec::new();
4391        for _ in 0..MAX_IN_FLIGHT_TRANSACTIONS {
4392            transactions.push(
4393                store
4394                    .new_transaction(lock_keys![], Options::default())
4395                    .await
4396                    .expect("new_transaction failed"),
4397            );
4398        }
4399
4400        ObjectStore::create_object(
4401            &store,
4402            transactions.last_mut().unwrap(),
4403            HandleOptions::default(),
4404            None,
4405        )
4406        .await
4407        .expect("create_object failed");
4408
4409        for transaction in transactions {
4410            transaction.commit().await.expect("commit failed");
4411        }
4412
4413        fs.close().await.expect("Close failed");
4414    }
4415
4416    #[fuchsia::test(threads = 2)]
4417    async fn test_race_object_id_cipher_roll_and_flush() {
4418        let fs = test_filesystem().await;
4419        let crypt = Arc::new(new_insecure_crypt());
4420
4421        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4422        let store = root_volume
4423            .new_volume(
4424                "test",
4425                NewChildStoreOptions {
4426                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4427                    ..Default::default()
4428                },
4429            )
4430            .await
4431            .expect("new_volume failed");
4432
4433        assert!(matches!(&*store.last_object_id.lock(), LastObjectId::Encrypted { .. }));
4434
4435        // Create some files so that our in-memory copy of StoreInfo has changes (the object
4436        // count) pending a flush.
4437        let root_dir_id = store.root_directory_object_id();
4438        let root_dir = Arc::new(Directory::open(&store, root_dir_id).await.expect("open failed"));
4439
4440        let _executor_tasks = testing::force_executor_threads_to_run(2).await;
4441
4442        for j in 0..100 {
4443            let mut transaction = store
4444                .new_transaction(
4445                    lock_keys![LockKey::object(store.store_object_id(), root_dir_id)],
4446                    Options::default(),
4447                )
4448                .await
4449                .expect("new_transaction failed");
4450            root_dir.create_child_file(&mut transaction, &format!("file {j}")).await.unwrap();
4451            transaction.commit().await.expect("commit failed");
4452
4453            let task = {
4454                let fs = fs.clone();
4455                fasync::Task::spawn(async move {
4456                    fs.journal().force_compact().await.unwrap();
4457                })
4458            };
4459
4460            // Hack the last object ID to force a roll of the object ID cipher.
4461            {
4462                let mut last_object_id = store.last_object_id.lock();
4463                let LastObjectId::Encrypted { id, .. } = &mut *last_object_id else {
4464                    unreachable!()
4465                };
4466                assert_eq!(*id >> 32, j);
4467                *id |= 0xffffffff;
4468            }
4469
4470            let mut transaction = store
4471                .new_transaction(
4472                    lock_keys![LockKey::object(
4473                        store.store_object_id(),
4474                        store.root_directory_object_id()
4475                    )],
4476                    Options::default(),
4477                )
4478                .await
4479                .expect("new_transaction failed");
4480            let root_directory = Directory::open(&store, store.root_directory_object_id())
4481                .await
4482                .expect("open failed");
4483            root_directory
4484                .create_child_file(&mut transaction, "test {j}")
4485                .await
4486                .expect("create_child_file failed");
4487            transaction.commit().await.expect("commit failed");
4488
4489            task.await;
4490
4491            // Check that the key has been changed.
4492            let new_store_info = store.load_store_info().await.unwrap();
4493
4494            let LastObjectIdInfo::Encrypted { id, key } = new_store_info.last_object_id else {
4495                unreachable!()
4496            };
4497            assert_eq!(id >> 32, j + 1);
4498            let LastObjectIdInfo::Encrypted { key: in_memory_key, .. } =
4499                store.store_info().unwrap().last_object_id
4500            else {
4501                unreachable!()
4502            };
4503            assert_eq!(key, in_memory_key);
4504        }
4505
4506        fs.close().await.expect("Close failed");
4507    }
4508
4509    #[fuchsia::test]
4510    async fn test_object_id_no_roll_for_unencrypted_store() {
4511        let fs = test_filesystem().await;
4512
4513        {
4514            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4515            let store = root_volume
4516                .new_volume("test", NewChildStoreOptions::default())
4517                .await
4518                .expect("new_volume failed");
4519
4520            // Hack the last object ID.
4521            {
4522                let mut last_object_id = store.last_object_id.lock();
4523                match &mut *last_object_id {
4524                    LastObjectId::Unencrypted { id } => {
4525                        assert_eq!(*id & OBJECT_ID_HI_MASK, 0);
4526                        *id |= 0xffffffff;
4527                    }
4528                    _ => unreachable!(),
4529                }
4530            }
4531
4532            let mut transaction = store
4533                .new_transaction(
4534                    lock_keys![LockKey::object(
4535                        store.store_object_id(),
4536                        store.root_directory_object_id()
4537                    )],
4538                    Options::default(),
4539                )
4540                .await
4541                .expect("new_transaction failed");
4542            let root_directory = Directory::open(&store, store.root_directory_object_id())
4543                .await
4544                .expect("open failed");
4545            let object = root_directory
4546                .create_child_file(&mut transaction, "test")
4547                .await
4548                .expect("create_child_file failed");
4549            transaction.commit().await.expect("commit failed");
4550
4551            assert_eq!(object.object_id(), 0x1_0000_0000);
4552
4553            // Check that there is still no key.
4554            assert_matches!(
4555                store.store_info().unwrap().last_object_id,
4556                LastObjectIdInfo::Unencrypted { .. }
4557            );
4558
4559            assert_eq!(store.last_object_id.lock().id(), 0x1_0000_0000);
4560        };
4561
4562        fs.close().await.expect("Close failed");
4563        let device = fs.take_device().await;
4564        device.reopen(false);
4565        let fs = FxFilesystem::open(device).await.expect("open failed");
4566        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4567        let store =
4568            root_volume.volume("test", StoreOptions::default()).await.expect("volume failed");
4569
4570        assert_eq!(store.last_object_id.lock().id(), 0x1_0000_0000);
4571    }
4572
4573    #[fuchsia::test]
4574    fn test_object_id_is_not_invalid_object_id() {
4575        let key = UnwrappedKey::new(vec![0; FXFS_KEY_SIZE]);
4576        // 1106634048 results in INVALID_OBJECT_ID with this key.
4577        let mut last_object_id =
4578            LastObjectId::Encrypted { id: 1106634047, cipher: Box::new(Ff1::new(&key)) };
4579        assert!(last_object_id.try_get_next().is_some());
4580        assert!(last_object_id.try_get_next().is_some());
4581    }
4582
4583    #[fuchsia::test]
4584    async fn test_last_object_id_is_correct_after_unlock() {
4585        let fs = test_filesystem().await;
4586        let crypt = Arc::new(new_insecure_crypt());
4587
4588        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4589        let store = root_volume
4590            .new_volume(
4591                "test",
4592                NewChildStoreOptions {
4593                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4594                    ..Default::default()
4595                },
4596            )
4597            .await
4598            .expect("new_volume failed");
4599
4600        let mut transaction = store
4601            .new_transaction(
4602                lock_keys![LockKey::object(
4603                    store.store_object_id(),
4604                    store.root_directory_object_id()
4605                )],
4606                Options::default(),
4607            )
4608            .await
4609            .expect("new_transaction failed");
4610        let root_directory =
4611            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
4612        root_directory
4613            .create_child_file(&mut transaction, "test")
4614            .await
4615            .expect("create_child_file failed");
4616        transaction.commit().await.expect("commit failed");
4617
4618        // Compact so that StoreInfo is written.
4619        fs.journal().force_compact().await.unwrap();
4620
4621        let last_object_id = store.last_object_id.lock().id();
4622
4623        store.lock().await.unwrap();
4624        store.unlock(crypt.clone()).await.unwrap();
4625
4626        assert_eq!(store.last_object_id.lock().id(), last_object_id);
4627    }
4628
4629    #[fuchsia::test(threads = 20)]
4630    async fn test_race_when_rolling_last_object_id_cipher() {
4631        // NOTE: This test is trying to test a race, so if it fails, it might be flaky.
4632
4633        const NUM_THREADS: usize = 20;
4634
4635        let fs = test_filesystem().await;
4636        let crypt = Arc::new(new_insecure_crypt());
4637
4638        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4639        let store = root_volume
4640            .new_volume(
4641                "test",
4642                NewChildStoreOptions {
4643                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4644                    ..Default::default()
4645                },
4646            )
4647            .await
4648            .expect("new_volume failed");
4649
4650        let store_id = store.store_object_id();
4651        let root_dir_id = store.root_directory_object_id();
4652
4653        let root_directory =
4654            Arc::new(Directory::open(&store, root_dir_id).await.expect("open failed"));
4655
4656        // Create directories.
4657        let mut directories = Vec::new();
4658        for _ in 0..NUM_THREADS {
4659            let mut transaction = fs
4660                .root_store()
4661                .new_transaction(
4662                    lock_keys![LockKey::object(store_id, root_dir_id,)],
4663                    Options::default(),
4664                )
4665                .await
4666                .expect("new_transaction failed");
4667            directories.push(
4668                root_directory
4669                    .create_child_dir(&mut transaction, "test")
4670                    .await
4671                    .expect("create_child_file failed"),
4672            );
4673            transaction.commit().await.expect("commit failed");
4674        }
4675
4676        // Hack the last object ID so that the next ID will require a roll.
4677        match &mut *store.last_object_id.lock() {
4678            LastObjectId::Encrypted { id, .. } => *id |= 0xffff_ffff,
4679            _ => unreachable!(),
4680        }
4681
4682        let scope = fasync::Scope::new();
4683
4684        let _executor_tasks = testing::force_executor_threads_to_run(NUM_THREADS).await;
4685
4686        for dir in directories {
4687            let fs = fs.clone();
4688            scope.spawn(async move {
4689                let mut transaction = fs
4690                    .root_store()
4691                    .new_transaction(
4692                        lock_keys![LockKey::object(store_id, dir.object_id(),)],
4693                        Options::default(),
4694                    )
4695                    .await
4696                    .expect("new_transaction failed");
4697                dir.create_child_file(&mut transaction, "test")
4698                    .await
4699                    .expect("create_child_file failed");
4700                transaction.commit().await.expect("commit failed");
4701            });
4702        }
4703
4704        scope.on_no_tasks().await;
4705
4706        assert_eq!(store.last_object_id.lock().id(), 0x1_0000_0000 + NUM_THREADS as u64 - 1);
4707    }
4708
4709    #[fuchsia::test(threads = 10)]
4710    async fn test_lock_store() {
4711        let fs = test_filesystem().await;
4712        let crypt = Arc::new(new_insecure_crypt());
4713
4714        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4715        let store = root_volume
4716            .new_volume(
4717                "test",
4718                NewChildStoreOptions {
4719                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4720                    ..NewChildStoreOptions::default()
4721                },
4722            )
4723            .await
4724            .expect("new_volume failed");
4725        let mut transaction = store
4726            .new_transaction(
4727                lock_keys![LockKey::object(
4728                    store.store_object_id(),
4729                    store.root_directory_object_id()
4730                )],
4731                Options::default(),
4732            )
4733            .await
4734            .expect("new_transaction failed");
4735        let root_directory =
4736            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
4737        root_directory
4738            .create_child_file(&mut transaction, "test")
4739            .await
4740            .expect("create_child_file failed");
4741        transaction.commit().await.expect("commit failed");
4742        store.lock().await.expect("lock failed");
4743
4744        store.unlock(crypt).await.expect("unlock failed");
4745        root_directory.lookup("test").await.expect("lookup failed").expect("not found");
4746    }
4747
4748    #[fuchsia::test(threads = 10)]
4749    async fn test_unlock_read_only() {
4750        let fs = test_filesystem().await;
4751        let crypt = Arc::new(new_insecure_crypt());
4752
4753        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4754        let store = root_volume
4755            .new_volume(
4756                "test",
4757                NewChildStoreOptions {
4758                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4759                    ..NewChildStoreOptions::default()
4760                },
4761            )
4762            .await
4763            .expect("new_volume failed");
4764        let mut transaction = store
4765            .new_transaction(
4766                lock_keys![LockKey::object(
4767                    store.store_object_id(),
4768                    store.root_directory_object_id()
4769                )],
4770                Options::default(),
4771            )
4772            .await
4773            .expect("new_transaction failed");
4774        let root_directory =
4775            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
4776        root_directory
4777            .create_child_file(&mut transaction, "test")
4778            .await
4779            .expect("create_child_file failed");
4780        transaction.commit().await.expect("commit failed");
4781        store.lock().await.expect("lock failed");
4782
4783        store.unlock_read_only(crypt.clone()).await.expect("unlock failed");
4784        root_directory.lookup("test").await.expect("lookup failed").expect("not found");
4785        store.lock_read_only();
4786        store.unlock_read_only(crypt).await.expect("unlock failed");
4787        root_directory.lookup("test").await.expect("lookup failed").expect("not found");
4788    }
4789
4790    #[fuchsia::test]
4791    async fn test_mutations_cipher_dropped_on_lock() {
4792        let fs = test_filesystem().await;
4793        let crypt = Arc::new(new_insecure_crypt());
4794
4795        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4796        let store = root_volume
4797            .new_volume(
4798                "test",
4799                NewChildStoreOptions {
4800                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4801                    ..NewChildStoreOptions::default()
4802                },
4803            )
4804            .await
4805            .expect("new_volume failed");
4806
4807        // When created/unlocked, mutations_cipher is present in LockState::Unlocked
4808        // and last_object_id is Encrypted.
4809        assert_matches!(*store.lock_state.lock(), LockState::Unlocked { .. });
4810        assert!(matches!(&*store.last_object_id.lock(), LastObjectId::Encrypted { .. }));
4811
4812        // When locked, LockState transitions to Locked, dropping mutations_cipher,
4813        // and last_object_id is reset to Pending.
4814        store.lock().await.expect("lock failed");
4815        assert_matches!(*store.lock_state.lock(), LockState::Locked);
4816        assert!(matches!(&*store.last_object_id.lock(), LastObjectId::Pending));
4817
4818        // When unlocked again, mutations_cipher is re-created in LockState::Unlocked
4819        // and last_object_id is restored to Encrypted.
4820        store.unlock(crypt).await.expect("unlock failed");
4821        assert_matches!(*store.lock_state.lock(), LockState::Unlocked { .. });
4822        assert!(matches!(&*store.last_object_id.lock(), LastObjectId::Encrypted { .. }));
4823    }
4824
4825    #[fuchsia::test(threads = 10)]
4826    async fn test_key_rolled_when_unlocked() {
4827        let fs = test_filesystem().await;
4828        let crypt = Arc::new(new_insecure_crypt());
4829
4830        let object_id;
4831        {
4832            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4833            let store = root_volume
4834                .new_volume(
4835                    "test",
4836                    NewChildStoreOptions {
4837                        options: StoreOptions {
4838                            crypt: Some(crypt.clone()),
4839                            ..StoreOptions::default()
4840                        },
4841                        ..Default::default()
4842                    },
4843                )
4844                .await
4845                .expect("new_volume failed");
4846            let mut transaction = store
4847                .new_transaction(
4848                    lock_keys![LockKey::object(
4849                        store.store_object_id(),
4850                        store.root_directory_object_id()
4851                    )],
4852                    Options::default(),
4853                )
4854                .await
4855                .expect("new_transaction failed");
4856            let root_directory = Directory::open(&store, store.root_directory_object_id())
4857                .await
4858                .expect("open failed");
4859            object_id = root_directory
4860                .create_child_file(&mut transaction, "test")
4861                .await
4862                .expect("create_child_file failed")
4863                .object_id();
4864            transaction.commit().await.expect("commit failed");
4865        }
4866
4867        fs.close().await.expect("Close failed");
4868        let mut device = fs.take_device().await;
4869
4870        // Repeatedly remount so that we can be sure that we can remount when there are many
4871        // mutations keys.
4872        for _ in 0..100 {
4873            device.reopen(false);
4874            let fs = FxFilesystem::open(device).await.expect("open failed");
4875            {
4876                let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4877                let store = root_volume
4878                    .volume(
4879                        "test",
4880                        StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
4881                    )
4882                    .await
4883                    .expect("open_volume failed");
4884
4885                // The key should get rolled every time we unlock.
4886                {
4887                    let lock_state = store.lock_state.lock();
4888                    let LockState::Unlocked { mutations_cipher, .. } = &*lock_state else {
4889                        panic!("Unexpected lock state: {lock_state:?}");
4890                    };
4891                    assert!(mutations_cipher.key_is_new());
4892                }
4893
4894                // Make sure there's an encrypted mutation.
4895                let handle =
4896                    ObjectStore::open_object(&store, object_id, HandleOptions::default(), None)
4897                        .await
4898                        .expect("open_object failed");
4899                let buffer = handle.allocate_buffer(100).await;
4900                handle
4901                    .write_or_append(Some(0), buffer.as_ref())
4902                    .await
4903                    .expect("write_or_append failed");
4904            }
4905            fs.close().await.expect("Close failed");
4906            device = fs.take_device().await;
4907        }
4908    }
4909
4910    #[test]
4911    fn test_store_info_max_serialized_size() {
4912        let info = StoreInfo {
4913            guid: [0xff; 16],
4914            last_object_id: LastObjectIdInfo::Encrypted {
4915                id: 0x1234567812345678,
4916                key: FxfsKey {
4917                    wrapping_key_id: 0x1234567812345678u128.to_le_bytes(),
4918                    key: WrappedKeyBytes::from([0xff; FXFS_WRAPPED_KEY_SIZE]),
4919                },
4920            },
4921            // Worst case, each layer should be 3/4 the size of the layer below it (because of the
4922            // compaction policy we're using).  If the smallest layer is 8,192 bytes, then 120
4923            // layers would take up a size that exceeds a 64 bit unsigned integer, so if this fits,
4924            // any size should fit.
4925            layers: vec![0x1234567812345678; 120],
4926            root_directory_object_id: 0x1234567812345678,
4927            graveyard_directory_object_id: 0x1234567812345678,
4928            object_count: 0x1234567812345678,
4929            mutations_key: Some(FxfsKey {
4930                wrapping_key_id: 0x1234567812345678u128.to_le_bytes(),
4931                key: WrappedKeyBytes::from([0xff; FXFS_WRAPPED_KEY_SIZE]),
4932            }),
4933            mutations_cipher_offset: 0x1234567812345678,
4934            encrypted_mutations_object_id: 0x1234567812345678,
4935            internal_directory_object_id: INVALID_OBJECT_ID,
4936        };
4937        let mut serialized_info = Vec::new();
4938        info.serialize_with_version(&mut serialized_info).unwrap();
4939        assert!(
4940            serialized_info.len() <= MAX_STORE_INFO_SERIALIZED_SIZE,
4941            "{}",
4942            serialized_info.len()
4943        );
4944    }
4945
4946    async fn reopen_after_crypt_failure_inner(read_only: bool) {
4947        let fs = test_filesystem().await;
4948        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
4949
4950        let store = {
4951            let crypt = Arc::new(new_insecure_crypt());
4952            let store = root_volume
4953                .new_volume(
4954                    "vol",
4955                    NewChildStoreOptions {
4956                        options: StoreOptions {
4957                            crypt: Some(crypt.clone()),
4958                            ..StoreOptions::default()
4959                        },
4960                        ..Default::default()
4961                    },
4962                )
4963                .await
4964                .expect("new_volume failed");
4965            let root_directory = Directory::open(&store, store.root_directory_object_id())
4966                .await
4967                .expect("open failed");
4968            let mut transaction = fs
4969                .root_store()
4970                .new_transaction(
4971                    lock_keys![LockKey::object(
4972                        store.store_object_id(),
4973                        root_directory.object_id()
4974                    )],
4975                    Options::default(),
4976                )
4977                .await
4978                .expect("new_transaction failed");
4979            root_directory
4980                .create_child_file(&mut transaction, "test")
4981                .await
4982                .expect("create_child_file failed");
4983            transaction.commit().await.expect("commit failed");
4984
4985            crypt.shutdown();
4986            let mut transaction = fs
4987                .root_store()
4988                .new_transaction(
4989                    lock_keys![LockKey::object(
4990                        store.store_object_id(),
4991                        root_directory.object_id()
4992                    )],
4993                    Options::default(),
4994                )
4995                .await
4996                .expect("new_transaction failed");
4997            root_directory
4998                .create_child_file(&mut transaction, "test2")
4999                .await
5000                .map(|_| ())
5001                .expect_err("create_child_file should fail");
5002            store.lock().await.expect("lock failed");
5003            store
5004        };
5005
5006        let crypt = Arc::new(new_insecure_crypt());
5007        if read_only {
5008            store.unlock_read_only(crypt).await.expect("unlock failed");
5009        } else {
5010            store.unlock(crypt).await.expect("unlock failed");
5011        }
5012        let root_directory =
5013            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
5014        root_directory.lookup("test").await.expect("lookup failed").expect("not found");
5015    }
5016
5017    #[fuchsia::test(threads = 10)]
5018    async fn test_reopen_after_crypt_failure() {
5019        reopen_after_crypt_failure_inner(false).await;
5020    }
5021
5022    #[fuchsia::test(threads = 10)]
5023    async fn test_reopen_read_only_after_crypt_failure() {
5024        reopen_after_crypt_failure_inner(true).await;
5025    }
5026
5027    #[fuchsia::test(threads = 10)]
5028    #[should_panic(expected = "Insufficient reservation space")]
5029    #[cfg(debug_assertions)]
5030    async fn large_transaction_causes_panic_in_debug_builds() {
5031        let fs = test_filesystem().await;
5032        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5033        let store = root_volume
5034            .new_volume("vol", NewChildStoreOptions::default())
5035            .await
5036            .expect("new_volume failed");
5037        let root_directory =
5038            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
5039        let mut transaction = fs
5040            .root_store()
5041            .new_transaction(
5042                lock_keys![LockKey::object(store.store_object_id(), root_directory.object_id())],
5043                Options::default(),
5044            )
5045            .await
5046            .expect("transaction");
5047        for i in 0..500 {
5048            root_directory
5049                .create_symlink(&mut transaction, b"link", &format!("{}", i))
5050                .await
5051                .expect("symlink");
5052        }
5053        assert_eq!(transaction.commit().await.expect("commit"), 0);
5054    }
5055
5056    #[fuchsia::test]
5057    async fn test_crypt_failure_does_not_fuse_journal() {
5058        let fs = test_filesystem().await;
5059
5060        {
5061            // Create two stores and a record for each store, so the journal will need to flush them
5062            // both later.
5063            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5064            let store1 = root_volume
5065                .new_volume(
5066                    "vol1",
5067                    NewChildStoreOptions {
5068                        options: StoreOptions {
5069                            crypt: Some(Arc::new(new_insecure_crypt())),
5070                            ..StoreOptions::default()
5071                        },
5072                        ..Default::default()
5073                    },
5074                )
5075                .await
5076                .expect("new_volume failed");
5077            let crypt = Arc::new(new_insecure_crypt());
5078            let store2 = root_volume
5079                .new_volume(
5080                    "vol2",
5081                    NewChildStoreOptions {
5082                        options: StoreOptions { crypt: Some(crypt.clone()) },
5083                        ..Default::default()
5084                    },
5085                )
5086                .await
5087                .expect("new_volume failed");
5088            for store in [&store1, &store2] {
5089                let root_directory = Directory::open(store, store.root_directory_object_id())
5090                    .await
5091                    .expect("open failed");
5092                let mut transaction = store
5093                    .new_transaction(
5094                        lock_keys![LockKey::object(
5095                            store.store_object_id(),
5096                            root_directory.object_id()
5097                        )],
5098                        Options::default(),
5099                    )
5100                    .await
5101                    .expect("new_transaction failed");
5102                root_directory
5103                    .create_child_file(&mut transaction, "test")
5104                    .await
5105                    .expect("create_child_file failed");
5106                transaction.commit().await.expect("commit failed");
5107            }
5108
5109            // Shut down the crypt instance for store2.
5110            crypt.shutdown();
5111
5112            // Compact. This flushes store2 (using cached key) and store1.  This consumes 1 cached
5113            // key from store2.
5114            fs.journal().force_compact().await.expect("compact failed");
5115
5116            // Write to store2 again. This should fail because store2 needs to top up its cached
5117            // keys, which requires calling the (dead) crypt service.
5118            let root_directory2 = Directory::open(&store2, store2.root_directory_object_id())
5119                .await
5120                .expect("open failed");
5121            let (child_id, _, _) = root_directory2
5122                .lookup("test")
5123                .await
5124                .expect("lookup failed")
5125                .expect("test file not found");
5126            let mut transaction2 = store2
5127                .new_transaction(
5128                    lock_keys![
5129                        LockKey::object(store2.store_object_id(), root_directory2.object_id()),
5130                        LockKey::object(store2.store_object_id(), child_id),
5131                    ],
5132                    Options::default(),
5133                )
5134                .await
5135                .expect("new_transaction failed");
5136            replace_child(&mut transaction2, None, (&root_directory2, "test"))
5137                .await
5138                .expect("replace_child failed");
5139            assert!(transaction2.commit().await.is_err());
5140
5141            // Write to store1 should still succeed (its crypt is not dead).
5142            let root_directory1 = Directory::open(&store1, store1.root_directory_object_id())
5143                .await
5144                .expect("open failed");
5145            let mut transaction1 = store1
5146                .new_transaction(
5147                    lock_keys![LockKey::object(
5148                        store1.store_object_id(),
5149                        root_directory1.object_id()
5150                    )],
5151                    Options::default(),
5152                )
5153                .await
5154                .expect("new_transaction failed");
5155            root_directory1
5156                .create_child_file(&mut transaction1, "test2")
5157                .await
5158                .expect("create_child_file failed");
5159            transaction1.commit().await.expect("commit failed");
5160
5161            // Compact again. Should succeed.
5162            fs.journal().force_compact().await.expect("compact failed");
5163        }
5164
5165        // Close and reopen to verify.
5166        fs.close().await.expect("close failed");
5167        let device = fs.take_device().await;
5168        device.reopen(false);
5169        let fs = FxFilesystem::open(device).await.expect("open failed");
5170        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5171
5172        // vol1 should have "test" and "test2".
5173        let store1 = root_volume
5174            .volume(
5175                "vol1",
5176                StoreOptions {
5177                    crypt: Some(Arc::new(new_insecure_crypt())),
5178                    ..StoreOptions::default()
5179                },
5180            )
5181            .await
5182            .expect("volume failed");
5183        let root_directory1 =
5184            Directory::open(&store1, store1.root_directory_object_id()).await.expect("open failed");
5185        assert!(root_directory1.lookup("test").await.expect("lookup failed").is_some());
5186        assert!(root_directory1.lookup("test2").await.expect("lookup failed").is_some());
5187
5188        // vol2 should only have "test".
5189        let store2 = root_volume
5190            .volume(
5191                "vol2",
5192                StoreOptions {
5193                    crypt: Some(Arc::new(new_insecure_crypt())),
5194                    ..StoreOptions::default()
5195                },
5196            )
5197            .await
5198            .expect("volume failed");
5199        let root_directory2 =
5200            Directory::open(&store2, store2.root_directory_object_id()).await.expect("open failed");
5201        assert!(root_directory2.lookup("test").await.expect("lookup failed").is_some());
5202        assert!(root_directory2.lookup("test2").await.expect("lookup failed").is_none());
5203
5204        fs.close().await.expect("close failed");
5205    }
5206
5207    #[fuchsia::test]
5208    async fn test_crypt_failure_during_unlock_race() {
5209        let fs = test_filesystem().await;
5210
5211        let store_object_id = {
5212            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5213            let store = root_volume
5214                .new_volume(
5215                    "vol",
5216                    NewChildStoreOptions {
5217                        options: StoreOptions { crypt: Some(Arc::new(new_insecure_crypt())) },
5218                        ..Default::default()
5219                    },
5220                )
5221                .await
5222                .expect("new_volume failed");
5223            let root_directory = Directory::open(&store, store.root_directory_object_id())
5224                .await
5225                .expect("open failed");
5226            let mut transaction = fs
5227                .root_store()
5228                .new_transaction(
5229                    lock_keys![LockKey::object(
5230                        store.store_object_id(),
5231                        root_directory.object_id()
5232                    )],
5233                    Options::default(),
5234                )
5235                .await
5236                .expect("new_transaction failed");
5237            root_directory
5238                .create_child_file(&mut transaction, "test")
5239                .await
5240                .expect("create_child_file failed");
5241            transaction.commit().await.expect("commit failed");
5242            store.store_object_id()
5243        };
5244
5245        fs.close().await.expect("close failed");
5246        let device = fs.take_device().await;
5247        device.reopen(false);
5248
5249        let fs = FxFilesystem::open(device).await.expect("open failed");
5250        {
5251            let fs_clone = fs.clone();
5252            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5253
5254            let crypt = Arc::new(new_insecure_crypt());
5255            let crypt_clone = crypt.clone();
5256            join!(
5257                async move {
5258                    // Unlock might fail, so ignore errors.
5259                    let _ =
5260                        root_volume.volume("vol", StoreOptions { crypt: Some(crypt_clone) }).await;
5261                },
5262                async move {
5263                    // Block until unlock is finished but before flushing due to unlock is finished, to
5264                    // maximize the chances of weirdness.
5265                    let keys = lock_keys![LockKey::flush(store_object_id)];
5266                    let _ = fs_clone.lock_manager().write_lock(keys).await;
5267                    crypt.shutdown();
5268                }
5269            );
5270        }
5271
5272        fs.close().await.expect("close failed");
5273        let device = fs.take_device().await;
5274        device.reopen(false);
5275
5276        let fs = FxFilesystem::open(device).await.expect("open failed");
5277        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5278        let store = root_volume
5279            .volume(
5280                "vol",
5281                StoreOptions {
5282                    crypt: Some(Arc::new(new_insecure_crypt())),
5283                    ..StoreOptions::default()
5284                },
5285            )
5286            .await
5287            .expect("open volume failed");
5288        let root_directory =
5289            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
5290        assert!(root_directory.lookup("test").await.expect("lookup failed").is_some());
5291
5292        fs.close().await.expect("close failed");
5293    }
5294
5295    #[fuchsia::test]
5296    async fn test_low_32_bit_object_ids() {
5297        let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
5298        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
5299
5300        {
5301            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
5302
5303            let store = root_vol
5304                .new_volume(
5305                    "test",
5306                    NewChildStoreOptions { low_32_bit_object_ids: true, ..Default::default() },
5307                )
5308                .await
5309                .expect("new_volume failed");
5310
5311            let root_dir = Directory::open(&store, store.root_directory_object_id())
5312                .await
5313                .expect("open failed");
5314
5315            let mut ids = std::collections::HashSet::new();
5316
5317            for i in 0..100 {
5318                let mut transaction = fs
5319                    .root_store()
5320                    .new_transaction(
5321                        lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
5322                        Options::default(),
5323                    )
5324                    .await
5325                    .expect("new_transaction failed");
5326
5327                for j in 0..100 {
5328                    let object = root_dir
5329                        .create_child_dir(&mut transaction, &format!("{i}.{j}"))
5330                        .await
5331                        .expect("create_child_file failed");
5332
5333                    assert!(object.object_id() < 1 << 32);
5334                    assert_ne!(object.object_id(), INVALID_OBJECT_ID);
5335                    assert!(ids.insert(object.object_id()));
5336                }
5337
5338                transaction.commit().await.expect("commit failed");
5339            }
5340
5341            assert_matches!(store.store_info().unwrap().last_object_id, LastObjectIdInfo::Low32Bit);
5342
5343            fsck_volume(&fs, store.store_object_id(), None).await.expect("fsck_volume failed");
5344        }
5345
5346        // Verify persistence
5347        fs.close().await.expect("Close failed");
5348        let device = fs.take_device().await;
5349        device.reopen(false);
5350        let fs = FxFilesystem::open(device).await.expect("open failed");
5351        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
5352        let store = root_vol.volume("test", StoreOptions::default()).await.expect("volume failed");
5353
5354        // Check that we can still create files and they have low 32-bit IDs.
5355        let root_dir =
5356            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
5357        let mut transaction = fs
5358            .root_store()
5359            .new_transaction(
5360                lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
5361                Options::default(),
5362            )
5363            .await
5364            .expect("new_transaction failed");
5365
5366        let object = root_dir
5367            .create_child_file(&mut transaction, "persistence_check")
5368            .await
5369            .expect("create_child_file failed");
5370        assert!(object.object_id() < 1 << 32);
5371
5372        transaction.commit().await.expect("commit failed");
5373
5374        assert_matches!(store.store_info().unwrap().last_object_id, LastObjectIdInfo::Low32Bit);
5375    }
5376
5377    struct StallingCrypt {
5378        delegate: Arc<dyn Crypt>,
5379        unwrap_stall_counter: AtomicIsize,
5380        unwrap_stalled_tx: Mutex<Option<oneshot::Sender<oneshot::Sender<()>>>>,
5381        create_key_stalled_tx: Mutex<Option<oneshot::Sender<oneshot::Sender<()>>>>,
5382    }
5383
5384    impl StallingCrypt {
5385        fn new(delegate: Arc<dyn Crypt>) -> Self {
5386            Self {
5387                delegate,
5388                unwrap_stall_counter: AtomicIsize::new(-1),
5389                unwrap_stalled_tx: Mutex::new(None),
5390                create_key_stalled_tx: Mutex::new(None),
5391            }
5392        }
5393
5394        fn stall_on_unwrap(
5395            self: &Arc<Self>,
5396            index: usize,
5397        ) -> oneshot::Receiver<oneshot::Sender<()>> {
5398            let (tx, rx) = oneshot::channel();
5399            self.unwrap_stall_counter.store(index as isize, Ordering::Relaxed);
5400            *self.unwrap_stalled_tx.lock() = Some(tx);
5401            rx
5402        }
5403
5404        fn stall_on_create_key(self: &Arc<Self>) -> oneshot::Receiver<oneshot::Sender<()>> {
5405            let (tx, rx) = oneshot::channel();
5406            *self.create_key_stalled_tx.lock() = Some(tx);
5407            rx
5408        }
5409    }
5410
5411    #[async_trait]
5412    impl Crypt for StallingCrypt {
5413        async fn create_key(
5414            &self,
5415            owner: u64,
5416            purpose: KeyPurpose,
5417        ) -> Result<(FxfsKey, UnwrappedKey), zx::Status> {
5418            if matches!(purpose, KeyPurpose::Data) {
5419                let stalled_tx = self.create_key_stalled_tx.lock().take();
5420                if let Some(tx) = stalled_tx {
5421                    let (continue_tx, continue_rx) = oneshot::channel();
5422                    let _ = tx.send(continue_tx);
5423                    let _ = continue_rx.await;
5424                }
5425            }
5426            self.delegate.create_key(owner, purpose).await
5427        }
5428
5429        async fn unwrap_key(
5430            &self,
5431            wrapped_key: &WrappedKey,
5432            owner: u64,
5433        ) -> Result<UnwrappedKey, zx::Status> {
5434            let count = self.unwrap_stall_counter.fetch_sub(1, Ordering::Relaxed);
5435            log::info!("unwrap_key called, count was {}, owner {}", count, owner);
5436            if count == 0 {
5437                log::info!("unwrap_key stalling");
5438                let stalled_tx = self.unwrap_stalled_tx.lock().take();
5439                if let Some(tx) = stalled_tx {
5440                    let (continue_tx, continue_rx) = oneshot::channel();
5441                    let _ = tx.send(continue_tx);
5442                    let _ = continue_rx.await;
5443                }
5444                log::info!("unwrap_key resumed");
5445            }
5446            self.delegate.unwrap_key(wrapped_key, owner).await
5447        }
5448
5449        async fn create_key_with_id(
5450            &self,
5451            owner: u64,
5452            wrapping_key_id: WrappingKeyId,
5453            object_type: ObjectType,
5454            flags: fidl_fuchsia_io::FscryptPolicyFlags,
5455        ) -> Result<(EncryptionKey, UnwrappedKey), zx::Status> {
5456            self.delegate.create_key_with_id(owner, wrapping_key_id, object_type, flags).await
5457        }
5458    }
5459
5460    #[fuchsia::test]
5461    async fn test_fsck_during_key_pre_cache_stall() {
5462        let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
5463        let fs = FxFilesystemBuilder::new().format(true).open(device).await.expect("open failed");
5464
5465        // Initialize crypt without blocker so new_volume doesn't stall.
5466        let crypt = Arc::new(StallingCrypt::new(Arc::new(new_insecure_crypt())));
5467
5468        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
5469        let store = root_vol
5470            .new_volume(
5471                "test",
5472                NewChildStoreOptions {
5473                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
5474                    ..Default::default()
5475                },
5476            )
5477            .await
5478            .expect("new_volume failed");
5479
5480        // Now install the blocker.
5481        let stalled_rx = crypt.stall_on_create_key();
5482
5483        // The volume was created, so it has 2 cached keys.
5484        // We want to exhaust them so that the next transaction tries to pre-cache more.
5485        {
5486            let mut lock_state = store.lock_state.lock();
5487            if let super::LockState::Unlocked { cached_keys, .. } = &mut *lock_state {
5488                cached_keys.clear();
5489            }
5490        }
5491
5492        // Start a transaction in the background. It will try to pre-cache keys,
5493        // and should stall on the crypt service.
5494        let store_clone = store.clone();
5495        let root_dir_id = store.root_directory_object_id();
5496
5497        let tx_join_handle = fasync::Task::spawn(async move {
5498            let root_dir = Directory::open(&store_clone, root_dir_id).await.expect("open failed");
5499            let mut transaction = store_clone
5500                .new_transaction(
5501                    lock_keys![LockKey::object(
5502                        store_clone.store_object_id(),
5503                        root_dir.object_id()
5504                    )],
5505                    Options::default(),
5506                )
5507                .await
5508                .expect("new_transaction failed");
5509            let name = "foo";
5510            root_dir
5511                .create_child_file(&mut transaction, name)
5512                .await
5513                .expect("create_child_file failed");
5514            transaction.commit().await.expect("commit failed");
5515        });
5516
5517        // Wait until the transaction has actually stalled on crypt.
5518        let continue_tx = stalled_rx.await.expect("stalled_rx failed");
5519
5520        // While it is blocked, we should still be able to run fsck.
5521        fsck(fs.clone()).await.expect("fsck failed");
5522
5523        // Unblock the crypt service so the transaction can complete.
5524        let _ = continue_tx.send(());
5525        tx_join_handle.await;
5526
5527        fs.close().await.expect("Close failed");
5528    }
5529
5530    #[fuchsia::test]
5531    async fn test_writes_to_other_store_not_blocked_by_stall() {
5532        let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
5533        let fs = FxFilesystemBuilder::new()
5534            .format(true)
5535            .journal_options(JournalOptions { reclaim_size: 32_768, ..Default::default() })
5536            .open(device)
5537            .await
5538            .expect("open failed");
5539
5540        // Initialize crypt for store1 with a blocker.
5541        // We will enable it after volume creation.
5542        let crypt1 = Arc::new(StallingCrypt::new(Arc::new(new_insecure_crypt())));
5543
5544        // Initialize normal crypt for store2.
5545        let crypt2 = Arc::new(new_insecure_crypt());
5546
5547        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
5548
5549        let store1 = root_vol
5550            .new_volume(
5551                "test1",
5552                NewChildStoreOptions {
5553                    options: StoreOptions {
5554                        crypt: Some(crypt1.clone()),
5555                        ..StoreOptions::default()
5556                    },
5557                    ..Default::default()
5558                },
5559            )
5560            .await
5561            .expect("new_volume failed");
5562
5563        store1.flush().await.expect("flush failed");
5564
5565        let store2 = root_vol
5566            .new_volume(
5567                "test2",
5568                NewChildStoreOptions {
5569                    options: StoreOptions {
5570                        crypt: Some(crypt2.clone()),
5571                        ..StoreOptions::default()
5572                    },
5573                    ..Default::default()
5574                },
5575            )
5576            .await
5577            .expect("new_volume failed");
5578
5579        // Now install the blocker on crypt1.
5580        let stalled_rx = crypt1.stall_on_create_key();
5581
5582        // Exhaust cached keys for store1 so next txn tries to pre-cache.
5583        {
5584            let mut lock_state = store1.lock_state.lock();
5585            if let super::LockState::Unlocked { cached_keys, .. } = &mut *lock_state {
5586                cached_keys.clear();
5587            }
5588        }
5589
5590        // Start a transaction on store1 in the background.
5591        // It should stall on crypt1.
5592        let store1_clone = store1.clone();
5593        let root_dir1_id = store1.root_directory_object_id();
5594
5595        let tx_join_handle = fasync::Task::spawn(async move {
5596            let root_dir1 =
5597                Directory::open(&store1_clone, root_dir1_id).await.expect("open failed");
5598            let mut transaction = store1_clone
5599                .new_transaction(
5600                    lock_keys![LockKey::object(
5601                        store1_clone.store_object_id(),
5602                        root_dir1.object_id()
5603                    )],
5604                    Options::default(),
5605                )
5606                .await
5607                .expect("new_transaction failed");
5608            root_dir1
5609                .create_child_file(&mut transaction, "foo")
5610                .await
5611                .expect("create_child_file failed");
5612            transaction.commit().await.expect("commit failed");
5613        });
5614
5615        // Wait until store1 transaction has actually stalled.
5616        let continue_tx = stalled_rx.await.expect("stalled_rx failed");
5617
5618        // While store1 is blocked, we should still be able to write to store2.
5619        let root_dir2 =
5620            Directory::open(&store2, store2.root_directory_object_id()).await.expect("open failed");
5621        let long_name = "a".repeat(255);
5622        let mut i = 0;
5623        while store2.counters.lock().num_flushes < 3 {
5624            if i > 200 {
5625                panic!(
5626                    "Failed to trigger 3 compactions after 200 transactions. Flushes: {}",
5627                    store2.counters.lock().num_flushes
5628                );
5629            }
5630            let mut transaction = store2
5631                .new_transaction(
5632                    lock_keys![LockKey::object(store2.store_object_id(), root_dir2.object_id())],
5633                    Options::default(),
5634                )
5635                .await
5636                .expect("new_transaction failed");
5637            root_dir2
5638                .create_child_file(&mut transaction, &format!("{}-{:03}", long_name, i))
5639                .await
5640                .expect("create_child_file failed");
5641            transaction.commit().await.expect("commit failed");
5642            i += 1;
5643            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
5644        }
5645
5646        // Unblock crypt1.
5647        let _ = continue_tx.send(());
5648        tx_join_handle.await;
5649
5650        fs.close().await.expect("Close failed");
5651    }
5652
5653    #[fuchsia::test]
5654    async fn test_concurrent_transactions_exhaust_cached_keys() {
5655        let fs = test_filesystem().await;
5656
5657        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5658        let crypt = Arc::new(new_insecure_crypt());
5659        let store = root_volume
5660            .new_volume(
5661                "vol",
5662                NewChildStoreOptions {
5663                    options: StoreOptions { crypt: Some(crypt.clone()), ..Default::default() },
5664                    ..Default::default()
5665                },
5666            )
5667            .await
5668            .expect("new_volume failed");
5669
5670        let root_directory =
5671            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
5672
5673        // Create three directories first, so we can run transactions concurrently
5674        // without blocking on directory locks.
5675        let mut transaction = store
5676            .new_transaction(
5677                lock_keys![LockKey::object(store.store_object_id(), root_directory.object_id())],
5678                Options::default(),
5679            )
5680            .await
5681            .expect("new_transaction failed");
5682        let dir1 = root_directory
5683            .create_child_dir(&mut transaction, "dir1")
5684            .await
5685            .expect("create_child_dir failed");
5686        let dir2 = root_directory
5687            .create_child_dir(&mut transaction, "dir2")
5688            .await
5689            .expect("create_child_dir failed");
5690        let dir3 = root_directory
5691            .create_child_dir(&mut transaction, "dir3")
5692            .await
5693            .expect("create_child_dir failed");
5694        transaction.commit().await.expect("commit failed");
5695
5696        // Start three transactions. They will all see that the cached keys are full
5697        // (size 2) after the first one tops them up, so they will all succeed to start
5698        // without calling crypt again (except the first one which tops up).
5699
5700        // Transaction 1
5701        let mut transaction1 = store
5702            .new_transaction(
5703                lock_keys![LockKey::object(store.store_object_id(), dir1.object_id())],
5704                Options::default(),
5705            )
5706            .await
5707            .expect("new_transaction 1 failed");
5708        dir1.create_child_file(&mut transaction1, "file1")
5709            .await
5710            .expect("create_child_file 1 failed");
5711
5712        // Transaction 2
5713        let mut transaction2 = store
5714            .new_transaction(
5715                lock_keys![LockKey::object(store.store_object_id(), dir2.object_id())],
5716                Options::default(),
5717            )
5718            .await
5719            .expect("new_transaction 2 failed");
5720        dir2.create_child_file(&mut transaction2, "file2")
5721            .await
5722            .expect("create_child_file 2 failed");
5723
5724        // Transaction 3
5725        let mut transaction3 = store
5726            .new_transaction(
5727                lock_keys![LockKey::object(store.store_object_id(), dir3.object_id())],
5728                Options::default(),
5729            )
5730            .await
5731            .expect("new_transaction 3 failed");
5732        dir3.create_child_file(&mut transaction3, "file3")
5733            .await
5734            .expect("create_child_file 3 failed");
5735
5736        // Now shut down the crypt service.
5737        crypt.shutdown();
5738
5739        // Commit transaction 1 and compact. This should succeed and consume 1 cached key.
5740        transaction1.commit().await.expect("commit 1 failed");
5741        fs.journal().force_compact().await.expect("compact 1 failed");
5742
5743        // Commit transaction 2 should FAIL because it tries to top up (since cache size is 1 < 2)
5744        // and the crypt service is dead.
5745        assert!(transaction2.commit().await.is_err());
5746
5747        fs.close().await.expect("Close failed");
5748    }
5749
5750    #[fuchsia::test(threads = 10)]
5751    async fn test_key_exhaustion_race() {
5752        let fs;
5753        let in_hook = AtomicBool::new(false);
5754        let (mut hooks, fs_hooks) = Hooks::new();
5755        fs = test_filesystem_with_hooks(fs_hooks).await;
5756
5757        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
5758        let crypt = Arc::new(new_insecure_crypt());
5759        let store = root_volume
5760            .new_volume(
5761                "vol",
5762                NewChildStoreOptions {
5763                    options: StoreOptions { crypt: Some(crypt.clone()), ..Default::default() },
5764                    ..Default::default()
5765                },
5766            )
5767            .await
5768            .expect("new_volume failed");
5769
5770        let root_directory =
5771            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
5772
5773        // Commit a transaction to ensure the store has dirty mutations and needs a flush.
5774        let mut transaction = store
5775            .new_transaction(
5776                lock_keys![LockKey::object(store.store_object_id(), root_directory.object_id())],
5777                Options::default(),
5778            )
5779            .await
5780            .expect("new_transaction failed");
5781        root_directory
5782            .create_child_dir(&mut transaction, "dir")
5783            .await
5784            .expect("create_child_dir failed");
5785        transaction.commit().await.expect("commit failed");
5786
5787        hooks.set_before_commit(|| {
5788            if in_hook.swap(true, Ordering::Relaxed) {
5789                return;
5790            }
5791
5792            // Run compaction. Since this hook runs before the next transaction acquires the
5793            // commit lock, this compaction will flush the mutations committed above and consume
5794            // one cached key.
5795            futures::executor::block_on(fs.journal().force_compact()).expect("compact failed");
5796        });
5797
5798        // Start a second transaction. When we commit this transaction, the hook we set up above
5799        // will trigger.
5800        let mut transaction = store
5801            .new_transaction(
5802                lock_keys![LockKey::object(store.store_object_id(), root_directory.object_id())],
5803                Options::default(),
5804            )
5805            .await
5806            .expect("new_transaction failed");
5807        root_directory
5808            .create_child_file(&mut transaction, "file1")
5809            .await
5810            .expect("create_child_file failed");
5811
5812        // When we commit:
5813        // 1. `prepare_commit` runs. It checks the key cache. If the limit is 1, it sees 1 key
5814        //    (which is >= the limit), so it does not top up. If the limit is 2, it sees 2 keys
5815        //    (which is >= the limit), so it also does not top up.
5816        // 2. The hook runs compaction. Compaction flushes the first transaction's mutations,
5817        //    consuming one cached key.
5818        // 3. This transaction commits. The store is marked as needing a flush. If the limit was
5819        //    1, the cache is now empty. If the limit was 2, the cache has 1 key left.
5820        transaction.commit().await.expect("commit failed");
5821
5822        // Run compaction again. It will try to flush the second transaction's mutations.
5823        // If the limit is 1, this will fail because the cache is empty.
5824        // If the limit is 2, `prepare_commit` would have topped up the cache to 2 keys, so
5825        // compaction would have left 1 key, and this compaction will succeed.
5826        fs.journal().force_compact().await.expect("compaction failed");
5827
5828        fs.close().await.expect("Close failed");
5829    }
5830
5831    #[fuchsia::test(threads = 10)]
5832    async fn test_unlock_pre_cache_stall_does_not_block_other_stores() {
5833        let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
5834        let fs = FxFilesystemBuilder::new()
5835            .format(true)
5836            .journal_options(JournalOptions { reclaim_size: 32_768, ..Default::default() })
5837            .open(device)
5838            .await
5839            .expect("open failed");
5840
5841        let crypt1 = Arc::new(StallingCrypt::new(Arc::new(new_insecure_crypt())));
5842        let crypt2 = Arc::new(new_insecure_crypt());
5843
5844        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
5845
5846        let store1 = root_vol
5847            .new_volume(
5848                "test1",
5849                NewChildStoreOptions {
5850                    options: StoreOptions {
5851                        crypt: Some(crypt1.clone()),
5852                        ..StoreOptions::default()
5853                    },
5854                    ..Default::default()
5855                },
5856            )
5857            .await
5858            .expect("new_volume failed");
5859
5860        // Write some data to store1 and lock it.
5861        {
5862            let mut transaction = store1
5863                .new_transaction(
5864                    lock_keys![LockKey::object(
5865                        store1.store_object_id(),
5866                        store1.root_directory_object_id()
5867                    )],
5868                    Options::default(),
5869                )
5870                .await
5871                .expect("new_transaction failed");
5872            let root_dir1 = Directory::open(&store1, store1.root_directory_object_id())
5873                .await
5874                .expect("open failed");
5875            root_dir1
5876                .create_child_file(&mut transaction, "foo")
5877                .await
5878                .expect("create_child_file failed");
5879            transaction.commit().await.expect("commit failed");
5880        }
5881        store1.lock().await.expect("lock failed");
5882
5883        let store2 = root_vol
5884            .new_volume(
5885                "test2",
5886                NewChildStoreOptions {
5887                    options: StoreOptions {
5888                        crypt: Some(crypt2.clone()),
5889                        ..StoreOptions::default()
5890                    },
5891                    ..Default::default()
5892                },
5893            )
5894            .await
5895            .expect("new_volume failed");
5896
5897        // Now install the blocker on crypt1.
5898        let stalled_rx = crypt1.stall_on_create_key();
5899
5900        // Start unlocking store1 in the background. It should stall on crypt1 during
5901        // pre_cache_keys.
5902        let store1_clone = store1.clone();
5903        let crypt1_clone = crypt1.clone();
5904        let unlock_join_handle = fasync::Task::spawn(async move {
5905            store1_clone.unlock(crypt1_clone).await.expect("unlock failed");
5906        });
5907
5908        // Wait until store1 unlock has actually stalled.
5909        let continue_tx = stalled_rx.await.expect("stalled_rx failed");
5910
5911        // While store1's unlock is blocked, we should still be able to write to store2
5912        // and trigger compactions.
5913        let root_dir2 =
5914            Directory::open(&store2, store2.root_directory_object_id()).await.expect("open failed");
5915        let long_name = "a".repeat(255);
5916        let mut i = 0;
5917        while store2.counters.lock().num_flushes < 3 {
5918            if i > 200 {
5919                panic!(
5920                    "Failed to trigger 3 compactions after 200 transactions. Flushes: {}",
5921                    store2.counters.lock().num_flushes
5922                );
5923            }
5924            let mut transaction = store2
5925                .new_transaction(
5926                    lock_keys![LockKey::object(store2.store_object_id(), root_dir2.object_id())],
5927                    Options::default(),
5928                )
5929                .await
5930                .expect("new_transaction failed");
5931            root_dir2
5932                .create_child_file(&mut transaction, &format!("{}-{:03}", long_name, i))
5933                .await
5934                .expect("create_child_file failed");
5935            transaction.commit().await.expect("commit failed");
5936            i += 1;
5937            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
5938        }
5939
5940        // Unblock crypt1.
5941        let _ = continue_tx.send(());
5942        unlock_join_handle.await;
5943
5944        fs.close().await.expect("Close failed");
5945    }
5946
5947    async fn run_unlock_stall_test(unwrap_stall_index: usize) -> bool {
5948        log::info!("Starting run_unlock_stall_test for index {}", unwrap_stall_index);
5949        let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
5950        let crypt1 = Arc::new(new_insecure_crypt());
5951        let crypt2 = Arc::new(new_insecure_crypt());
5952        let (store1_id, device) = {
5953            let fs = FxFilesystemBuilder::new()
5954                .format(true)
5955                .journal_options(JournalOptions { reclaim_size: 32_768, ..Default::default() })
5956                .open(device)
5957                .await
5958                .expect("open failed");
5959
5960            let store1_id = {
5961                let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
5962
5963                let store1 = root_vol
5964                    .new_volume(
5965                        "test1",
5966                        NewChildStoreOptions {
5967                            options: StoreOptions {
5968                                crypt: Some(crypt1.clone()),
5969                                ..StoreOptions::default()
5970                            },
5971                            ..Default::default()
5972                        },
5973                    )
5974                    .await
5975                    .expect("new_volume failed");
5976
5977                root_vol
5978                    .new_volume(
5979                        "test2",
5980                        NewChildStoreOptions {
5981                            options: StoreOptions {
5982                                crypt: Some(crypt2.clone()),
5983                                ..StoreOptions::default()
5984                            },
5985                            ..Default::default()
5986                        },
5987                    )
5988                    .await
5989                    .expect("new_volume failed");
5990
5991                let root_dir1 = Directory::open(&store1, store1.root_directory_object_id())
5992                    .await
5993                    .expect("open failed");
5994
5995                for i in 0..8 {
5996                    let mut transaction = store1
5997                        .new_transaction(
5998                            lock_keys![LockKey::object(
5999                                store1.store_object_id(),
6000                                root_dir1.object_id()
6001                            )],
6002                            Options::default(),
6003                        )
6004                        .await
6005                        .expect("new_transaction failed");
6006                    let name = format!("file_{i}");
6007                    root_dir1
6008                        .create_child_file(&mut transaction, &name)
6009                        .await
6010                        .expect("create_child_file failed");
6011                    transaction.commit().await.expect("commit failed");
6012                }
6013
6014                // Flush store1 to roll the key.
6015                store1.flush().await.expect("flush failed");
6016
6017                // Now write something after the flush, so it is encrypted with the new key
6018                // and remains in the journal.
6019                {
6020                    let mut transaction = store1
6021                        .new_transaction(
6022                            lock_keys![LockKey::object(
6023                                store1.store_object_id(),
6024                                root_dir1.object_id()
6025                            )],
6026                            Options::default(),
6027                        )
6028                        .await
6029                        .expect("new_transaction failed");
6030                    root_dir1
6031                        .create_child_file(&mut transaction, "file_after_flush")
6032                        .await
6033                        .expect("create_child_file failed");
6034                    transaction.commit().await.expect("commit failed");
6035                }
6036                store1.store_object_id()
6037            };
6038
6039            fs.close().await.expect("Close failed");
6040            let device = fs.take_device().await;
6041            (store1_id, device)
6042        };
6043        device.reopen(false);
6044
6045        log::info!("Opening filesystem");
6046        let fs = FxFilesystemBuilder::new()
6047            .journal_options(JournalOptions { reclaim_size: 32_768, ..Default::default() })
6048            .open(device)
6049            .await
6050            .expect("open failed");
6051        log::info!("Opened filesystem, getting root volume");
6052        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
6053        log::info!("Got root volume");
6054
6055        log::info!("Re-opening store2");
6056        // Re-open store2.
6057        let store2 = root_vol
6058            .volume("test2", StoreOptions { crypt: Some(crypt2), ..StoreOptions::default() })
6059            .await
6060            .expect("volume failed");
6061        log::info!("Re-opened store2");
6062
6063        // Prepare stalling crypt for store1.
6064        let stalling_crypt1 = Arc::new(StallingCrypt::new(crypt1));
6065        let stalled_rx = stalling_crypt1.stall_on_unwrap(unwrap_stall_index);
6066
6067        // Get store1 handle from object manager without unlocking.
6068        let store1 = fs.object_manager().store(store1_id).unwrap();
6069
6070        log::info!("Spawning store1 unlock task");
6071        // Start unlocking store1 in the background.
6072        let store1_clone = store1.clone();
6073        let stalling_crypt1_clone = stalling_crypt1.clone();
6074        let unlock_join_handle = fasync::Task::spawn(async move {
6075            log::info!("store1 unlock task started");
6076            store1_clone.unlock(stalling_crypt1_clone).await.expect("unlock failed");
6077            log::info!("store1 unlock task finished");
6078        });
6079
6080        // Wait for either the stall to occur or unlock to complete.
6081        let continue_tx = futures::select! {
6082            res = stalled_rx.fuse() => {
6083                log::info!("Stalled at index {}", unwrap_stall_index);
6084                Some(res.expect("stalled_rx failed"))
6085            }
6086            _ = unlock_join_handle.fuse() => {
6087                log::info!("Unlock completed without stalling at index {}", unwrap_stall_index);
6088                None
6089            }
6090        };
6091
6092        let stalled = continue_tx.is_some();
6093
6094        if let Some(continue_tx) = continue_tx {
6095            log::info!("Writing to store2 while stalled at index {}", unwrap_stall_index);
6096            // While store1's unlock is blocked, try to write to store2 and trigger compactions.
6097            let root_dir2 = Directory::open(&store2, store2.root_directory_object_id())
6098                .await
6099                .expect("open failed");
6100            let long_name = "a".repeat(255);
6101            let mut i = 0;
6102            while store2.counters.lock().num_flushes < 3 {
6103                if i > 200 {
6104                    panic!(
6105                        "Failed to trigger 3 compactions after 200 transactions. Flushes: {}",
6106                        store2.counters.lock().num_flushes
6107                    );
6108                }
6109                let mut transaction = store2
6110                    .new_transaction(
6111                        lock_keys![LockKey::object(
6112                            store2.store_object_id(),
6113                            root_dir2.object_id()
6114                        )],
6115                        Options::default(),
6116                    )
6117                    .await
6118                    .expect("new_transaction failed");
6119                root_dir2
6120                    .create_child_file(&mut transaction, &format!("{}-{:03}", long_name, i))
6121                    .await
6122                    .expect("create_child_file failed");
6123                transaction.commit().await.expect("commit failed");
6124                i += 1;
6125                fasync::Timer::new(std::time::Duration::from_millis(5)).await;
6126            }
6127
6128            log::info!("Unblocking store1 at index {}", unwrap_stall_index);
6129            // Unblock store1 unlock.
6130            let _ = continue_tx.send(());
6131        }
6132
6133        fs.close().await.expect("Close failed");
6134        log::info!(
6135            "Finished run_unlock_stall_test for index {}, stalled: {}",
6136            unwrap_stall_index,
6137            stalled
6138        );
6139        stalled
6140    }
6141
6142    #[fuchsia::test(threads = 10)]
6143    async fn test_unlock_replay_stall_does_not_block_other_stores() {
6144        let mut unwrap_stall_index = 0;
6145        loop {
6146            let stalled = run_unlock_stall_test(unwrap_stall_index).await;
6147            if !stalled {
6148                break;
6149            }
6150            unwrap_stall_index += 1;
6151        }
6152    }
6153
6154    #[fuchsia::test(threads = 10)]
6155    async fn test_unlock_flush_race() {
6156        let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
6157        let crypt = Arc::new(new_insecure_crypt());
6158
6159        // Phase 1: Format and set up store1 with some mutations.
6160        let (store1_id, device) = {
6161            let (store1_id, fs) = {
6162                let fs = FxFilesystemBuilder::new()
6163                    .format(true)
6164                    .open(device)
6165                    .await
6166                    .expect("open failed");
6167                let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
6168                let store1 = root_vol
6169                    .new_volume(
6170                        "test1",
6171                        NewChildStoreOptions {
6172                            options: StoreOptions {
6173                                crypt: Some(crypt.clone()),
6174                                ..StoreOptions::default()
6175                            },
6176                            ..Default::default()
6177                        },
6178                    )
6179                    .await
6180                    .expect("new_volume failed");
6181
6182                let root_dir = Directory::open(&store1, store1.root_directory_object_id())
6183                    .await
6184                    .expect("open failed");
6185                let mut transaction = store1
6186                    .new_transaction(
6187                        lock_keys![LockKey::object(store1.store_object_id(), root_dir.object_id())],
6188                        Options::default(),
6189                    )
6190                    .await
6191                    .expect("new_transaction failed");
6192                root_dir
6193                    .create_child_file(&mut transaction, "file1")
6194                    .await
6195                    .expect("create_child_file failed");
6196                transaction.commit().await.expect("commit failed");
6197
6198                (store1.store_object_id(), fs)
6199            };
6200            fs.close().await.expect("Close failed");
6201            let device = fs.take_device().await;
6202            (store1_id, device)
6203        };
6204        device.reopen(false);
6205
6206        // Phase 2: Reopen and unlock with a race.
6207        let fs;
6208        let store1;
6209        let (mut hooks, fs_hooks) = Hooks::new();
6210        fs = FxFilesystemBuilder::new().hooks(fs_hooks).open(device).await.expect("open failed");
6211
6212        store1 = fs.object_manager().store(store1_id).unwrap();
6213
6214        // Set up the callback to trigger a flush of store1 during unlock (when flush
6215        // lock is dropped).
6216        hooks.set_unlock_resources_acquired(|| {
6217            futures::executor::block_on(store1.flush()).expect("flush failed");
6218        });
6219
6220        store1.unlock(crypt).await.expect("unlock failed");
6221
6222        fs.journal().force_compact().await.expect("compact failed");
6223
6224        fsck(fs.clone()).await.expect("fsck failed");
6225
6226        fs.close().await.expect("Close failed");
6227    }
6228
6229    #[fuchsia::test]
6230    async fn test_open_object_with_illegal_key_in_root_store() {
6231        let fs = test_filesystem().await;
6232        let crypt = Arc::new(new_insecure_crypt());
6233        let object_id = {
6234            let store = fs.root_store();
6235            let mut transaction = store
6236                .new_transaction(lock_keys![], Options::default())
6237                .await
6238                .expect("new_transaction failed");
6239
6240            let reserved_id =
6241                store.get_next_object_id(&transaction).await.expect("get_next_object_id failed");
6242            let object_id = reserved_id.get();
6243
6244            let (fxfs_key, unwrapped_key) =
6245                crypt.create_key(object_id, KeyPurpose::Data).await.expect("create_key failed");
6246
6247            let handle = ObjectStore::create_object_with_id(
6248                &store,
6249                &mut transaction,
6250                reserved_id,
6251                HandleOptions::default(),
6252                Some(ObjectEncryptionOptions {
6253                    permanent: false,
6254                    key_id: 0,
6255                    key: EncryptionKey::Fxfs(fxfs_key),
6256                    unwrapped_key,
6257                }),
6258            )
6259            .expect("create_object_with_id failed");
6260            transaction.commit().await.expect("commit failed");
6261
6262            let buf = handle.allocate_buffer(fs.block_size().get() as usize).await;
6263            handle.write_or_append(None, buf.as_ref()).await.expect("Write some data");
6264
6265            // Manually overwrite keys to be appended with an illegal key type. It won't even
6266            // attempt to use the key since nothing references it, but it will still cause the
6267            // failure.
6268            let old_keys = store.get_keys(object_id).await.expect("Old keys should work fine");
6269            let mut new_keys =
6270                vec![(1, EncryptionKey::FscryptInoLblk32File { key_identifier: [0; 16] })];
6271            for key in old_keys.iter() {
6272                new_keys.push(key.clone());
6273            }
6274            let mut transaction = store
6275                .new_transaction(
6276                    lock_keys![LockKey::Object {
6277                        store_object_id: store.store_object_id,
6278                        object_id
6279                    }],
6280                    Options::default(),
6281                )
6282                .await
6283                .expect("new_transaction failed");
6284            transaction.add(
6285                store.store_object_id(),
6286                Mutation::replace_or_insert_object(
6287                    ObjectKey::keys(object_id),
6288                    ObjectValue::keys(new_keys.into()),
6289                ),
6290            );
6291            transaction.commit().await.expect("commit failed");
6292
6293            object_id
6294        };
6295
6296        fs.close().await.expect("close failed");
6297        let device = fs.take_device().await;
6298        device.reopen(false);
6299        let fs = FxFilesystem::open(device).await.expect("open failed");
6300
6301        let store = fs.root_store();
6302        assert!(
6303            ObjectStore::open_object(&store, object_id, HandleOptions::default(), Some(crypt))
6304                .await
6305                .is_err()
6306        );
6307
6308        let res = store.get_keys(object_id).await;
6309        assert!(matches!(res, Err(e) if FxfsError::IntegrityError.matches(&e)));
6310    }
6311
6312    // Tests mixing the old encryption format with the new by writing an encrypted mutations object
6313    // in the old format and manually attaching it to the object store.
6314    #[fuchsia::test]
6315    async fn test_encrypted_mutations_mixed_versions() {
6316        let fs = test_filesystem().await;
6317        let crypt = Arc::new(new_insecure_crypt());
6318
6319        // Create a new encrypted child store "test" and flush so its root directory is in layers.
6320        let (store_object_id, root_dir_id, unwrapped_key) = {
6321            let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
6322            let store = root_volume
6323                .new_volume(
6324                    "test",
6325                    NewChildStoreOptions {
6326                        options: StoreOptions {
6327                            crypt: Some(crypt.clone()),
6328                            ..StoreOptions::default()
6329                        },
6330                        ..Default::default()
6331                    },
6332                )
6333                .await
6334                .expect("new_volume failed");
6335            let store_object_id = store.store_object_id();
6336            let root_dir_id = store.root_directory_object_id();
6337            let store_info = store.store_info().unwrap();
6338            let key = store_info.mutations_key.as_ref().unwrap().clone();
6339            let unwrapped_key = crypt
6340                .unwrap_key(&fxfs_crypto::WrappedKey::Fxfs(key.into()), store_object_id)
6341                .await
6342                .expect("unwrap_key failed");
6343
6344            // Flush the filesystem so that the store creation and root directory are saved to layers.
6345            fs.object_manager()
6346                .flush(FlushReason::Journal(ForceMajor::False))
6347                .await
6348                .expect("flush failed");
6349
6350            (store_object_id, root_dir_id, unwrapped_key)
6351        };
6352
6353        // Generate mutations for the encrypted mutations object, then serialize and encrypt them
6354        // as the old version (ChaCha20) encrypted mutations object.
6355        let store = fs.object_manager().store(store_object_id).unwrap();
6356        let (old_file_oid, chacha_offset) = {
6357            let mut transaction = store
6358                .new_transaction(
6359                    lock_keys![LockKey::object(store_object_id, root_dir_id)],
6360                    Options::default(),
6361                )
6362                .await
6363                .expect("new_transaction failed");
6364            let root_directory =
6365                Directory::open(&store, root_dir_id).await.expect("open root dir failed");
6366            let old_file = root_directory
6367                .create_child_file(&mut transaction, "old_file")
6368                .await
6369                .expect("create_child_file failed");
6370            let oid = old_file.object_id();
6371
6372            let mutations = transaction.take_mutations();
6373            let mut old_data = Vec::new();
6374            for item in &mutations {
6375                item.mutation.serialize_into(&mut old_data).expect("serialize failed");
6376            }
6377
6378            let mut stream_cipher = StreamCipher::new(&unwrapped_key, 0);
6379            stream_cipher.encrypt(&mut old_data);
6380
6381            let old_version = Version { major: 56, minor: 0 };
6382            let encrypted_mutations = EncryptedMutations {
6383                transactions: vec![(
6384                    JournalCheckpoint { file_offset: 0, checksum: 0, version: old_version },
6385                    mutations.len() as u64,
6386                )],
6387                data: old_data,
6388                mutations_key_roll: Vec::new(),
6389            };
6390
6391            // Write `encrypted_mutations` to parent store using DirectWriter and attach to StoreInfo.
6392            let parent_store = store.parent_store().unwrap();
6393            let mut create_txn = parent_store
6394                .new_transaction(lock_keys![], Options::default())
6395                .await
6396                .expect("new_transaction failed");
6397            let handle = ObjectStore::create_object(
6398                &parent_store,
6399                &mut create_txn,
6400                HandleOptions::default(),
6401                None,
6402            )
6403            .await
6404            .expect("create_object failed");
6405            create_txn.commit().await.expect("commit failed");
6406
6407            let mut writer = DirectWriter::new(&handle, Options::default()).await;
6408            let mut buffer = Vec::new();
6409            encrypted_mutations.serialize_with_version(&mut buffer).expect("serialize failed");
6410            writer.write_bytes(&buffer).await.expect("write_bytes failed");
6411            writer.complete().await.expect("writer complete failed");
6412
6413            let mut store_info = store.load_store_info().await.unwrap();
6414            store_info.encrypted_mutations_object_id = handle.object_id();
6415            store_info.mutations_cipher_offset = 0;
6416
6417            let mut update_txn = parent_store
6418                .new_transaction(
6419                    lock_keys![LockKey::object(
6420                        parent_store.store_object_id(),
6421                        store.store_info_handle_object_id().unwrap(),
6422                    )],
6423                    Options::default(),
6424                )
6425                .await
6426                .expect("new_transaction failed");
6427            store
6428                .write_store_info(&mut update_txn, &store_info)
6429                .await
6430                .expect("write_store_info failed");
6431            *store.store_info.lock() = Some(store_info);
6432            update_txn.commit().await.expect("commit failed");
6433
6434            (oid, stream_cipher.offset())
6435        };
6436
6437        // Roll the mutations key for AES picking up where ChaCha20 left off, and write a
6438        // transaction in the journal.
6439        let (new_wrapped_key, new_unwrapped_key) = crypt
6440            .create_key(store_object_id, KeyPurpose::Metadata)
6441            .await
6442            .expect("create_key failed");
6443        if let LockState::Unlocked { mutations_cipher, .. } = &mut *store.lock_state.lock() {
6444            *mutations_cipher = JournalCipher::new_aes256_xts(&new_unwrapped_key, chacha_offset);
6445        } else {
6446            panic!("Unexpected lock state");
6447        }
6448        store.store_info.lock().as_mut().unwrap().mutations_key = Some(new_wrapped_key);
6449
6450        let new_file_oid = {
6451            let mut transaction = store
6452                .new_transaction(
6453                    lock_keys![LockKey::object(store_object_id, root_dir_id)],
6454                    Options::default(),
6455                )
6456                .await
6457                .expect("new_transaction failed");
6458            let root_directory =
6459                Directory::open(&store, root_dir_id).await.expect("open root dir failed");
6460            let new_file = root_directory
6461                .create_child_file(&mut transaction, "new_file")
6462                .await
6463                .expect("create_child_file failed");
6464            transaction.commit().await.expect("commit failed");
6465            new_file.object_id()
6466        };
6467
6468        // Close and reopen the filesystem *without* flushing the child store.
6469        fs.close().await.expect("close failed");
6470        let device = fs.take_device().await;
6471        device.reopen(false);
6472        let fs = FxFilesystem::open(device).await.expect("FS open failed");
6473
6474        // Mount/unlock the encrypted volume and verify that both files exist and can be opened.
6475        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
6476        let volume = root_volume
6477            .volume("test", StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() })
6478            .await
6479            .expect("volume failed");
6480
6481        let root_directory =
6482            Directory::open(&volume, root_dir_id).await.expect("open root dir failed");
6483
6484        assert_eq!(
6485            root_directory.lookup("old_file").await.expect("lookup old_file failed"),
6486            Some((old_file_oid, ObjectDescriptor::File, false))
6487        );
6488
6489        assert_eq!(
6490            root_directory.lookup("new_file").await.expect("lookup new_file failed"),
6491            Some((new_file_oid, ObjectDescriptor::File, false))
6492        );
6493
6494        let old_file =
6495            ObjectStore::open_object(&volume, old_file_oid, HandleOptions::default(), None)
6496                .await
6497                .expect("open old_file failed");
6498        assert_eq!(old_file.object_id(), old_file_oid);
6499
6500        let new_file =
6501            ObjectStore::open_object(&volume, new_file_oid, HandleOptions::default(), None)
6502                .await
6503                .expect("open new_file failed");
6504        assert_eq!(new_file.object_id(), new_file_oid);
6505    }
6506
6507    #[test_case(false; "unencrypted")]
6508    #[test_case(true; "encrypted")]
6509    #[fuchsia::test]
6510    async fn test_large_mutations_in_transaction(encrypted: bool) {
6511        let fs = test_filesystem().await;
6512        let crypt = Arc::new(new_insecure_crypt());
6513
6514        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
6515        let store = root_volume
6516            .new_volume(
6517                "test",
6518                NewChildStoreOptions {
6519                    options: StoreOptions {
6520                        crypt: if encrypted { Some(crypt.clone()) } else { None },
6521                        ..StoreOptions::default()
6522                    },
6523                    ..Default::default()
6524                },
6525            )
6526            .await
6527            .expect("new_volume failed");
6528        let store_object_id = store.store_object_id();
6529        let root_dir_id = store.root_directory_object_id();
6530
6531        let root_dir = Directory::open(&store, root_dir_id).await.expect("open root dir failed");
6532        let mut transaction = store
6533            .new_transaction(
6534                lock_keys![LockKey::object(store_object_id, root_dir_id)],
6535                Options::default(),
6536            )
6537            .await
6538            .expect("new_transaction failed");
6539
6540        let file = root_dir
6541            .create_child_file(&mut transaction, "test_file")
6542            .await
6543            .expect("create_child_file failed");
6544
6545        // Add multiple inline extended attributes to the file, limiting each to
6546        // MAX_INLINE_XATTR_SIZE, until the total serialized size exceeds
6547        // DEFAULT_MAX_SERIALIZED_RECORD_SIZE.
6548        let mut total_size = 0;
6549        let mut i = 0;
6550        while total_size <= DEFAULT_MAX_SERIALIZED_RECORD_SIZE {
6551            let name = format!("attr{i}").into_bytes();
6552            let mutation = Mutation::replace_or_insert_object(
6553                ObjectKey::extended_attribute(file.object_id(), name),
6554                ObjectValue::inline_extended_attribute(vec![0u8; MAX_INLINE_XATTR_SIZE]),
6555            );
6556            let mut buf = Vec::new();
6557            mutation.serialize_into(&mut buf).unwrap();
6558            total_size += buf.len() as u64;
6559            transaction.add(store_object_id, mutation);
6560            i += 1;
6561        }
6562
6563        transaction.commit().await.expect("commit should succeed");
6564
6565        if encrypted {
6566            store.lock().await.expect("lock failed");
6567            store.unlock(crypt).await.expect("unlock failed");
6568        }
6569
6570        for j in 0..i {
6571            let name = format!("attr{j}").into_bytes();
6572            let item = store
6573                .tree
6574                .find(&ObjectKey::extended_attribute(file.object_id(), name))
6575                .await
6576                .expect("find failed")
6577                .expect("attr not found");
6578            assert_eq!(
6579                item.value,
6580                ObjectValue::inline_extended_attribute(vec![0u8; MAX_INLINE_XATTR_SIZE])
6581            );
6582        }
6583    }
6584}