Skip to main content

fxfs/object_store/
transaction.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
5use crate::checksum::Checksum;
6use crate::filesystem::FxFilesystem;
7use crate::log::*;
8use crate::lsm_tree::types::Item;
9use crate::object_handle::INVALID_OBJECT_ID;
10use crate::object_store::allocator::{AllocatorItem, Hold, Reservation};
11use crate::object_store::object_manager::{ObjectManager, reserved_space_from_journal_usage};
12use crate::object_store::object_record::{
13    BytesAndNodes, FxfsKey, FxfsKeyV49, ObjectItem, ObjectItemV56, ObjectItemV59, ObjectKey,
14    ObjectKeyData, ObjectValue, ProjectProperty,
15};
16use crate::object_store::{AttributeId, AttributeKey, ProjectId};
17use crate::serialized_types::{Migrate, Versioned, migrate_nodefault, migrate_to_version};
18use anyhow::Error;
19use either::{Either, Left, Right};
20use fprint::TypeFingerprint;
21use fuchsia_sync::Mutex;
22use futures::future::poll_fn;
23use futures::pin_mut;
24use rustc_hash::FxHashMap as HashMap;
25use scopeguard::ScopeGuard;
26use serde::{Deserialize, Serialize};
27use std::cell::UnsafeCell;
28use std::cmp::Ordering;
29use std::collections::hash_map::Entry;
30use std::collections::{BTreeSet, btree_set};
31use std::iter::Peekable;
32use std::marker::PhantomPinned;
33use std::ops::{Deref, DerefMut, Range};
34use std::sync::Arc;
35use std::task::{Poll, Waker};
36use std::{fmt, mem};
37
38/// How metadata and allocator reservations should be handled for a transaction.
39///
40/// A transaction consumes space in two different ways:
41///
42/// 1. The transaction will consume space in the journal and eventually in an LSM tree (the
43///    "metadata" cost of the transaction).
44/// 2. The transaction might also include additional allocations or deallocations.  These may come
45///    directly from the global allocator, or they might come from a Reservation object, depending
46///    on the context.  (See specific enum cases for details.)
47///
48/// The basic decision tree for which variant to use is as follows:
49///
50/// - Is there an existing reservation which the transaction should draw from?
51///   - YES: Use [`ReservationOptions::Hold`]
52///   -  NO: Is the transaction allocating or deallocating extents for internal filesystem metadata
53///          (e.g. LSM tree compaction, journal/superblock growth, or graveyard tombstoning in the
54///          root/root-parent stores)?
55///          - YES: Use [`ReservationOptions::BorrowedMetadataAndData`]
56///          -  NO: Is the transaction expected to be space-neutral or space-saving?
57///                 - YES: Use [`ReservationOptions::BorrowedMetadata`]
58///                 -  NO: Use [`ReservationOptions::New`]
59#[derive(Clone, Copy, Default)]
60pub enum ReservationOptions<'a> {
61    /// A new reservation is created to accommodate the metadata of a maximally-sized transaction,
62    /// and the transaction's metadata is paid for from this reservation.  Allocations and
63    /// deallocations within the transaction will go directly to the global allocator (and therefore
64    /// might fail if there's insufficient space).
65    ///
66    /// This should be the default case for most transactions which permit failure when there's no
67    /// space left in the global allocator.
68    #[default]
69    New,
70
71    /// Metadata space is borrowed directly from the global reservation held by the ObjectManager.
72    /// Deallocations in the transaction will be returned to the allocator, rather than to the
73    /// global reservation.
74    ///
75    /// This must be used for transactions which will either not affect net space usage after
76    /// compaction (e.g. setting attributes on an object), or will reduce space (e.g. unlinking a
77    /// file or deleting a range of extents).  By extension, this should not be used for
78    /// transactions which include allocations.
79    ///
80    /// When this is used appropriately, it is guaranteed that transactions will not fail due to
81    /// lack of space.
82    ///
83    /// When choosing between this and BorrowedMetadataAndData, the main consideration is whether
84    /// the freed space from a deallocation must go back to the global reservation, or whether it
85    /// can go back to the global pool.  An example of where it must go back to the global
86    /// reservation would be when purging a layer file during compaction; that space was borrowed
87    /// from the metadata reservation, so it needs to go back there as well (and in that case
88    /// BorrowedMetadataAndData should be used instead).
89    BorrowedMetadata,
90
91    /// Metadata space is borrowed directly from the global reservation held by the ObjectManager.
92    /// Allocations and deallocations in the transaction will be from the global reservation as
93    /// well.
94    ///
95    /// This must be used for transactions that allocate or deallocate extents for internal
96    /// filesystem metadata, such as journal and superblock growth, LSM tree compaction, and
97    /// graveyard tombstoning of objects in the root and root-parent stores.  These transactions
98    /// will never fail due to running out of space (barring any bugs resulting in a deficit in the
99    /// global reservation).
100    BorrowedMetadataAndData,
101
102    /// Use the specified reservation for the transaction's metadata (placing a hold on it) as well
103    /// as any allocations or deallocations in the transaction.
104    Hold(&'a Reservation),
105}
106
107/// This allows for special handling of certain transactions such as deletes and the
108/// extension of Journal extents. For most other use cases it is appropriate to use
109/// `default()` here.
110#[derive(Clone, Copy, Default)]
111pub struct Options<'a> {
112    /// If true, don't check for low journal space.  This should be true for any transactions that
113    /// might alleviate journal space (i.e. compaction).
114    pub skip_journal_checks: bool,
115
116    /// How metadata and allocator reservations should be handled for the transaction.
117    pub reservation: ReservationOptions<'a>,
118
119    /// If set, indicates that this transaction is nested within `parent_transaction` and shares
120    /// its in-flight transaction slot.  The use case for this is a nested transaction, where a
121    /// separate transaction must be committed while constructing another transaction (for example,
122    /// when rolling the object ID cipher).
123    ///
124    /// NOTE: The locks acquired by the nested transaction do not cooperate in any way with the
125    /// locks already held by parent_transaction, so care must still be taken to avoid deadlocks.
126    pub parent_transaction: Option<&'a Transaction<'a>>,
127}
128
129// This is the amount of space that we reserve for metadata when we are creating a new transaction.
130// A transaction should not take more than this.  This is expressed in terms of space occupied in
131// the journal; transactions must not take up more space in the journal than the number below.  The
132// amount chosen here must be large enough for the maximum possible transaction that can be created,
133// so transactions always need to be bounded which might involve splitting an operation up into
134// smaller transactions.
135pub const TRANSACTION_MAX_JOURNAL_USAGE: u64 = 24_576;
136pub const TRANSACTION_METADATA_MAX_AMOUNT: u64 =
137    reserved_space_from_journal_usage(TRANSACTION_MAX_JOURNAL_USAGE);
138
139#[must_use]
140pub struct TransactionLocks<'a>(pub WriteGuard<'a>);
141
142/// The journal consists of these records which will be replayed at mount time.  Within a
143/// transaction, these are stored as a set which allows some mutations to be deduplicated and found
144/// (and we require custom comparison functions below).  For example, we need to be able to find
145/// object size changes.
146pub type Mutation = MutationV59;
147
148#[derive(
149    Clone, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize, TypeFingerprint, Versioned,
150)]
151#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
152pub enum MutationV59 {
153    ObjectStore(ObjectStoreMutationV59),
154    EncryptedObjectStore(#[serde(with = "crate::zerocopy_serialization")] Box<[u8]>),
155    Allocator(AllocatorMutationV32),
156    /// Indicates the beginning of a flush. This would typically involve sealing a tree.
157    BeginFlush,
158    /// Indicates the end of a flush. This would typically involve replacing the immutable layers
159    /// with compacted ones.
160    EndFlush,
161    /// Volume has been deleted. Requires we remove it from the set of managed ObjectStore.
162    DeleteVolume,
163    UpdateBorrowed(u64),
164    UpdateMutationsKey(UpdateMutationsKey),
165    CreateInternalDir(u64),
166}
167
168#[derive(Migrate, Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
169#[migrate_to_version(MutationV59)]
170pub enum MutationV57 {
171    ObjectStore(ObjectStoreMutationV56),
172    EncryptedObjectStore(#[serde(with = "crate::zerocopy_serialization")] Box<[u8]>),
173    Allocator(AllocatorMutationV32),
174    BeginFlush,
175    EndFlush,
176    DeleteVolume,
177    UpdateBorrowed(u64),
178    UpdateMutationsKey(UpdateMutationsKey),
179    CreateInternalDir(u64),
180}
181
182#[derive(Migrate, Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
183#[migrate_to_version(MutationV57)]
184pub enum MutationV56 {
185    ObjectStore(ObjectStoreMutationV56),
186    EncryptedObjectStore(#[serde(with = "crate::zerocopy_serialization")] Box<[u8]>),
187    Allocator(AllocatorMutationV32),
188    BeginFlush,
189    EndFlush,
190    DeleteVolume,
191    UpdateBorrowed(u64),
192    UpdateMutationsKey(UpdateMutationsKey),
193    CreateInternalDir(u64),
194}
195
196impl Mutation {
197    pub fn insert_object(key: ObjectKey, value: ObjectValue) -> Self {
198        Mutation::ObjectStore(ObjectStoreMutation {
199            item: Item::new(key, value),
200            op: Operation::Insert,
201        })
202    }
203
204    pub fn replace_or_insert_object(key: ObjectKey, value: ObjectValue) -> Self {
205        Mutation::ObjectStore(ObjectStoreMutation {
206            item: Item::new(key, value),
207            op: Operation::ReplaceOrInsert,
208        })
209    }
210
211    pub fn merge_object(key: ObjectKey, value: ObjectValue) -> Self {
212        Mutation::ObjectStore(ObjectStoreMutation {
213            item: Item::new(key, value),
214            op: Operation::Merge,
215        })
216    }
217
218    pub fn update_mutations_key(key: FxfsKey) -> Self {
219        Mutation::UpdateMutationsKey(key.into())
220    }
221}
222
223// We have custom comparison functions for mutations that just use the key, rather than the key and
224// value that would be used by default so that we can deduplicate and find mutations (see
225// get_object_mutation below).
226pub type ObjectStoreMutation = ObjectStoreMutationV59;
227
228#[derive(Clone, Debug, Serialize, Deserialize, TypeFingerprint)]
229#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
230pub struct ObjectStoreMutationV59 {
231    pub item: ObjectItemV59,
232    pub op: Operation,
233}
234
235#[derive(Migrate, Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
236#[migrate_to_version(ObjectStoreMutationV59)]
237#[migrate_nodefault]
238pub struct ObjectStoreMutationV56 {
239    pub item: ObjectItemV56,
240    pub op: Operation,
241}
242
243/// The different LSM tree operations that can be performed as part of a mutation.
244pub type Operation = OperationV32;
245
246#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, TypeFingerprint)]
247#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
248pub enum OperationV32 {
249    Insert,
250    ReplaceOrInsert,
251    Merge,
252}
253
254impl Ord for ObjectStoreMutation {
255    fn cmp(&self, other: &Self) -> Ordering {
256        self.item.key.cmp(&other.item.key)
257    }
258}
259
260impl PartialOrd for ObjectStoreMutation {
261    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
262        Some(self.cmp(other))
263    }
264}
265
266impl PartialEq for ObjectStoreMutation {
267    fn eq(&self, other: &Self) -> bool {
268        self.item.key.eq(&other.item.key)
269    }
270}
271
272impl Eq for ObjectStoreMutation {}
273
274impl Ord for AllocatorItem {
275    fn cmp(&self, other: &Self) -> Ordering {
276        self.key.cmp(&other.key)
277    }
278}
279
280impl PartialOrd for AllocatorItem {
281    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
282        Some(self.cmp(other))
283    }
284}
285
286/// Same as std::ops::Range but with Ord and PartialOrd support, sorted first by start of the range,
287/// then by the end.
288pub type DeviceRange = DeviceRangeV32;
289
290#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TypeFingerprint)]
291#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
292pub struct DeviceRangeV32(pub Range<u64>);
293
294impl Deref for DeviceRange {
295    type Target = Range<u64>;
296
297    fn deref(&self) -> &Self::Target {
298        &self.0
299    }
300}
301
302impl DerefMut for DeviceRange {
303    fn deref_mut(&mut self) -> &mut Self::Target {
304        &mut self.0
305    }
306}
307
308impl From<Range<u64>> for DeviceRange {
309    fn from(range: Range<u64>) -> Self {
310        Self(range)
311    }
312}
313
314impl Into<Range<u64>> for DeviceRange {
315    fn into(self) -> Range<u64> {
316        self.0
317    }
318}
319
320impl Ord for DeviceRange {
321    fn cmp(&self, other: &Self) -> Ordering {
322        self.start.cmp(&other.start).then(self.end.cmp(&other.end))
323    }
324}
325
326impl PartialOrd for DeviceRange {
327    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
328        Some(self.cmp(other))
329    }
330}
331
332pub type AllocatorMutation = AllocatorMutationV32;
333
334#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize, TypeFingerprint)]
335#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
336pub enum AllocatorMutationV32 {
337    Allocate {
338        device_range: DeviceRangeV32,
339        owner_object_id: u64,
340    },
341    Deallocate {
342        device_range: DeviceRangeV32,
343        owner_object_id: u64,
344    },
345    SetLimit {
346        owner_object_id: u64,
347        bytes: u64,
348    },
349    /// Marks all extents with a given owner_object_id for deletion.
350    /// Used to free space allocated to encrypted ObjectStore where we may not have the key.
351    /// Note that the actual deletion time is undefined so this should never be used where an
352    /// ObjectStore is still in use due to a high risk of corruption. Similarly, owner_object_id
353    /// should never be reused for the same reasons.
354    MarkForDeletion(u64),
355}
356
357pub type UpdateMutationsKey = UpdateMutationsKeyV49;
358
359#[derive(Clone, Debug, Serialize, Deserialize, TypeFingerprint)]
360pub struct UpdateMutationsKeyV49(pub FxfsKeyV49);
361
362impl From<UpdateMutationsKey> for FxfsKey {
363    fn from(outer: UpdateMutationsKey) -> Self {
364        outer.0
365    }
366}
367
368impl From<FxfsKey> for UpdateMutationsKey {
369    fn from(inner: FxfsKey) -> Self {
370        Self(inner)
371    }
372}
373
374#[cfg(fuzz)]
375impl<'a> arbitrary::Arbitrary<'a> for UpdateMutationsKey {
376    fn arbitrary(u: &mut arbitrary::Unstructured<'a>) -> arbitrary::Result<Self> {
377        Ok(UpdateMutationsKey::from(FxfsKey::arbitrary(u).unwrap()))
378    }
379}
380
381impl Ord for UpdateMutationsKey {
382    fn cmp(&self, other: &Self) -> Ordering {
383        (self as *const UpdateMutationsKey).cmp(&(other as *const _))
384    }
385}
386
387impl PartialOrd for UpdateMutationsKey {
388    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
389        Some(self.cmp(other))
390    }
391}
392
393impl Eq for UpdateMutationsKey {}
394
395impl PartialEq for UpdateMutationsKey {
396    fn eq(&self, other: &Self) -> bool {
397        std::ptr::eq(self, other)
398    }
399}
400
401/// When creating a transaction, locks typically need to be held to prevent two or more writers
402/// trying to make conflicting mutations at the same time.  LockKeys are used for this.
403/// NOTE: Ordering is important here!  The lock manager sorts the list of locks in a transaction
404/// to acquire them in a consistent order, but there is a special case for the Flush lock.
405/// The Flush lock is taken when we flush an LSM tree (e.g. an object store), and is held for
406/// several transactions.  As such, it must come first in the lock acquisition ordering, so that
407/// other transactions using the Flush lock have the same ordering as in flushing.
408#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Copy)]
409pub enum LockKey {
410    /// Used to lock flushing an object.
411    Flush {
412        object_id: u64,
413    },
414
415    /// Used to lock changes to a particular object attribute (e.g. writes).
416    ObjectAttribute {
417        store_object_id: u64,
418        object_id: u64,
419        attribute_id: AttributeId,
420    },
421
422    /// Used to lock changes to a particular object (e.g. adding a child to a directory).
423    Object {
424        store_object_id: u64,
425        object_id: u64,
426    },
427
428    ProjectId {
429        store_object_id: u64,
430        project_id: ProjectId,
431    },
432
433    /// Used to lock any truncate operations for a file.
434    Truncate {
435        store_object_id: u64,
436        object_id: u64,
437    },
438
439    /// A lock used when getting or creating the internal directory.
440    InternalDirectory {
441        store_object_id: u64,
442    },
443
444    /// Used to serialize pre caching of keys.  The lock ordering is different for this: it is
445    /// acquired in `Filesystem::commit_transaction` and happens *after* other keys have been
446    /// promoted to write locks, but before the commit lock.
447    PreCacheKeys {
448        store_object_id: u64,
449    },
450}
451
452impl LockKey {
453    pub const fn object_attribute(
454        store_object_id: u64,
455        object_id: u64,
456        attribute_id: AttributeId,
457    ) -> Self {
458        LockKey::ObjectAttribute { store_object_id, object_id, attribute_id }
459    }
460
461    pub const fn object(store_object_id: u64, object_id: u64) -> Self {
462        LockKey::Object { store_object_id, object_id }
463    }
464
465    pub const fn flush(object_id: u64) -> Self {
466        LockKey::Flush { object_id }
467    }
468
469    pub const fn truncate(store_object_id: u64, object_id: u64) -> Self {
470        LockKey::Truncate { store_object_id, object_id }
471    }
472
473    pub const fn pre_cache_keys(store_object_id: u64) -> Self {
474        LockKey::PreCacheKeys { store_object_id }
475    }
476}
477
478/// A container for holding `LockKey` objects. Can store a single `LockKey` inline.
479#[derive(Clone, Debug)]
480pub enum LockKeys {
481    None,
482    Inline(LockKey),
483    Vec(Vec<LockKey>),
484}
485
486impl LockKeys {
487    pub fn with_capacity(capacity: usize) -> Self {
488        if capacity > 1 { LockKeys::Vec(Vec::with_capacity(capacity)) } else { LockKeys::None }
489    }
490
491    pub fn push(&mut self, key: LockKey) {
492        match self {
493            Self::None => *self = LockKeys::Inline(key),
494            Self::Inline(inline) => {
495                *self = LockKeys::Vec(vec![*inline, key]);
496            }
497            Self::Vec(vec) => vec.push(key),
498        }
499    }
500
501    pub fn truncate(&mut self, len: usize) {
502        match self {
503            Self::None => {}
504            Self::Inline(_) => {
505                if len == 0 {
506                    *self = Self::None;
507                }
508            }
509            Self::Vec(vec) => vec.truncate(len),
510        }
511    }
512
513    fn len(&self) -> usize {
514        match self {
515            Self::None => 0,
516            Self::Inline(_) => 1,
517            Self::Vec(vec) => vec.len(),
518        }
519    }
520
521    fn contains(&self, key: &LockKey) -> bool {
522        match self {
523            Self::None => false,
524            Self::Inline(single) => single == key,
525            Self::Vec(vec) => vec.contains(key),
526        }
527    }
528
529    fn sort_unstable(&mut self) {
530        match self {
531            Self::Vec(vec) => vec.sort_unstable(),
532            _ => {}
533        }
534    }
535
536    fn dedup(&mut self) {
537        match self {
538            Self::Vec(vec) => vec.dedup(),
539            _ => {}
540        }
541    }
542
543    fn iter(&self) -> LockKeysIter<'_> {
544        match self {
545            LockKeys::None => LockKeysIter::None,
546            LockKeys::Inline(key) => LockKeysIter::Inline(key),
547            LockKeys::Vec(keys) => LockKeysIter::Vec(keys.iter()),
548        }
549    }
550}
551
552enum LockKeysIter<'a> {
553    None,
554    Inline(&'a LockKey),
555    Vec(std::slice::Iter<'a, LockKey>),
556}
557
558impl<'a> Iterator for LockKeysIter<'a> {
559    type Item = &'a LockKey;
560    fn next(&mut self) -> Option<Self::Item> {
561        match self {
562            Self::None => None,
563            Self::Inline(inline) => {
564                let next = *inline;
565                *self = Self::None;
566                Some(next)
567            }
568            Self::Vec(vec) => vec.next(),
569        }
570    }
571}
572
573impl Default for LockKeys {
574    fn default() -> Self {
575        LockKeys::None
576    }
577}
578
579#[macro_export]
580macro_rules! lock_keys {
581    () => {
582        $crate::object_store::transaction::LockKeys::None
583    };
584    ($lock_key:expr $(,)?) => {
585        $crate::object_store::transaction::LockKeys::Inline($lock_key)
586    };
587    ($($lock_keys:expr),+ $(,)?) => {
588        $crate::object_store::transaction::LockKeys::Vec(vec![$($lock_keys),+])
589    };
590}
591pub use lock_keys;
592
593/// Mutations in a transaction can be associated with an object so that when mutations are applied,
594/// updates can be applied to in-memory structures. For example, we cache object sizes, so when a
595/// size change is applied, we can update the cached object size.
596pub trait AssociatedObject: Send + Sync {
597    fn will_apply_mutation(&self, _mutation: &Mutation, _object_id: u64, _manager: &ObjectManager) {
598    }
599}
600
601pub enum AssocObj<'a> {
602    None,
603    Borrowed(&'a dyn AssociatedObject),
604    Owned(Box<dyn AssociatedObject>),
605}
606
607impl AssocObj<'_> {
608    pub fn map<R, F: FnOnce(&dyn AssociatedObject) -> R>(&self, f: F) -> Option<R> {
609        match self {
610            AssocObj::None => None,
611            AssocObj::Borrowed(b) => Some(f(*b)),
612            AssocObj::Owned(o) => Some(f(o.as_ref())),
613        }
614    }
615}
616
617pub struct TxnMutation<'a> {
618    // This, at time of writing, is either the object ID of an object store, or the object ID of the
619    // allocator.  In the case of an object mutation, there's another object ID in the mutation
620    // record that would be for the object actually being changed.
621    pub object_id: u64,
622
623    // The actual mutation.  This gets serialized to the journal.
624    pub mutation: Mutation,
625
626    // An optional associated object for the mutation.  During replay, there will always be no
627    // associated object.
628    pub associated_object: AssocObj<'a>,
629}
630
631// We store TxnMutation in a set, and for that, we only use object_id and mutation and not the
632// associated object or checksum.
633//
634// WARNING: It is critical that `object_id` (which corresponds to the store ID for store mutations)
635// remains the primary key for sorting. The commit pipeline (`commit_transaction`) relies on
636// the mutations being sorted by `object_id` to acquire store locks in a consistent order,
637// preventing deadlocks.
638impl Ord for TxnMutation<'_> {
639    fn cmp(&self, other: &Self) -> Ordering {
640        self.object_id.cmp(&other.object_id).then_with(|| self.mutation.cmp(&other.mutation))
641    }
642}
643
644impl PartialOrd for TxnMutation<'_> {
645    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
646        Some(self.cmp(other))
647    }
648}
649
650impl PartialEq for TxnMutation<'_> {
651    fn eq(&self, other: &Self) -> bool {
652        self.object_id.eq(&other.object_id) && self.mutation.eq(&other.mutation)
653    }
654}
655
656impl Eq for TxnMutation<'_> {}
657
658impl std::fmt::Debug for TxnMutation<'_> {
659    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
660        f.debug_struct("TxnMutation")
661            .field("object_id", &self.object_id)
662            .field("mutation", &self.mutation)
663            .finish()
664    }
665}
666
667/// An iterator over mutations belonging to a single object store within a transaction.
668/// It wraps a mutable reference to a peekable iterator over transaction mutations and yields
669/// mutations until the object ID changes from the object ID peeked at creation time.
670pub struct ObjectMutationIterator<'a, 'b> {
671    iter: &'a mut Peekable<btree_set::Iter<'b, TxnMutation<'b>>>,
672    object_id: u64,
673}
674
675impl<'a, 'b> ObjectMutationIterator<'a, 'b> {
676    pub fn new(iter: &'a mut Peekable<btree_set::Iter<'b, TxnMutation<'b>>>) -> Option<Self> {
677        let object_id = iter.peek()?.object_id;
678        Some(Self { iter, object_id })
679    }
680
681    pub fn object_id(&self) -> u64 {
682        self.object_id
683    }
684}
685
686impl<'b> Iterator for ObjectMutationIterator<'_, 'b> {
687    type Item = &'b Mutation;
688
689    fn next(&mut self) -> Option<Self::Item> {
690        if self.iter.peek().is_some_and(|m| m.object_id == self.object_id) {
691            Some(&self.iter.next().unwrap().mutation)
692        } else {
693            None
694        }
695    }
696}
697
698impl Drop for ObjectMutationIterator<'_, '_> {
699    // Need to have a drop to iterate to the end so that we always finish the current object before
700    // releasing the borrow.
701    fn drop(&mut self) {
702        for _ in self.by_ref() {}
703    }
704}
705
706pub enum MetadataReservation<'a> {
707    /// Metadata space for this transaction is being borrowed from ObjectManager's metadata
708    /// reservation, but allocations/deallocations in the transaction do not use it.
709    BorrowedMetadata,
710
711    /// Metadata space for this transaction is being borrowed from ObjectManager's metadata
712    /// reservation, and allocations/deallocations in the transaction also use ObjectManager's
713    /// metadata reservation.
714    BorrowedMetadataAndData,
715
716    /// A metadata reservation was made when the transaction was created.
717    Reservation(Reservation),
718
719    /// The metadata space is being _held_ within an existing reservation.
720    Hold(Hold<'a>),
721}
722
723/// A transaction groups mutation records to be committed as a group.
724pub struct Transaction<'a> {
725    fs: Arc<FxFilesystem>,
726
727    // The mutations that make up this transaction.
728    mutations: BTreeSet<TxnMutation<'a>>,
729
730    // The locks that this transaction currently holds.
731    txn_locks: LockKeys,
732
733    /// The reservation for the metadata for this transaction.
734    pub metadata_reservation: MetadataReservation<'a>,
735
736    // Keep track of objects explicitly created by this transaction. No locks are required for them.
737    // Addressed by (owner_object_id, object_id).
738    new_objects: BTreeSet<(u64, u64)>,
739
740    /// Any data checksums which should be evaluated when replaying this transaction.
741    checksums: Vec<(Range<u64>, Vec<Checksum>, bool)>,
742
743    /// Set if this transaction contains data (i.e. includes any extent mutations).
744    includes_write: bool,
745
746    /// If true, don't check for low journal space.
747    pub skip_journal_checks: bool,
748
749    /// If set, indicates that this transaction is nested within `parent_transaction` and shares
750    /// its in-flight transaction slot.
751    parent_transaction: Option<&'a Transaction<'a>>,
752}
753
754impl<'a> Transaction<'a> {
755    /// Creates a new transaction.  `txn_locks` are read locks that can be upgraded to write locks
756    /// at commit time.
757    pub async fn new(
758        fs: Arc<FxFilesystem>,
759        options: Options<'a>,
760        txn_locks: LockKeys,
761    ) -> Result<Transaction<'a>, Error> {
762        if options.parent_transaction.is_none() {
763            fs.add_transaction(options.skip_journal_checks).await;
764        }
765        // We support three options for metadata space reservation:
766        //
767        //   1. We can borrow from the filesystem's metadata reservation.  This should only be
768        //      be used on the understanding that eventually, potentially after a full compaction,
769        //      there should be no net increase in space used.  For example, unlinking an object
770        //      should eventually decrease the amount of space used and setting most attributes
771        //      should not result in any change.
772        //
773        //   2. A reservation is provided in which case we'll place a hold on some of it for
774        //      metadata.
775        //
776        //   3. No reservation is supplied, so we try and reserve space with the allocator now,
777        //      and will return NoSpace if that fails.
778        let metadata_reservation = if fs.options().image_builder_mode.is_some() {
779            MetadataReservation::BorrowedMetadata
780        } else {
781            match options.reservation {
782                ReservationOptions::New => {
783                    MetadataReservation::Reservation(fs.allocator().reserve(None, 0).unwrap())
784                }
785                ReservationOptions::BorrowedMetadata => MetadataReservation::BorrowedMetadata,
786                ReservationOptions::BorrowedMetadataAndData => {
787                    MetadataReservation::BorrowedMetadataAndData
788                }
789                ReservationOptions::Hold(reservation) => {
790                    MetadataReservation::Hold(reservation.reserve(0).unwrap())
791                }
792            }
793        };
794        let mut transaction = Transaction {
795            fs: fs.clone(),
796            mutations: BTreeSet::new(),
797            txn_locks: LockKeys::default(),
798            metadata_reservation,
799            new_objects: BTreeSet::new(),
800            checksums: Vec::new(),
801            includes_write: false,
802            skip_journal_checks: options.skip_journal_checks,
803            parent_transaction: options.parent_transaction,
804        };
805        fs.add_transaction_reservation(&mut transaction).await?;
806
807        transaction.txn_locks = {
808            let lock_manager = fs.lock_manager();
809            let mut write_guard = lock_manager.txn_lock(txn_locks).await;
810            std::mem::take(&mut write_guard.0.lock_keys)
811        };
812        Ok(transaction)
813    }
814
815    /// Returns the allocator reservation to be used for allocations and deallocations in this
816    /// transaction, if any.
817    pub fn allocator_reservation(&self) -> Option<&Reservation> {
818        match &self.metadata_reservation {
819            MetadataReservation::BorrowedMetadataAndData => {
820                Some(self.fs.object_manager().metadata_reservation())
821            }
822            MetadataReservation::Hold(hold) => Some(hold.owner()),
823            MetadataReservation::BorrowedMetadata | MetadataReservation::Reservation(_) => None,
824        }
825    }
826
827    pub fn mutations(&self) -> &BTreeSet<TxnMutation<'a>> {
828        &self.mutations
829    }
830
831    pub fn take_mutations(&mut self) -> BTreeSet<TxnMutation<'a>> {
832        self.new_objects.clear();
833        mem::take(&mut self.mutations)
834    }
835
836    /// Adds a mutation to this transaction.  If the mutation already exists, it is replaced and the
837    /// old mutation is returned.
838    pub fn add(&mut self, object_id: u64, mutation: Mutation) -> Option<Mutation> {
839        self.add_with_object(object_id, mutation, AssocObj::None)
840    }
841
842    /// Adds `delta` to a `BytesAndNodes` merge mutation for `key`, folding into any delta already
843    /// staged in this transaction.  Mutations in a transaction are deduplicated by `ObjectKey`
844    /// alone, so adding a second delta for the same key would otherwise replace the first.
845    pub fn merge_bytes_and_nodes(
846        &mut self,
847        store_object_id: u64,
848        key: ObjectKey,
849        delta: BytesAndNodes,
850    ) {
851        let delta = match self.get_object_mutation(store_object_id, key.clone()) {
852            Some(ObjectStoreMutation {
853                item: Item { value: ObjectValue::BytesAndNodes { bytes, nodes }, .. },
854                ..
855            }) => delta + BytesAndNodes { bytes: *bytes, nodes: *nodes },
856            _ => delta,
857        };
858        if delta.is_zero() {
859            self.remove(store_object_id, Mutation::merge_object(key, ObjectValue::None));
860        } else {
861            self.add_with_object_internal(
862                store_object_id,
863                Mutation::merge_object(key, delta.into()),
864                AssocObj::None,
865            );
866        }
867    }
868
869    /// Removes a mutation that matches `mutation`.
870    pub fn remove(&mut self, object_id: u64, mutation: Mutation) {
871        let txn_mutation = TxnMutation { object_id, mutation, associated_object: AssocObj::None };
872        if self.mutations.remove(&txn_mutation) {
873            if let Mutation::ObjectStore(ObjectStoreMutation {
874                item:
875                    ObjectItem {
876                        key: ObjectKey { object_id: new_object_id, data: ObjectKeyData::Object },
877                        ..
878                    },
879                op: Operation::Insert,
880            }) = txn_mutation.mutation
881            {
882                self.new_objects.remove(&(object_id, new_object_id));
883            }
884        }
885    }
886
887    /// Adds a mutation with an associated object. If the mutation already exists, it is replaced
888    /// and the old mutation is returned.
889    pub fn add_with_object(
890        &mut self,
891        object_id: u64,
892        mutation: Mutation,
893        associated_object: AssocObj<'a>,
894    ) -> Option<Mutation> {
895        debug_assert!(
896            !matches!(
897                mutation,
898                Mutation::ObjectStore(ObjectStoreMutation {
899                    item: ObjectItem {
900                        key: ObjectKey {
901                            data: ObjectKeyData::Project { property: ProjectProperty::Usage, .. },
902                            ..
903                        },
904                        ..
905                    },
906                    ..
907                })
908            ),
909            "Use merge_bytes_and_nodes"
910        );
911        self.add_with_object_internal(object_id, mutation, associated_object)
912    }
913
914    /// Adds a mutation with an associated object. If the mutation already exists, it is replaced
915    /// and the old mutation is returned.
916    fn add_with_object_internal(
917        &mut self,
918        object_id: u64,
919        mutation: Mutation,
920        associated_object: AssocObj<'a>,
921    ) -> Option<Mutation> {
922        assert!(object_id != INVALID_OBJECT_ID);
923        let mut is_allocate = false;
924        match &mutation {
925            Mutation::ObjectStore(ObjectStoreMutation {
926                item:
927                    Item {
928                        key:
929                            ObjectKey {
930                                data: ObjectKeyData::Attribute(_, AttributeKey::Extent(_)), ..
931                            },
932                        ..
933                    },
934                ..
935            }) => {
936                self.includes_write = true;
937            }
938            Mutation::Allocator(AllocatorMutation::Allocate { .. }) => {
939                is_allocate = true;
940            }
941            _ => {}
942        }
943        let txn_mutation = TxnMutation { object_id, mutation, associated_object };
944        self.verify_locks(&txn_mutation);
945        let old = self.mutations.replace(txn_mutation).map(|m| m.mutation);
946        if is_allocate {
947            assert!(
948                !matches!(self.metadata_reservation, MetadataReservation::BorrowedMetadata)
949                    || self.fs.options().image_builder_mode.is_some(),
950                "Allocations are not allowed in BorrowedMetadata transactions"
951            );
952        }
953        old
954    }
955
956    pub fn add_checksum(&mut self, range: Range<u64>, checksums: Vec<Checksum>, first_write: bool) {
957        self.checksums.push((range, checksums, first_write));
958    }
959
960    pub fn includes_write(&self) -> bool {
961        self.includes_write
962    }
963
964    pub fn checksums(&self) -> &[(Range<u64>, Vec<Checksum>, bool)] {
965        &self.checksums
966    }
967
968    pub fn take_checksums(&mut self) -> Vec<(Range<u64>, Vec<Checksum>, bool)> {
969        std::mem::replace(&mut self.checksums, Vec::new())
970    }
971
972    fn verify_locks(&mut self, mutation: &TxnMutation<'_>) {
973        // It was considered to change the locks from Vec to BTreeSet since we'll now be searching
974        // through it, but given the small set that these locks usually comprise, it probably isn't
975        // worth it.
976        match mutation {
977            TxnMutation {
978                mutation:
979                    Mutation::ObjectStore {
980                        0: ObjectStoreMutation { item: ObjectItem { key, .. }, op },
981                    },
982                object_id: store_object_id,
983                ..
984            } => {
985                match &key.data {
986                    ObjectKeyData::Attribute(..) => {
987                        // TODO(https://fxbug.dev/42073914): Check lock requirements.
988                    }
989                    ObjectKeyData::Child { .. }
990                    | ObjectKeyData::EncryptedChild(_)
991                    | ObjectKeyData::EncryptedCasefoldChild(_)
992                    | ObjectKeyData::CasefoldChild { .. }
993                    | ObjectKeyData::LegacyCasefoldChild(_) => {
994                        let id = key.object_id;
995                        if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
996                            && !self.new_objects.contains(&(*store_object_id, id))
997                        {
998                            debug_assert!(
999                                false,
1000                                "Not holding required lock for object {id} \
1001                                in store {store_object_id}"
1002                            );
1003                            error!(
1004                                "Not holding required lock for object {id} in store \
1005                                {store_object_id}"
1006                            )
1007                        }
1008                    }
1009                    ObjectKeyData::GraveyardEntry { .. } => {
1010                        // TODO(https://fxbug.dev/42073911): Check lock requirements.
1011                    }
1012                    ObjectKeyData::GraveyardAttributeEntry { .. } => {
1013                        // TODO(https://fxbug.dev/122974): Check lock requirements.
1014                    }
1015                    ObjectKeyData::Keys => {
1016                        let id = key.object_id;
1017                        if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
1018                            && !self.new_objects.contains(&(*store_object_id, id))
1019                        {
1020                            debug_assert!(
1021                                false,
1022                                "Not holding required lock for object {id} \
1023                                in store {store_object_id}"
1024                            );
1025                            error!(
1026                                "Not holding required lock for object {id} in store \
1027                                {store_object_id}"
1028                            )
1029                        }
1030                    }
1031                    ObjectKeyData::Object => match op {
1032                        // Insert implies the caller expects no object with which to race
1033                        Operation::Insert => {
1034                            self.new_objects.insert((*store_object_id, key.object_id));
1035                        }
1036                        Operation::Merge | Operation::ReplaceOrInsert => {
1037                            let id = key.object_id;
1038                            if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
1039                                && !self.new_objects.contains(&(*store_object_id, id))
1040                            {
1041                                debug_assert!(
1042                                    false,
1043                                    "Not holding required lock for object {id} \
1044                                    in store {store_object_id}"
1045                                );
1046                                error!(
1047                                    "Not holding required lock for object {id} in store \
1048                                    {store_object_id}"
1049                                )
1050                            }
1051                        }
1052                    },
1053                    ObjectKeyData::Project { project_id, property: ProjectProperty::Limit } => {
1054                        if !self.txn_locks.contains(&LockKey::ProjectId {
1055                            store_object_id: *store_object_id,
1056                            project_id: *project_id,
1057                        }) {
1058                            debug_assert!(
1059                                false,
1060                                "Not holding required lock for project limit id {project_id} \
1061                                in store {store_object_id}"
1062                            );
1063                            error!(
1064                                "Not holding required lock for project limit id {project_id} in \
1065                                store {store_object_id}"
1066                            )
1067                        }
1068                    }
1069                    ObjectKeyData::Project { property: ProjectProperty::Usage, .. } => match op {
1070                        Operation::Insert | Operation::ReplaceOrInsert => {
1071                            panic!(
1072                                "Project usage is all handled by merging deltas, no inserts or \
1073                                replacements should be used"
1074                            );
1075                        }
1076                        // Merges are all handled like atomic +/- and serialized by the tree locks.
1077                        Operation::Merge => {}
1078                    },
1079                    ObjectKeyData::ExtendedAttribute { .. } => {
1080                        let id = key.object_id;
1081                        if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
1082                            && !self.new_objects.contains(&(*store_object_id, id))
1083                        {
1084                            debug_assert!(
1085                                false,
1086                                "Not holding required lock for object {id} \
1087                                in store {store_object_id} while mutating extended attribute"
1088                            );
1089                            error!(
1090                                "Not holding required lock for object {id} in store \
1091                                {store_object_id} while mutating extended attribute"
1092                            )
1093                        }
1094                    }
1095                }
1096            }
1097            TxnMutation { mutation: Mutation::DeleteVolume, object_id, .. } => {
1098                if !self.txn_locks.contains(&LockKey::flush(*object_id)) {
1099                    debug_assert!(false, "Not holding required lock for DeleteVolume");
1100                    error!("Not holding required lock for DeleteVolume");
1101                }
1102            }
1103            _ => {}
1104        }
1105    }
1106
1107    /// Returns true if this transaction has no mutations.
1108    pub fn is_empty(&self) -> bool {
1109        self.mutations.is_empty()
1110    }
1111
1112    /// Searches for an existing object mutation within the transaction that has the given key and
1113    /// returns it if found.
1114    pub fn get_object_mutation(
1115        &self,
1116        store_object_id: u64,
1117        key: ObjectKey,
1118    ) -> Option<&ObjectStoreMutation> {
1119        if let Some(TxnMutation { mutation: Mutation::ObjectStore(mutation), .. }) =
1120            self.mutations.get(&TxnMutation {
1121                object_id: store_object_id,
1122                mutation: Mutation::insert_object(key, ObjectValue::None),
1123                associated_object: AssocObj::None,
1124            })
1125        {
1126            Some(mutation)
1127        } else {
1128            None
1129        }
1130    }
1131
1132    /// Commits a transaction.  If successful, returns the journal offset of the transaction.
1133    pub async fn commit(mut self) -> Result<u64, Error> {
1134        debug!(txn:? = &self; "Commit");
1135        self.fs.clone().commit_transaction(&mut self, |x| x).await
1136    }
1137
1138    /// Commits and then runs the callback whilst locks are held.  The callback accepts a single
1139    /// parameter which is the journal offset of the transaction.
1140    pub async fn commit_with_callback<R: Send>(
1141        mut self,
1142        f: impl FnOnce(u64) -> R + Send,
1143    ) -> Result<R, Error> {
1144        debug!(txn:? = &self; "Commit");
1145        self.fs.clone().commit_transaction(&mut self, f).await
1146    }
1147
1148    /// Commits the transaction, but allows the transaction to be used again.  The locks are not
1149    /// dropped (but transaction locks will get downgraded to read locks).
1150    pub async fn commit_and_continue(&mut self) -> Result<(), Error> {
1151        debug!(txn:? = self; "Commit");
1152        self.fs.clone().commit_transaction(self, |_| {}).await?;
1153        assert!(self.mutations.is_empty());
1154        assert!(self.new_objects.is_empty());
1155        assert!(self.checksums.is_empty());
1156        self.includes_write = false;
1157        self.fs.lock_manager().downgrade_locks(&self.txn_locks);
1158        self.fs.clone().add_transaction_reservation(self).await
1159    }
1160
1161    /// Prepares to commit by upgrading transaction locks to write locks and waiting for active
1162    /// readers to finish.
1163    pub async fn commit_prepare(&self) {
1164        self.fs.lock_manager().commit_prepare(self).await;
1165    }
1166}
1167
1168impl Drop for Transaction<'_> {
1169    fn drop(&mut self) {
1170        // Call the filesystem implementation of drop_transaction which should, as a minimum, call
1171        // LockManager's drop_transaction to ensure the locks are released.
1172        debug!(txn:? = &self; "Drop");
1173        if self.parent_transaction.is_none() {
1174            self.fs.sub_transaction();
1175        }
1176        self.fs.clone().drop_transaction(self);
1177    }
1178}
1179
1180impl std::fmt::Debug for Transaction<'_> {
1181    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1182        f.debug_struct("Transaction")
1183            .field("mutations", &self.mutations)
1184            .field("txn_locks", &self.txn_locks)
1185            .field("reservation", &self.allocator_reservation())
1186            .finish()
1187    }
1188}
1189
1190pub enum BorrowedOrOwned<'a, T> {
1191    Borrowed(&'a T),
1192    Owned(T),
1193}
1194
1195impl<T> Deref for BorrowedOrOwned<'_, T> {
1196    type Target = T;
1197
1198    fn deref(&self) -> &Self::Target {
1199        match self {
1200            BorrowedOrOwned::Borrowed(b) => b,
1201            BorrowedOrOwned::Owned(o) => &o,
1202        }
1203    }
1204}
1205
1206impl<'a, T> From<&'a T> for BorrowedOrOwned<'a, T> {
1207    fn from(value: &'a T) -> Self {
1208        BorrowedOrOwned::Borrowed(value)
1209    }
1210}
1211
1212impl<T> From<T> for BorrowedOrOwned<'_, T> {
1213    fn from(value: T) -> Self {
1214        BorrowedOrOwned::Owned(value)
1215    }
1216}
1217
1218/// LockManager holds the locks that transactions might have taken.  A TransactionManager
1219/// implementation would typically have one of these.
1220///
1221/// Three different kinds of locks are supported.  There are read locks and write locks, which are
1222/// as one would expect.  The third kind of lock is a _transaction_ lock (which is also known as an
1223/// upgradeable read lock).  When first acquired, these block other writes (including other
1224/// transaction locks) but do not block reads.  When it is time to commit a transaction, these locks
1225/// are upgraded to full write locks (without ever dropping the lock) and then dropped after
1226/// committing (unless commit_and_continue is used).  This way, reads are only blocked for the
1227/// shortest possible time.  It follows that write locks should be used sparingly.  Locks are
1228/// granted in order with one exception: when a lock is in the initial _transaction_ lock state
1229/// (LockState::Locked), all read locks are allowed even if there are other tasks waiting for the
1230/// lock.  The reason for this is because we allow read locks to be taken by tasks that have taken a
1231/// _transaction_ lock (i.e. recursion is allowed).  In other cases, such as when a writer is
1232/// waiting and there are only readers, readers will queue up behind the writer.
1233///
1234/// To summarize:
1235///
1236/// +-------------------------+-----------------+----------------+------------------+
1237/// |                         | While read_lock | While txn_lock | While write_lock |
1238/// |                         | is held         | is held        | is held          |
1239/// +-------------------------+-----------------+----------------+------------------+
1240/// | Can acquire read_lock?  | true            | true           | false            |
1241/// +-------------------------+-----------------+----------------+------------------+
1242/// | Can acquire txn_lock?   | true            | false          | false            |
1243/// +-------------------------+-----------------+----------------+------------------+
1244/// | Can acquire write_lock? | false           | false          | false            |
1245/// +-------------------------+-----------------+----------------+------------------+
1246pub struct LockManager {
1247    locks: Mutex<Locks>,
1248}
1249
1250struct Locks {
1251    keys: HashMap<LockKey, LockEntry>,
1252}
1253
1254impl Locks {
1255    fn drop_lock(&mut self, key: LockKey, state: LockState) {
1256        if let Entry::Occupied(mut occupied) = self.keys.entry(key) {
1257            let entry = occupied.get_mut();
1258            let wake = match state {
1259                LockState::ReadLock => {
1260                    entry.read_count -= 1;
1261                    entry.read_count == 0
1262                }
1263                // drop_write_locks currently depends on us treating Locked and WriteLock the same.
1264                LockState::Locked | LockState::WriteLock => {
1265                    entry.state = LockState::ReadLock;
1266                    true
1267                }
1268            };
1269            if wake {
1270                // SAFETY: The lock in `LockManager::locks` is held.
1271                unsafe {
1272                    entry.wake();
1273                }
1274                if entry.can_remove() {
1275                    occupied.remove_entry();
1276                }
1277            }
1278        } else {
1279            unreachable!();
1280        }
1281    }
1282
1283    fn drop_read_locks(&mut self, lock_keys: LockKeys) {
1284        for lock in lock_keys.iter() {
1285            self.drop_lock(*lock, LockState::ReadLock);
1286        }
1287    }
1288
1289    fn drop_write_locks(&mut self, lock_keys: LockKeys) {
1290        for lock in lock_keys.iter() {
1291            // This is a bit hacky, but this works for locks in either the Locked or WriteLock
1292            // states.
1293            self.drop_lock(*lock, LockState::WriteLock);
1294        }
1295    }
1296
1297    // Downgrades locks from WriteLock to Locked.
1298    fn downgrade_locks(&mut self, lock_keys: &LockKeys) {
1299        for lock in lock_keys.iter() {
1300            // SAFETY: The lock in `LockManager::locks` is held.
1301            unsafe {
1302                self.keys.get_mut(lock).unwrap().downgrade_lock();
1303            }
1304        }
1305    }
1306}
1307
1308#[derive(Debug)]
1309struct LockEntry {
1310    // In the states that allow readers (ReadLock, Locked), this count can be non-zero
1311    // to indicate the number of active readers.
1312    read_count: u64,
1313
1314    // The state of the lock (see below).
1315    state: LockState,
1316
1317    // A doubly-linked list of wakers that should be woken when they have been granted the lock.
1318    // New wakers are usually chained on to tail, with the exception being the case where a lock in
1319    // state Locked is to be upgraded to WriteLock, but can't because there are readers.  It might
1320    // be possible to use intrusive-collections in the future.
1321    head: *const LockWaker,
1322    tail: *const LockWaker,
1323}
1324
1325unsafe impl Send for LockEntry {}
1326
1327// Represents a node in the waker list.  It is only safe to access the members wrapped by UnsafeCell
1328// when LockManager's `locks` member is locked.
1329struct LockWaker {
1330    // The next and previous pointers in the doubly-linked list.
1331    next: UnsafeCell<*const LockWaker>,
1332    prev: UnsafeCell<*const LockWaker>,
1333
1334    // Holds the lock key for this waker.  This is required so that we can find the associated
1335    // `LockEntry`.
1336    key: LockKey,
1337
1338    // The underlying waker that should be used to wake the task.
1339    waker: UnsafeCell<WakerState>,
1340
1341    // The target state for this waker.
1342    target_state: LockState,
1343
1344    // True if this is an upgrade.
1345    is_upgrade: bool,
1346
1347    // We need to be pinned because these form part of the linked list.
1348    _pin: PhantomPinned,
1349}
1350
1351enum WakerState {
1352    // This is the initial state before the waker has been first polled.
1353    Pending,
1354
1355    // Once polled, this contains the actual waker.
1356    Registered(Waker),
1357
1358    // The waker has been woken and has been granted the lock.
1359    Woken,
1360}
1361
1362impl WakerState {
1363    fn is_woken(&self) -> bool {
1364        matches!(self, WakerState::Woken)
1365    }
1366}
1367
1368unsafe impl Send for LockWaker {}
1369unsafe impl Sync for LockWaker {}
1370
1371impl LockWaker {
1372    // Waits for the waker to be woken.
1373    async fn wait(&self, manager: &LockManager) {
1374        // We must guard against the future being dropped.
1375        let waker_guard = scopeguard::guard((), |_| {
1376            let mut locks = manager.locks.lock();
1377            // SAFETY: We've acquired the lock.
1378            unsafe {
1379                if (*self.waker.get()).is_woken() {
1380                    // We were woken, but didn't actually run, so we must drop the lock.
1381                    if self.is_upgrade {
1382                        locks.keys.get_mut(&self.key).unwrap().downgrade_lock();
1383                    } else {
1384                        locks.drop_lock(self.key, self.target_state);
1385                    }
1386                } else {
1387                    // We haven't been woken but we've been dropped so we must remove ourself from
1388                    // the waker list.
1389                    locks.keys.get_mut(&self.key).unwrap().remove_waker(self);
1390                }
1391            }
1392        });
1393
1394        poll_fn(|cx| {
1395            let _locks = manager.locks.lock();
1396            // SAFETY: We've acquired the lock.
1397            unsafe {
1398                if (*self.waker.get()).is_woken() {
1399                    Poll::Ready(())
1400                } else {
1401                    *self.waker.get() = WakerState::Registered(cx.waker().clone());
1402                    Poll::Pending
1403                }
1404            }
1405        })
1406        .await;
1407
1408        ScopeGuard::into_inner(waker_guard);
1409    }
1410}
1411
1412#[derive(Copy, Clone, Debug, PartialEq)]
1413enum LockState {
1414    // In this state, there are only readers.
1415    ReadLock,
1416
1417    // This state is used for transactions to lock other writers (including other transactions), but
1418    // it still allows readers.
1419    Locked,
1420
1421    // A writer has exclusive access; all other readers and writers are blocked.
1422    WriteLock,
1423}
1424
1425impl LockManager {
1426    pub fn new() -> Self {
1427        LockManager { locks: Mutex::new(Locks { keys: HashMap::default() }) }
1428    }
1429
1430    /// Acquires the locks.  It is the caller's responsibility to ensure that drop_transaction is
1431    /// called when a transaction is dropped i.e. the filesystem's drop_transaction method should
1432    /// call LockManager's drop_transaction method.
1433    pub async fn txn_lock<'a>(&'a self, lock_keys: LockKeys) -> TransactionLocks<'a> {
1434        TransactionLocks(
1435            debug_assert_not_too_long!(self.lock(lock_keys, LockState::Locked)).right().unwrap(),
1436        )
1437    }
1438
1439    // `state` indicates the kind of lock required.  ReadLock means acquire a read lock.  Locked
1440    // means lock other writers, but still allow readers.  WriteLock means acquire a write lock.
1441    async fn lock<'a>(
1442        &'a self,
1443        mut lock_keys: LockKeys,
1444        target_state: LockState,
1445    ) -> Either<ReadGuard<'a>, WriteGuard<'a>> {
1446        let mut guard = match &target_state {
1447            LockState::ReadLock => Left(ReadGuard {
1448                manager: self.into(),
1449                lock_keys: LockKeys::with_capacity(lock_keys.len()),
1450            }),
1451            LockState::Locked | LockState::WriteLock => Right(WriteGuard {
1452                manager: self.into(),
1453                lock_keys: LockKeys::with_capacity(lock_keys.len()),
1454            }),
1455        };
1456        let guard_keys = match &mut guard {
1457            Left(g) => &mut g.lock_keys,
1458            Right(g) => &mut g.lock_keys,
1459        };
1460        lock_keys.sort_unstable();
1461        lock_keys.dedup();
1462        for lock in lock_keys.iter() {
1463            let lock_waker = None;
1464            pin_mut!(lock_waker);
1465            {
1466                let mut locks = self.locks.lock();
1467                match locks.keys.entry(*lock) {
1468                    Entry::Vacant(vacant) => {
1469                        vacant.insert(LockEntry {
1470                            read_count: if let LockState::ReadLock = target_state {
1471                                guard_keys.push(*lock);
1472                                1
1473                            } else {
1474                                guard_keys.push(*lock);
1475                                0
1476                            },
1477                            state: target_state,
1478                            head: std::ptr::null(),
1479                            tail: std::ptr::null(),
1480                        });
1481                    }
1482                    Entry::Occupied(mut occupied) => {
1483                        let entry = occupied.get_mut();
1484                        // SAFETY: We've acquired the lock.
1485                        if unsafe { entry.is_allowed(target_state, entry.head.is_null()) } {
1486                            if let LockState::ReadLock = target_state {
1487                                entry.read_count += 1;
1488                                guard_keys.push(*lock);
1489                            } else {
1490                                entry.state = target_state;
1491                                guard_keys.push(*lock);
1492                            }
1493                        } else {
1494                            // Initialise a waker and push it on the tail of the list.
1495                            // SAFETY: `lock_waker` isn't used prior to this point.
1496                            unsafe {
1497                                *lock_waker.as_mut().get_unchecked_mut() = Some(LockWaker {
1498                                    next: UnsafeCell::new(std::ptr::null()),
1499                                    prev: UnsafeCell::new(entry.tail),
1500                                    key: *lock,
1501                                    waker: UnsafeCell::new(WakerState::Pending),
1502                                    target_state: target_state,
1503                                    is_upgrade: false,
1504                                    _pin: PhantomPinned,
1505                                });
1506                            }
1507                            let waker = (*lock_waker).as_ref().unwrap();
1508                            if entry.tail.is_null() {
1509                                entry.head = waker;
1510                            } else {
1511                                // SAFETY: We've acquired the lock.
1512                                unsafe {
1513                                    *(*entry.tail).next.get() = waker;
1514                                }
1515                            }
1516                            entry.tail = waker;
1517                        }
1518                    }
1519                }
1520            }
1521            if let Some(waker) = &*lock_waker {
1522                waker.wait(self).await;
1523                guard_keys.push(*lock);
1524            }
1525        }
1526        guard
1527    }
1528
1529    /// This should be called by the filesystem's drop_transaction implementation.
1530    pub fn drop_transaction(&self, transaction: &mut Transaction<'_>) {
1531        let mut locks = self.locks.lock();
1532        locks.drop_write_locks(std::mem::take(&mut transaction.txn_locks));
1533    }
1534
1535    /// Prepares to commit by waiting for readers to finish.
1536    pub async fn commit_prepare(&self, transaction: &Transaction<'_>) {
1537        self.commit_prepare_keys(&transaction.txn_locks).await;
1538    }
1539
1540    async fn commit_prepare_keys(&self, lock_keys: &LockKeys) {
1541        for lock in lock_keys.iter() {
1542            let lock_waker = None;
1543            pin_mut!(lock_waker);
1544            {
1545                let mut locks = self.locks.lock();
1546                let entry = locks.keys.get_mut(lock).unwrap();
1547                // Callers may invoke `Transaction::commit_prepare` explicitly before committing
1548                // (e.g. to upgrade locks and verify invariants after disk I/O has finished).
1549                // Skipping keys already in `LockState::WriteLock` makes `commit_prepare`
1550                // idempotent when called again during `commit_transaction`.
1551                if entry.state == LockState::WriteLock {
1552                    continue;
1553                }
1554                assert_eq!(entry.state, LockState::Locked);
1555
1556                if entry.read_count == 0 {
1557                    entry.state = LockState::WriteLock;
1558                } else {
1559                    // Initialise a waker and push it on the head of the list.
1560                    // SAFETY: `lock_waker` isn't used prior to this point.
1561                    unsafe {
1562                        *lock_waker.as_mut().get_unchecked_mut() = Some(LockWaker {
1563                            next: UnsafeCell::new(entry.head),
1564                            prev: UnsafeCell::new(std::ptr::null()),
1565                            key: *lock,
1566                            waker: UnsafeCell::new(WakerState::Pending),
1567                            target_state: LockState::WriteLock,
1568                            is_upgrade: true,
1569                            _pin: PhantomPinned,
1570                        });
1571                    }
1572                    let waker = (*lock_waker).as_ref().unwrap();
1573                    if entry.head.is_null() {
1574                        entry.tail = (*lock_waker).as_ref().unwrap();
1575                    } else {
1576                        // SAFETY: We've acquired the lock.
1577                        unsafe {
1578                            *(*entry.head).prev.get() = waker;
1579                        }
1580                    }
1581                    entry.head = waker;
1582                }
1583            }
1584
1585            if let Some(waker) = &*lock_waker {
1586                waker.wait(self).await;
1587            }
1588        }
1589    }
1590
1591    /// Acquires a read lock for the given keys.  Read locks are only blocked whilst a transaction
1592    /// is being committed for the same locks.  They are only necessary where consistency is
1593    /// required between different mutations within a transaction.  For example, a write might
1594    /// change the size and extents for an object, in which case a read lock is required so that
1595    /// observed size and extents are seen together or not at all.
1596    pub async fn read_lock<'a>(&'a self, lock_keys: LockKeys) -> ReadGuard<'a> {
1597        debug_assert_not_too_long!(self.lock(lock_keys, LockState::ReadLock)).left().unwrap()
1598    }
1599
1600    /// Acquires a write lock for the given keys.  Write locks provide exclusive access to the
1601    /// requested lock keys.
1602    pub async fn write_lock<'a>(&'a self, lock_keys: LockKeys) -> WriteGuard<'a> {
1603        debug_assert_not_too_long!(self.lock(lock_keys, LockState::WriteLock)).right().unwrap()
1604    }
1605
1606    /// Downgrades locks from the WriteLock state to Locked state.  This will panic if the locks are
1607    /// not in the WriteLock state.
1608    pub fn downgrade_locks(&self, lock_keys: &LockKeys) {
1609        self.locks.lock().downgrade_locks(lock_keys);
1610    }
1611}
1612
1613// These unsafe functions require that `locks` in LockManager is locked.
1614impl LockEntry {
1615    unsafe fn wake(&mut self) {
1616        // If the lock's state is WriteLock, or there's nothing waiting, return early.
1617        if self.head.is_null() || self.state == LockState::WriteLock {
1618            return;
1619        }
1620
1621        let waker = unsafe { &*self.head };
1622
1623        if waker.is_upgrade {
1624            if self.read_count > 0 {
1625                return;
1626            }
1627        } else if !unsafe { self.is_allowed(waker.target_state, true) } {
1628            return;
1629        }
1630
1631        unsafe { self.pop_and_wake() };
1632
1633        // If the waker was a write lock, we can't wake any more up, but otherwise, we can keep
1634        // waking up readers.
1635        if waker.target_state == LockState::WriteLock {
1636            return;
1637        }
1638
1639        while !self.head.is_null() && unsafe { (*self.head).target_state } == LockState::ReadLock {
1640            unsafe { self.pop_and_wake() };
1641        }
1642    }
1643
1644    unsafe fn pop_and_wake(&mut self) {
1645        let waker = unsafe { &*self.head };
1646
1647        // Pop the waker.
1648        self.head = unsafe { *waker.next.get() };
1649        if self.head.is_null() {
1650            self.tail = std::ptr::null()
1651        } else {
1652            unsafe { *(*self.head).prev.get() = std::ptr::null() };
1653        }
1654
1655        // Adjust our state accordingly.
1656        if waker.target_state == LockState::ReadLock {
1657            self.read_count += 1;
1658        } else {
1659            self.state = waker.target_state;
1660        }
1661
1662        // Now wake the task.
1663        if let WakerState::Registered(waker) =
1664            std::mem::replace(unsafe { &mut *waker.waker.get() }, WakerState::Woken)
1665        {
1666            waker.wake();
1667        }
1668    }
1669
1670    fn can_remove(&self) -> bool {
1671        self.state == LockState::ReadLock && self.read_count == 0
1672    }
1673
1674    unsafe fn remove_waker(&mut self, waker: &LockWaker) {
1675        unsafe {
1676            let is_first = (*waker.prev.get()).is_null();
1677            if is_first {
1678                self.head = *waker.next.get();
1679            } else {
1680                *(**waker.prev.get()).next.get() = *waker.next.get();
1681            }
1682            if (*waker.next.get()).is_null() {
1683                self.tail = *waker.prev.get();
1684            } else {
1685                *(**waker.next.get()).prev.get() = *waker.prev.get();
1686            }
1687            if is_first {
1688                // We must call wake in case we erased a pending write lock and readers can now
1689                // proceed.
1690                self.wake();
1691            }
1692        }
1693    }
1694
1695    // Returns whether or not a lock with given `target_state` can proceed.  `is_head` should be
1696    // true if this is something at the head of the waker list (or the waker list is empty) and
1697    // false if there are other items on the waker list that are prior.
1698    unsafe fn is_allowed(&self, target_state: LockState, is_head: bool) -> bool {
1699        match self.state {
1700            LockState::ReadLock => {
1701                // Allow ReadLock and Locked so long as nothing else is waiting.
1702                (self.read_count == 0
1703                    || target_state == LockState::Locked
1704                    || target_state == LockState::ReadLock)
1705                    && is_head
1706            }
1707            LockState::Locked => {
1708                // Always allow reads unless there's an upgrade waiting.  We have to
1709                // always allow reads in this state because tasks that have locks in
1710                // the Locked state can later try and acquire ReadLock.
1711                target_state == LockState::ReadLock
1712                    && (is_head || unsafe { !(*self.head).is_upgrade })
1713            }
1714            LockState::WriteLock => false,
1715        }
1716    }
1717
1718    unsafe fn downgrade_lock(&mut self) {
1719        assert_eq!(std::mem::replace(&mut self.state, LockState::Locked), LockState::WriteLock);
1720        unsafe { self.wake() };
1721    }
1722}
1723
1724#[must_use]
1725pub struct ReadGuard<'a> {
1726    manager: LockManagerRef<'a>,
1727    lock_keys: LockKeys,
1728}
1729
1730impl ReadGuard<'_> {
1731    pub fn fs(&self) -> Option<&Arc<FxFilesystem>> {
1732        if let LockManagerRef::Owned(fs) = &self.manager { Some(fs) } else { None }
1733    }
1734
1735    pub fn into_owned(mut self, fs: Arc<FxFilesystem>) -> ReadGuard<'static> {
1736        ReadGuard {
1737            manager: LockManagerRef::Owned(fs),
1738            lock_keys: std::mem::replace(&mut self.lock_keys, LockKeys::None),
1739        }
1740    }
1741}
1742
1743impl Drop for ReadGuard<'_> {
1744    fn drop(&mut self) {
1745        let mut locks = self.manager.locks.lock();
1746        locks.drop_read_locks(std::mem::take(&mut self.lock_keys));
1747    }
1748}
1749
1750impl fmt::Debug for ReadGuard<'_> {
1751    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1752        f.debug_struct("ReadGuard")
1753            .field("manager", &(&self.manager as *const _))
1754            .field("lock_keys", &self.lock_keys)
1755            .finish()
1756    }
1757}
1758
1759#[must_use]
1760pub struct WriteGuard<'a> {
1761    manager: LockManagerRef<'a>,
1762    lock_keys: LockKeys,
1763}
1764
1765impl WriteGuard<'_> {
1766    pub fn into_owned(mut self, fs: Arc<FxFilesystem>) -> WriteGuard<'static> {
1767        WriteGuard {
1768            manager: LockManagerRef::Owned(fs),
1769            lock_keys: std::mem::replace(&mut self.lock_keys, LockKeys::None),
1770        }
1771    }
1772}
1773
1774impl Drop for WriteGuard<'_> {
1775    fn drop(&mut self) {
1776        let mut locks = self.manager.locks.lock();
1777        locks.drop_write_locks(std::mem::take(&mut self.lock_keys));
1778    }
1779}
1780
1781impl fmt::Debug for WriteGuard<'_> {
1782    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1783        f.debug_struct("WriteGuard")
1784            .field("manager", &(&self.manager as *const _))
1785            .field("lock_keys", &self.lock_keys)
1786            .finish()
1787    }
1788}
1789
1790enum LockManagerRef<'a> {
1791    Borrowed(&'a LockManager),
1792    Owned(Arc<FxFilesystem>),
1793}
1794
1795impl Deref for LockManagerRef<'_> {
1796    type Target = LockManager;
1797
1798    fn deref(&self) -> &Self::Target {
1799        match self {
1800            LockManagerRef::Borrowed(m) => m,
1801            LockManagerRef::Owned(f) => f.lock_manager(),
1802        }
1803    }
1804}
1805
1806impl<'a> From<&'a LockManager> for LockManagerRef<'a> {
1807    fn from(value: &'a LockManager) -> Self {
1808        LockManagerRef::Borrowed(value)
1809    }
1810}
1811
1812#[cfg(test)]
1813mod tests {
1814    use super::{
1815        AssocObj, AttributeId, LockKey, LockKeys, LockManager, LockState, Mutation,
1816        ObjectMutationIterator, Options, ReservationOptions, TxnMutation,
1817    };
1818    use crate::filesystem::FxFilesystem;
1819    use crate::object_store::{BytesAndNodes, ObjectKey};
1820    use fuchsia_async as fasync;
1821    use fuchsia_sync::Mutex;
1822    use futures::channel::oneshot::channel;
1823    use futures::future::FutureExt;
1824    use futures::stream::FuturesUnordered;
1825    use futures::{StreamExt, join, pin_mut};
1826    use std::collections::BTreeSet;
1827    use std::task::Poll;
1828    use std::time::Duration;
1829    use storage_device::DeviceHolder;
1830    use storage_device::fake_device::FakeDevice;
1831
1832    #[fuchsia::test]
1833    async fn test_simple() {
1834        let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1835        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1836        let mut t = fs
1837            .root_store()
1838            .new_transaction(lock_keys![], Options::default())
1839            .await
1840            .expect("new_transaction failed");
1841        t.add(1, Mutation::BeginFlush);
1842        assert!(!t.is_empty());
1843    }
1844
1845    #[fuchsia::test]
1846    async fn test_locks() {
1847        let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1848        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1849        let (send1, recv1) = channel();
1850        let (send2, recv2) = channel();
1851        let (send3, recv3) = channel();
1852        let done = Mutex::new(false);
1853        let mut futures = FuturesUnordered::new();
1854        futures.push(
1855            async {
1856                let _t = fs
1857                    .root_store()
1858                    .new_transaction(
1859                        lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1860                        Options::default(),
1861                    )
1862                    .await
1863                    .expect("new_transaction failed");
1864                send1.send(()).unwrap(); // Tell the next future to continue.
1865                send3.send(()).unwrap(); // Tell the last future to continue.
1866                recv2.await.unwrap();
1867                // This is a halting problem so all we can do is sleep.
1868                fasync::Timer::new(Duration::from_millis(100)).await;
1869                assert!(!*done.lock());
1870            }
1871            .boxed(),
1872        );
1873        futures.push(
1874            async {
1875                recv1.await.unwrap();
1876                // This should not block since it is a different key.
1877                let _t = fs
1878                    .root_store()
1879                    .new_transaction(
1880                        lock_keys![LockKey::object_attribute(2, 2, AttributeId::TEST_ID)],
1881                        Options::default(),
1882                    )
1883                    .await
1884                    .expect("new_transaction failed");
1885                // Tell the first future to continue.
1886                send2.send(()).unwrap();
1887            }
1888            .boxed(),
1889        );
1890        futures.push(
1891            async {
1892                // This should block until the first future has completed.
1893                recv3.await.unwrap();
1894                let _t = fs
1895                    .root_store()
1896                    .new_transaction(
1897                        lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1898                        Options::default(),
1899                    )
1900                    .await;
1901                *done.lock() = true;
1902            }
1903            .boxed(),
1904        );
1905        while let Some(()) = futures.next().await {}
1906    }
1907
1908    #[fuchsia::test]
1909    async fn test_read_lock_after_write_lock() {
1910        let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1911        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1912        let (send1, recv1) = channel();
1913        let (send2, recv2) = channel();
1914        let done = Mutex::new(false);
1915        join!(
1916            async {
1917                let t = fs
1918                    .root_store()
1919                    .new_transaction(
1920                        lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1921                        Options::default(),
1922                    )
1923                    .await
1924                    .expect("new_transaction failed");
1925                send1.send(()).unwrap(); // Tell the next future to continue.
1926                recv2.await.unwrap();
1927                t.commit().await.expect("commit failed");
1928                *done.lock() = true;
1929            },
1930            async {
1931                recv1.await.unwrap();
1932                // Reads should not be blocked until the transaction is committed.
1933                let _guard = fs
1934                    .lock_manager()
1935                    .read_lock(lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)])
1936                    .await;
1937                // Tell the first future to continue.
1938                send2.send(()).unwrap();
1939                // It shouldn't proceed until we release our read lock, but it's a halting
1940                // problem, so sleep.
1941                fasync::Timer::new(Duration::from_millis(100)).await;
1942                assert!(!*done.lock());
1943            },
1944        );
1945    }
1946
1947    #[fuchsia::test]
1948    async fn test_write_lock_after_read_lock() {
1949        let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1950        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1951        let (send1, recv1) = channel();
1952        let (send2, recv2) = channel();
1953        let done = Mutex::new(false);
1954        join!(
1955            async {
1956                // Reads should not be blocked until the transaction is committed.
1957                let _guard = fs
1958                    .lock_manager()
1959                    .read_lock(lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)])
1960                    .await;
1961                // Tell the next future to continue and then wait.
1962                send1.send(()).unwrap();
1963                recv2.await.unwrap();
1964                // It shouldn't proceed until we release our read lock, but it's a halting
1965                // problem, so sleep.
1966                fasync::Timer::new(Duration::from_millis(100)).await;
1967                assert!(!*done.lock());
1968            },
1969            async {
1970                recv1.await.unwrap();
1971                let t = fs
1972                    .root_store()
1973                    .new_transaction(
1974                        lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1975                        Options::default(),
1976                    )
1977                    .await
1978                    .expect("new_transaction failed");
1979                send2.send(()).unwrap(); // Tell the first future to continue;
1980                t.commit().await.expect("commit failed");
1981                *done.lock() = true;
1982            },
1983        );
1984    }
1985
1986    #[fuchsia::test]
1987    async fn test_drop_uncommitted_transaction() {
1988        let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1989        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1990        let key = lock_keys![LockKey::object(1, 1)];
1991
1992        // Dropping while there's a reader.
1993        {
1994            let _write_lock = fs
1995                .root_store()
1996                .new_transaction(key.clone(), Options::default())
1997                .await
1998                .expect("new_transaction failed");
1999            let _read_lock = fs.lock_manager().read_lock(key.clone()).await;
2000        }
2001        // Dropping while there's no reader.
2002        {
2003            let _write_lock = fs
2004                .root_store()
2005                .new_transaction(key.clone(), Options::default())
2006                .await
2007                .expect("new_transaction failed");
2008        }
2009        // Make sure we can take the lock again (i.e. it was actually released).
2010        fs.root_store()
2011            .new_transaction(key.clone(), Options::default())
2012            .await
2013            .expect("new_transaction failed");
2014    }
2015
2016    #[fuchsia::test]
2017    async fn test_drop_waiting_write_lock() {
2018        let manager = LockManager::new();
2019        let keys = lock_keys![LockKey::object(1, 1)];
2020        {
2021            let _guard = manager.lock(keys.clone(), LockState::ReadLock).await;
2022            if let Poll::Ready(_) =
2023                futures::poll!(manager.lock(keys.clone(), LockState::WriteLock).boxed())
2024            {
2025                assert!(false);
2026            }
2027        }
2028        let _ = manager.lock(keys, LockState::WriteLock).await;
2029    }
2030
2031    #[fuchsia::test]
2032    async fn test_write_lock_blocks_everything() {
2033        let manager = LockManager::new();
2034        let keys = lock_keys![LockKey::object(1, 1)];
2035        {
2036            let _guard = manager.lock(keys.clone(), LockState::WriteLock).await;
2037            if let Poll::Ready(_) =
2038                futures::poll!(manager.lock(keys.clone(), LockState::WriteLock).boxed())
2039            {
2040                assert!(false);
2041            }
2042            if let Poll::Ready(_) =
2043                futures::poll!(manager.lock(keys.clone(), LockState::ReadLock).boxed())
2044            {
2045                assert!(false);
2046            }
2047        }
2048        {
2049            let _guard = manager.lock(keys.clone(), LockState::WriteLock).await;
2050        }
2051        {
2052            let _guard = manager.lock(keys, LockState::ReadLock).await;
2053        }
2054    }
2055
2056    #[fuchsia::test]
2057    async fn test_downgrade_locks() {
2058        let manager = LockManager::new();
2059        let keys = lock_keys![LockKey::object(1, 1)];
2060        let _guard = manager.txn_lock(keys.clone()).await;
2061        manager.commit_prepare_keys(&keys).await;
2062
2063        // Use FuturesUnordered so that we can check that the waker is woken.
2064        let mut read_lock: FuturesUnordered<_> =
2065            std::iter::once(manager.read_lock(keys.clone())).collect();
2066
2067        // Trying to acquire a read lock now should be blocked.
2068        assert!(futures::poll!(read_lock.next()).is_pending());
2069
2070        manager.downgrade_locks(&keys);
2071
2072        // After downgrading, it should be possible to take a read lock.
2073        assert!(futures::poll!(read_lock.next()).is_ready());
2074    }
2075
2076    #[fuchsia::test]
2077    async fn test_dropped_write_lock_wakes() {
2078        let manager = LockManager::new();
2079        let keys = lock_keys![LockKey::object(1, 1)];
2080        let _guard = manager.lock(keys.clone(), LockState::ReadLock).await;
2081        let mut read_lock = FuturesUnordered::new();
2082        read_lock.push(manager.lock(keys.clone(), LockState::ReadLock));
2083
2084        {
2085            let write_lock = manager.lock(keys, LockState::WriteLock);
2086            pin_mut!(write_lock);
2087
2088            // The write lock should be blocked because of the read lock.
2089            assert!(futures::poll!(write_lock).is_pending());
2090
2091            // Another read lock should be blocked because of the write lock.
2092            assert!(futures::poll!(read_lock.next()).is_pending());
2093        }
2094
2095        // Dropping the write lock should allow the read lock to proceed.
2096        assert!(futures::poll!(read_lock.next()).is_ready());
2097    }
2098
2099    #[fuchsia::test]
2100    async fn test_drop_upgrade() {
2101        let manager = LockManager::new();
2102        let keys = lock_keys![LockKey::object(1, 1)];
2103        let _guard = manager.lock(keys.clone(), LockState::Locked).await;
2104
2105        {
2106            let commit_prepare = manager.commit_prepare_keys(&keys);
2107            pin_mut!(commit_prepare);
2108            let _read_guard = manager.lock(keys.clone(), LockState::ReadLock).await;
2109            assert!(futures::poll!(commit_prepare).is_pending());
2110
2111            // Now we test dropping read_guard which should wake commit_prepare and
2112            // then dropping commit_prepare.
2113        }
2114
2115        // We should be able to still commit_prepare.
2116        manager.commit_prepare_keys(&keys).await;
2117    }
2118
2119    #[fuchsia::test]
2120    async fn test_woken_upgrade_blocks_reads() {
2121        let manager = LockManager::new();
2122        let keys = lock_keys![LockKey::object(1, 1)];
2123        // Start with a transaction lock.
2124        let guard = manager.lock(keys.clone(), LockState::Locked).await;
2125
2126        // Take a read lock.
2127        let read1 = manager.lock(keys.clone(), LockState::ReadLock).await;
2128
2129        // Try and upgrade the transaction lock, which should not be possible because of the read.
2130        let commit_prepare = manager.commit_prepare_keys(&keys);
2131        pin_mut!(commit_prepare);
2132        assert!(futures::poll!(commit_prepare.as_mut()).is_pending());
2133
2134        // Taking another read should also be blocked.
2135        let read2 = manager.lock(keys.clone(), LockState::ReadLock);
2136        pin_mut!(read2);
2137        assert!(futures::poll!(read2.as_mut()).is_pending());
2138
2139        // Drop the first read and the upgrade should complete.
2140        std::mem::drop(read1);
2141        assert!(futures::poll!(commit_prepare).is_ready());
2142
2143        // But the second read should still be blocked.
2144        assert!(futures::poll!(read2.as_mut()).is_pending());
2145
2146        // If we drop the write lock now, the read should be unblocked.
2147        std::mem::drop(guard);
2148        assert!(futures::poll!(read2).is_ready());
2149    }
2150
2151    static LOCK_KEY_1: LockKey = LockKey::flush(1);
2152    static LOCK_KEY_2: LockKey = LockKey::flush(2);
2153    static LOCK_KEY_3: LockKey = LockKey::flush(3);
2154
2155    // The keys, storage method, and capacity must all match.
2156    fn assert_lock_keys_equal(value: &LockKeys, expected: &LockKeys) {
2157        match (value, expected) {
2158            (LockKeys::None, LockKeys::None) => {}
2159            (LockKeys::Inline(key1), LockKeys::Inline(key2)) => {
2160                if key1 != key2 {
2161                    panic!("{key1:?} != {key2:?}");
2162                }
2163            }
2164            (LockKeys::Vec(vec1), LockKeys::Vec(vec2)) => {
2165                if vec1 != vec2 {
2166                    panic!("{vec1:?} != {vec2:?}");
2167                }
2168                if vec1.capacity() != vec2.capacity() {
2169                    panic!(
2170                        "LockKeys have different capacity: {} != {}",
2171                        vec1.capacity(),
2172                        vec2.capacity()
2173                    );
2174                }
2175            }
2176            (_, _) => panic!("{value:?} != {expected:?}"),
2177        }
2178    }
2179
2180    // Only the keys must match. Storage method and capacity don't matter.
2181    fn assert_lock_keys_equivalent(value: &LockKeys, expected: &LockKeys) {
2182        let value: Vec<_> = value.iter().collect();
2183        let expected: Vec<_> = expected.iter().collect();
2184        assert_eq!(value, expected);
2185    }
2186
2187    #[test]
2188    fn test_lock_keys_macro() {
2189        assert_lock_keys_equal(&lock_keys![], &LockKeys::None);
2190        assert_lock_keys_equal(&lock_keys![LOCK_KEY_1], &LockKeys::Inline(LOCK_KEY_1));
2191        assert_lock_keys_equal(
2192            &lock_keys![LOCK_KEY_1, LOCK_KEY_2],
2193            &LockKeys::Vec(vec![LOCK_KEY_1, LOCK_KEY_2]),
2194        );
2195    }
2196
2197    #[test]
2198    fn test_lock_keys_with_capacity() {
2199        assert_lock_keys_equal(&LockKeys::with_capacity(0), &LockKeys::None);
2200        assert_lock_keys_equal(&LockKeys::with_capacity(1), &LockKeys::None);
2201        assert_lock_keys_equal(&LockKeys::with_capacity(2), &LockKeys::Vec(Vec::with_capacity(2)));
2202    }
2203
2204    #[test]
2205    fn test_lock_keys_len() {
2206        assert_eq!(lock_keys![].len(), 0);
2207        assert_eq!(lock_keys![LOCK_KEY_1].len(), 1);
2208        assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].len(), 2);
2209    }
2210
2211    #[test]
2212    fn test_lock_keys_contains() {
2213        assert_eq!(lock_keys![].contains(&LOCK_KEY_1), false);
2214        assert_eq!(lock_keys![LOCK_KEY_1].contains(&LOCK_KEY_1), true);
2215        assert_eq!(lock_keys![LOCK_KEY_1].contains(&LOCK_KEY_2), false);
2216        assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].contains(&LOCK_KEY_1), true);
2217        assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].contains(&LOCK_KEY_2), true);
2218        assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].contains(&LOCK_KEY_3), false);
2219    }
2220
2221    #[test]
2222    fn test_lock_keys_push() {
2223        let mut keys = lock_keys![];
2224        keys.push(LOCK_KEY_1);
2225        assert_lock_keys_equal(&keys, &LockKeys::Inline(LOCK_KEY_1));
2226        keys.push(LOCK_KEY_2);
2227        assert_lock_keys_equal(&keys, &LockKeys::Vec(vec![LOCK_KEY_1, LOCK_KEY_2]));
2228        keys.push(LOCK_KEY_3);
2229        assert_lock_keys_equivalent(
2230            &keys,
2231            &LockKeys::Vec(vec![LOCK_KEY_1, LOCK_KEY_2, LOCK_KEY_3]),
2232        );
2233    }
2234
2235    #[test]
2236    fn test_lock_keys_sort_unstable() {
2237        let mut keys = lock_keys![];
2238        keys.sort_unstable();
2239        assert_lock_keys_equal(&keys, &lock_keys![]);
2240
2241        let mut keys = lock_keys![LOCK_KEY_1];
2242        keys.sort_unstable();
2243        assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1]);
2244
2245        let mut keys = lock_keys![LOCK_KEY_2, LOCK_KEY_1];
2246        keys.sort_unstable();
2247        assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1, LOCK_KEY_2]);
2248    }
2249
2250    #[test]
2251    fn test_lock_keys_dedup() {
2252        let mut keys = lock_keys![];
2253        keys.dedup();
2254        assert_lock_keys_equal(&keys, &lock_keys![]);
2255
2256        let mut keys = lock_keys![LOCK_KEY_1];
2257        keys.dedup();
2258        assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1]);
2259
2260        let mut keys = lock_keys![LOCK_KEY_1, LOCK_KEY_1];
2261        keys.dedup();
2262        assert_lock_keys_equivalent(&keys, &lock_keys![LOCK_KEY_1]);
2263    }
2264
2265    #[test]
2266    fn test_lock_keys_truncate() {
2267        let mut keys = lock_keys![];
2268        keys.truncate(5);
2269        assert_lock_keys_equal(&keys, &lock_keys![]);
2270        keys.truncate(0);
2271        assert_lock_keys_equal(&keys, &lock_keys![]);
2272
2273        let mut keys = lock_keys![LOCK_KEY_1];
2274        keys.truncate(5);
2275        assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1]);
2276        keys.truncate(0);
2277        assert_lock_keys_equal(&keys, &lock_keys![]);
2278
2279        let mut keys = lock_keys![LOCK_KEY_1, LOCK_KEY_2];
2280        keys.truncate(5);
2281        assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1, LOCK_KEY_2]);
2282        keys.truncate(1);
2283        // Although there's only 1 key after truncate the key is not stored inline.
2284        assert_lock_keys_equivalent(&keys, &lock_keys![LOCK_KEY_1]);
2285    }
2286
2287    #[test]
2288    fn test_lock_keys_iter() {
2289        assert_eq!(lock_keys![].iter().collect::<Vec<_>>(), Vec::<&LockKey>::new());
2290
2291        assert_eq!(lock_keys![LOCK_KEY_1].iter().collect::<Vec<_>>(), vec![&LOCK_KEY_1]);
2292
2293        assert_eq!(
2294            lock_keys![LOCK_KEY_1, LOCK_KEY_2].iter().collect::<Vec<_>>(),
2295            vec![&LOCK_KEY_1, &LOCK_KEY_2]
2296        );
2297    }
2298
2299    #[test]
2300    fn test_object_mutation_iterator() {
2301        let mut mutations = BTreeSet::new();
2302        mutations.insert(TxnMutation {
2303            object_id: 1,
2304            mutation: Mutation::BeginFlush,
2305            associated_object: AssocObj::None,
2306        });
2307        mutations.insert(TxnMutation {
2308            object_id: 1,
2309            mutation: Mutation::EndFlush,
2310            associated_object: AssocObj::None,
2311        });
2312        mutations.insert(TxnMutation {
2313            object_id: 2,
2314            mutation: Mutation::DeleteVolume,
2315            associated_object: AssocObj::None,
2316        });
2317
2318        let mut iter = mutations.iter().peekable();
2319
2320        {
2321            let mut obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 1");
2322            assert_eq!(obj_iter.object_id(), 1);
2323            assert_eq!(obj_iter.next(), Some(&Mutation::BeginFlush));
2324            assert_eq!(obj_iter.next(), Some(&Mutation::EndFlush));
2325            assert_eq!(obj_iter.next(), None);
2326        }
2327
2328        {
2329            let mut obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 2");
2330            assert_eq!(obj_iter.object_id(), 2);
2331            assert_eq!(obj_iter.next(), Some(&Mutation::DeleteVolume));
2332            assert_eq!(obj_iter.next(), None);
2333        }
2334
2335        // No more objects
2336        assert!(ObjectMutationIterator::new(&mut iter).is_none());
2337    }
2338
2339    #[test]
2340    fn test_object_mutation_iterator_drop_drains_remaining() {
2341        let mut mutations = BTreeSet::new();
2342        // Object 1 with 3 mutations
2343        mutations.insert(TxnMutation {
2344            object_id: 1,
2345            mutation: Mutation::BeginFlush,
2346            associated_object: AssocObj::None,
2347        });
2348        mutations.insert(TxnMutation {
2349            object_id: 1,
2350            mutation: Mutation::EndFlush,
2351            associated_object: AssocObj::None,
2352        });
2353        mutations.insert(TxnMutation {
2354            object_id: 1,
2355            mutation: Mutation::DeleteVolume,
2356            associated_object: AssocObj::None,
2357        });
2358        // Object 2 with 2 mutations
2359        mutations.insert(TxnMutation {
2360            object_id: 2,
2361            mutation: Mutation::BeginFlush,
2362            associated_object: AssocObj::None,
2363        });
2364        mutations.insert(TxnMutation {
2365            object_id: 2,
2366            mutation: Mutation::EndFlush,
2367            associated_object: AssocObj::None,
2368        });
2369
2370        let mut iter = mutations.iter().peekable();
2371
2372        {
2373            let mut obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 1");
2374            assert_eq!(obj_iter.object_id(), 1);
2375            assert_eq!(obj_iter.next(), Some(&Mutation::BeginFlush));
2376            // Drop without reading EndFlush or DeleteVolume
2377        }
2378
2379        {
2380            let obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 2");
2381            assert_eq!(obj_iter.object_id(), 2);
2382            // Drop immediately without reading any mutations for object 2
2383        }
2384
2385        // Dropping obj_iter for object 2 should have drained all object 2 mutations as well.
2386        assert!(ObjectMutationIterator::new(&mut iter).is_none());
2387    }
2388
2389    #[fuchsia::test]
2390    async fn test_merge_bytes_and_nodes() {
2391        let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
2392        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2393        let root_store = fs.root_store();
2394        let store_id = root_store.store_object_id();
2395        let mut t = root_store
2396            .new_transaction(lock_keys![], Options::default())
2397            .await
2398            .expect("new_transaction failed");
2399
2400        let key = ObjectKey::project_usage(
2401            root_store.root_directory_object_id(),
2402            crate::object_store::ProjectId::new(1).unwrap(),
2403        );
2404
2405        // Initial delta.
2406        t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: 100, nodes: 2 });
2407        let mutation = t.get_object_mutation(store_id, key.clone()).expect("mutation expected");
2408        assert_eq!(mutation.op, super::Operation::Merge);
2409        assert_eq!(
2410            mutation.item.value,
2411            crate::object_store::object_record::ObjectValue::BytesAndNodes { bytes: 100, nodes: 2 }
2412        );
2413
2414        // Second delta for the same key folds into the first rather than replacing it.
2415        t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: 50, nodes: 1 });
2416        let mutation = t.get_object_mutation(store_id, key.clone()).expect("mutation expected");
2417        assert_eq!(mutation.op, super::Operation::Merge);
2418        assert_eq!(
2419            mutation.item.value,
2420            crate::object_store::object_record::ObjectValue::BytesAndNodes { bytes: 150, nodes: 3 }
2421        );
2422
2423        // Negative delta brings it back down.
2424        t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: -50, nodes: -1 });
2425        let mutation = t.get_object_mutation(store_id, key.clone()).expect("mutation expected");
2426        assert_eq!(mutation.op, super::Operation::Merge);
2427        assert_eq!(
2428            mutation.item.value,
2429            crate::object_store::object_record::ObjectValue::BytesAndNodes { bytes: 100, nodes: 2 }
2430        );
2431
2432        // Folding to zero removes the mutation completely.
2433        t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: -100, nodes: -2 });
2434        assert!(t.get_object_mutation(store_id, key.clone()).is_none());
2435
2436        // Calling with (0, 0) when no mutation exists is a no-op.
2437        t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: 0, nodes: 0 });
2438        assert!(t.get_object_mutation(store_id, key.clone()).is_none());
2439    }
2440
2441    #[fuchsia::test]
2442    async fn test_commit_and_continue_checks_journal_space() {
2443        use crate::filesystem::FxFilesystemBuilder;
2444        use crate::hooks::Hooks;
2445        use crate::object_store::journal::JournalOptions;
2446        use crate::object_store::volume::root_volume;
2447        use crate::object_store::{NewChildStoreOptions, ObjectKey, ObjectValue};
2448        use storage_device::DeviceHolder;
2449        use storage_device::fake_device::FakeDevice;
2450
2451        let reclaim_size = 65536;
2452        let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
2453        let (mut hooks, fs_hooks) = Hooks::new();
2454        let fs = FxFilesystemBuilder::new()
2455            .hooks(fs_hooks)
2456            .journal_options(JournalOptions { reclaim_size, ..Default::default() })
2457            .format(true)
2458            .open(device)
2459            .await
2460            .expect("open failed");
2461
2462        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
2463        let store = root_volume
2464            .new_volume("test", NewChildStoreOptions::default())
2465            .await
2466            .expect("new_volume failed");
2467
2468        fs.journal().force_compact().await.expect("force_compact failed");
2469        fs.journal().pause_compactions().await;
2470
2471        let journal_clone = fs.journal().clone();
2472        hooks.set_waiting_for_journal_space(move || {
2473            journal_clone.resume_compactions();
2474        });
2475
2476        // Reserve all remaining free space so the journal has no free space to grow into beyond
2477        // the metadata reservation.
2478        let _reservation = fs.allocator().reserve_with(None, |limit| limit);
2479
2480        let mut transaction = store
2481            .new_transaction(
2482                lock_keys![LockKey::object(store.store_object_id(), 1000)],
2483                Options { reservation: ReservationOptions::BorrowedMetadata, ..Default::default() },
2484            )
2485            .await
2486            .expect("new_transaction failed");
2487
2488        // Repeatedly write mutations and call commit_and_continue so the journal must extend past
2489        // its reserved space unless commit_and_continue checks for journal space and waits for
2490        // compaction.
2491        for _ in 0..64 {
2492            for i in 0..16 {
2493                transaction.add(
2494                    store.store_object_id(),
2495                    Mutation::replace_or_insert_object(
2496                        ObjectKey::extended_attribute(1000, format!("attr_{i}").into_bytes()),
2497                        ObjectValue::inline_extended_attribute(vec![0u8; 1024]),
2498                    ),
2499                );
2500            }
2501            transaction.commit_and_continue().await.expect("commit_and_continue failed");
2502        }
2503        transaction.commit().await.expect("commit failed");
2504
2505        fs.close().await.expect("close failed");
2506    }
2507
2508    #[fuchsia::test]
2509    async fn test_commit_and_continue_resets_state() {
2510        use crate::object_store::volume::root_volume;
2511        use crate::object_store::{ExtentValue, NewChildStoreOptions, ObjectKey, ObjectValue};
2512        use storage_device::DeviceHolder;
2513        use storage_device::fake_device::FakeDevice;
2514
2515        let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
2516        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2517        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
2518        let store = root_volume
2519            .new_volume("test", NewChildStoreOptions::default())
2520            .await
2521            .expect("new_volume failed");
2522
2523        let mut transaction = store
2524            .new_transaction(
2525                lock_keys![LockKey::object(store.store_object_id(), 1000)],
2526                Options::default(),
2527            )
2528            .await
2529            .expect("new_transaction failed");
2530
2531        assert!(!transaction.includes_write());
2532        transaction.add(
2533            store.store_object_id(),
2534            Mutation::merge_object(
2535                ObjectKey::extent(1000, AttributeId::DATA, 0..4096),
2536                ObjectValue::Extent(ExtentValue::deleted_extent()),
2537            ),
2538        );
2539        transaction.add_checksum(0..4096, vec![0], true);
2540        assert!(transaction.includes_write());
2541        assert!(!transaction.checksums().is_empty());
2542
2543        transaction.commit_and_continue().await.expect("commit_and_continue failed");
2544
2545        assert!(!transaction.includes_write());
2546        assert!(transaction.checksums().is_empty());
2547
2548        transaction.commit().await.expect("commit failed");
2549        fs.close().await.expect("close failed");
2550    }
2551}