Skip to main content

fxfs/object_store/
allocator.rs

1// Copyright 2021 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5//! # The Allocator
6//!
7//! The allocator in Fxfs is filesystem-wide entity responsible for managing the allocation of
8//! regions of the device to "owners" (which are `ObjectStore`).
9//!
10//! Allocations are tracked in an LSMTree with coalescing used to merge neighboring allocations
11//! with the same properties (owner and reference count). As of writing, reference counting is not
12//! used. (Reference counts are intended for future use if/when snapshot support is implemented.)
13//!
14//! There are a number of important implementation features of the allocator that warrant further
15//! documentation.
16//!
17//! ## Byte limits
18//!
19//! Fxfs is a multi-volume filesystem. Fuchsia with fxblob currently uses two primary volumes -
20//! an unencrypted volume for blob storage and an encrypted volume for data storage.
21//! Byte limits ensure that no one volume can consume all available storage. This is important
22//! as software updates must always be possible (blobs) and conversely configuration data should
23//! always be writable (data).
24//!
25//! ## Reservation tracking
26//!
27//! Fxfs on Fuchsia leverages write-back caching which allows us to buffer writes in RAM until
28//! memory pressure, an explicit flush or a file close requires us to persist data to storage.
29//!
30//! To ensure that we do not over-commit in such cases (e.g. by writing many files to RAM only to
31//! find out tha there is insufficient disk to store them all), the allocator includes a simple
32//! reservation tracking mechanism.
33//!
34//! Reservation tracking is implemented hierarchically such that a reservation can portion out
35//! sub-reservations, each of which may be "reserved" or "forgotten" when it comes time to actually
36//! allocate space.
37//!
38//! ## Fixed locations
39//!
40//! The only fixed location files in Fxfs are the first 512kiB extents of the two superblocks that
41//! exist as the first things on the disk (i.e. The first 1MiB). The allocator manages these
42//! locations, but does so using a 'mark_allocated' method distinct from all other allocations in
43//! which the location is left up to the allocator.
44//!
45//! ## Deallocated unusable regions
46//!
47//! It is not legal to reuse a deallocated disk region until after a flush. Transactions
48//! are not guaranteed to be persisted until after a successful flush so any reuse of
49//! a region before then followed by power loss may lead to corruption.
50//!
51//! e.g. Consider if we deallocate a file, reuse its extent for another file, then crash after
52//! writing to the new file but not yet flushing the journal. At next mount we will have lost the
53//! transaction despite overwriting the original file's data, leading to inconsistency errors (data
54//! corruption).
55//!
56//! These regions are currently tracked in RAM in the allocator.
57//!
58//! ## Allocated but uncommitted regions
59//!
60//! These are allocated regions that are not (yet) persisted to disk. They are regular
61//! file allocations, but are not stored on persistent storage until their transaction is committed
62//! and safely flushed to disk.
63//!
64//! ## TRIMed unusable regions
65//!
66//! We periodically TRIM unallocated regions to give the SSD controller insight into which
67//! parts of the device contain data. We must avoid using temporary TRIM allocations that are held
68//! while we perform these operations.
69//!
70//! ## Volume deletion
71//!
72//! We make use of an optimisation in the case where an entire volume is deleted. In such cases,
73//! rather than individually deallocate all disk regions associated with that volume, we make
74//! note of the deletion and perform special merge operation on the next LSMTree compaction that
75//! filters out allocations for the deleted volume.
76//!
77//! This is designed to make dropping of volumes significantly cheaper, but it does add some
78//! additional complexity if implementing an allocator that implements data structures to track
79//! free space (rather than just allocated space).
80//!
81//! ## Image generation
82//!
83//! The Fuchsia build process requires building an initial filesystem image. In the case of
84//! fxblob-based boards, this is an Fxfs filesystem containing a volume with the base set of
85//! blobs required to bootstrap the system. When we build such an image, we want it to be as compact
86//! as possible as we're potentially packaging it up for distribution. To that end, our allocation
87//! strategy (or at least the strategy used for image generation) should prefer to allocate from the
88//! start of the device wherever possible.
89pub mod merge;
90pub mod strategy;
91
92use crate::drop_event::DropEvent;
93use crate::errors::FxfsError;
94use crate::filesystem::{ApplyContext, ApplyMode, FxFilesystem, JournalingObject, SyncOptions};
95use crate::log::*;
96use crate::lsm_tree::skip_list_layer::SkipListLayer;
97use crate::lsm_tree::types::{
98    FuzzyHash, Item, ItemRef, Layer, LayerIterator, LayerKey, MergeType, OrdLowerBound,
99    OrdUpperBound, SortByU64, Value,
100};
101use crate::lsm_tree::{LSMTree, Query, compact_with_iterator, layer_from_handle, open_layers};
102use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, ReadObjectHandle};
103use crate::object_store::object_manager::ReservationUpdate;
104use crate::object_store::transaction::{
105    AllocatorMutation, AssocObj, LockKey, Mutation, Options, ReservationOptions, Transaction,
106    WriteGuard, lock_keys,
107};
108use crate::object_store::{
109    DataObjectHandle, DirectWriter, Extent, HandleOptions, ObjectStore, ReservedId, tree,
110};
111use crate::range::RangeExt;
112use crate::round::round_div;
113use crate::serialized_types::{
114    DEFAULT_MAX_SERIALIZED_RECORD_SIZE, LATEST_VERSION, Version, Versioned, VersionedLatest,
115};
116use anyhow::{Context, Error, anyhow, bail, ensure};
117use async_trait::async_trait;
118use either::Either::{Left, Right};
119use event_listener::EventListener;
120use fprint::TypeFingerprint;
121use fuchsia_inspect::HistogramProperty;
122use fuchsia_sync::Mutex;
123use futures::FutureExt;
124use futures::future::BoxFuture;
125use fxfs_macros::SerializeKey;
126use merge::{filter_marked_for_deletion, filter_tombstones, merge};
127use serde::{Deserialize, Serialize};
128use std::borrow::Borrow;
129use std::collections::{BTreeMap, HashSet, VecDeque};
130use std::hash::Hash;
131use std::marker::PhantomData;
132use std::num::{NonZero, Saturating};
133use std::ops::Range;
134use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
135use std::sync::{Arc, Weak};
136use storage_units::BlockSize;
137
138/// This trait is implemented by things that own reservations.
139pub trait ReservationOwner: Send + Sync {
140    /// Report that bytes are being released from the reservation back to the the |ReservationOwner|
141    /// where |owner_object_id| is the owner under the root object store associated with the
142    /// reservation.
143    fn release_reservation(&self, owner_object_id: Option<u64>, amount: u64);
144}
145
146/// A reservation guarantees that when it comes time to actually allocate, it will not fail due to
147/// lack of space.  Sub-reservations (a.k.a. holds) are possible which effectively allows part of a
148/// reservation to be set aside until it's time to commit.  Reservations do offer some
149/// thread-safety, but some responsibility is born by the caller: e.g. calling `forget` and
150/// `reserve` at the same time from different threads is unsafe. Reservations are have an
151/// |owner_object_id| which associates it with an object under the root object store that the
152/// reservation is accounted against.
153pub struct ReservationImpl<T: Borrow<U>, U: ReservationOwner + ?Sized> {
154    owner: T,
155    owner_object_id: Option<u64>,
156    inner: Mutex<ReservationInner>,
157    phantom: PhantomData<U>,
158}
159
160#[derive(Debug, Default)]
161struct ReservationInner {
162    // Amount currently held by this reservation.
163    amount: u64,
164
165    // Amount reserved by sub-reservations.
166    reserved: u64,
167}
168
169impl<T: Borrow<U>, U: ReservationOwner + ?Sized> std::fmt::Debug for ReservationImpl<T, U> {
170    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
171        self.inner.lock().fmt(f)
172    }
173}
174
175impl<T: Borrow<U> + Clone + Send + Sync, U: ReservationOwner + ?Sized> ReservationImpl<T, U> {
176    pub fn new(owner: T, owner_object_id: Option<u64>, amount: u64) -> Self {
177        Self {
178            owner,
179            owner_object_id,
180            inner: Mutex::new(ReservationInner { amount, reserved: 0 }),
181            phantom: PhantomData,
182        }
183    }
184
185    pub fn owner(&self) -> &T {
186        &self.owner
187    }
188
189    pub fn owner_object_id(&self) -> Option<u64> {
190        self.owner_object_id
191    }
192
193    /// Returns the total amount of the reservation, not accounting for anything that might be held.
194    pub fn amount(&self) -> u64 {
195        self.inner.lock().amount
196    }
197
198    /// Adds more to the reservation.
199    pub fn add(&self, amount: u64) {
200        self.inner.lock().amount += amount;
201    }
202
203    /// Returns the entire amount of the reservation.  The caller is responsible for maintaining
204    /// consistency, i.e. updating counters, etc, and there can be no sub-reservations (an assert
205    /// will fire otherwise).
206    pub fn forget(&self) -> u64 {
207        let mut inner = self.inner.lock();
208        assert_eq!(inner.reserved, 0);
209        std::mem::take(&mut inner.amount)
210    }
211
212    /// Takes some of the reservation.  The caller is responsible for maintaining consistency,
213    /// i.e. updating counters, etc.  This will assert that the amount being forgotten does not
214    /// exceed the available reservation amount; the caller should ensure that this is the case.
215    pub fn forget_some(&self, amount: u64) {
216        let mut inner = self.inner.lock();
217        inner.amount -= amount;
218        assert!(inner.reserved <= inner.amount);
219    }
220
221    /// Returns a partial amount of the reservation. |amount| is passed the |limit| and should
222    /// return the amount, which can be zero.
223    fn reserve_with(&self, amount: impl FnOnce(u64) -> u64) -> ReservationImpl<&Self, Self> {
224        let mut inner = self.inner.lock();
225        let taken = amount(inner.amount - inner.reserved);
226        inner.reserved += taken;
227        ReservationImpl::new(self, self.owner_object_id, taken)
228    }
229
230    /// Reserves *exactly* amount if possible.
231    pub fn reserve(&self, amount: u64) -> Option<ReservationImpl<&Self, Self>> {
232        let mut inner = self.inner.lock();
233        if inner.amount - inner.reserved < amount {
234            None
235        } else {
236            inner.reserved += amount;
237            Some(ReservationImpl::new(self, self.owner_object_id, amount))
238        }
239    }
240
241    /// Commits a previously reserved amount from this reservation.  The caller is responsible for
242    /// ensuring the amount was reserved.
243    pub fn commit(&self, amount: u64) {
244        let mut inner = self.inner.lock();
245        inner.reserved -= amount;
246        inner.amount -= amount;
247    }
248
249    /// Returns some of the reservation.
250    pub fn give_back(&self, amount: u64) {
251        self.owner.borrow().release_reservation(self.owner_object_id, amount);
252        let mut inner = self.inner.lock();
253        inner.amount -= amount;
254        assert!(inner.reserved <= inner.amount);
255    }
256
257    /// Moves `amount` from this reservation to another reservation.
258    pub fn move_to<V: Borrow<W> + Clone + Send + Sync, W: ReservationOwner + ?Sized>(
259        &self,
260        other: &ReservationImpl<V, W>,
261        amount: u64,
262    ) {
263        assert_eq!(self.owner_object_id, other.owner_object_id());
264        let mut inner = self.inner.lock();
265        if let Some(amount) = inner.amount.checked_sub(amount) {
266            inner.amount = amount;
267        } else {
268            std::mem::drop(inner);
269            panic!("Insufficient reservation space");
270        }
271        other.add(amount);
272    }
273}
274
275impl<T: Borrow<U>, U: ReservationOwner + ?Sized> Drop for ReservationImpl<T, U> {
276    fn drop(&mut self) {
277        let inner = self.inner.get_mut();
278        assert_eq!(inner.reserved, 0);
279        let owner_object_id = self.owner_object_id;
280        if inner.amount > 0 {
281            self.owner
282                .borrow()
283                .release_reservation(owner_object_id, std::mem::take(&mut inner.amount));
284        }
285    }
286}
287
288impl<T: Borrow<U> + Send + Sync, U: ReservationOwner + ?Sized> ReservationOwner
289    for ReservationImpl<T, U>
290{
291    fn release_reservation(&self, owner_object_id: Option<u64>, amount: u64) {
292        // Sub-reservations should belong to the same owner (or lack thereof).
293        assert_eq!(owner_object_id, self.owner_object_id);
294        let mut inner = self.inner.lock();
295        assert!(inner.reserved >= amount, "{} >= {}", inner.reserved, amount);
296        inner.reserved -= amount;
297    }
298}
299
300pub type Reservation = ReservationImpl<Arc<dyn ReservationOwner>, dyn ReservationOwner>;
301
302pub type Hold<'a> = ReservationImpl<&'a Reservation, Reservation>;
303
304/// Our allocator implementation tracks extents with a reference count.  At time of writing, these
305/// reference counts should never exceed 1, but that might change with snapshots and clones.
306pub type AllocatorKey = AllocatorKeyV32;
307
308#[derive(
309    Clone,
310    Debug,
311    Deserialize,
312    Eq,
313    Hash,
314    Ord,
315    PartialEq,
316    PartialOrd,
317    Serialize,
318    TypeFingerprint,
319    SerializeKey,
320    Versioned,
321)]
322#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
323pub struct AllocatorKeyV32 {
324    pub device_range: Extent,
325}
326
327impl SortByU64 for AllocatorKey {
328    fn get_leading_u64(&self) -> u64 {
329        self.device_range.end / crate::object_store::extent::MIN_BLOCK_SIZE
330    }
331}
332
333const EXTENT_HASH_BUCKET_SIZE: BlockSize = BlockSize::SIZE_1MIB;
334
335pub struct AllocatorKeyPartitionIterator {
336    device_range: Range<u64>,
337}
338
339impl Iterator for AllocatorKeyPartitionIterator {
340    type Item = u64;
341
342    fn next(&mut self) -> Option<Self::Item> {
343        if self.device_range.start >= self.device_range.end {
344            None
345        } else {
346            let start = self.device_range.start;
347            self.device_range.start = start.saturating_add(EXTENT_HASH_BUCKET_SIZE.get());
348            let end = std::cmp::min(self.device_range.start, self.device_range.end);
349            let key = AllocatorKey { device_range: Extent(start..end) };
350            let hash = crate::stable_hash::stable_hash(key);
351            Some(hash)
352        }
353    }
354
355    fn size_hint(&self) -> (usize, Option<usize>) {
356        let len = if self.device_range.start >= self.device_range.end {
357            0
358        } else {
359            let diff = self.device_range.end - self.device_range.start;
360            let count = EXTENT_HASH_BUCKET_SIZE.align_up_to_blocks(diff);
361            usize::try_from(count).unwrap_or(usize::MAX)
362        };
363        (len, Some(len))
364    }
365}
366
367impl ExactSizeIterator for AllocatorKeyPartitionIterator {}
368
369impl FuzzyHash for AllocatorKey {
370    fn fuzzy_hash(&self) -> impl ExactSizeIterator<Item = u64> {
371        AllocatorKeyPartitionIterator {
372            device_range: EXTENT_HASH_BUCKET_SIZE.align_down(self.device_range.start)
373                ..EXTENT_HASH_BUCKET_SIZE.align_up(self.device_range.end).unwrap_or(u64::MAX),
374        }
375    }
376
377    fn is_range_key(&self) -> bool {
378        true
379    }
380}
381
382impl AllocatorKey {
383    /// Returns a search key for `merge_into` that finds any touching predecessor to coalesce with.
384    ///
385    /// Extents in the tree are sorted by `cmp_upper_bound` as `(end, -start)`:
386    ///
387    /// ```text
388    /// Tree:       [  Predecessor  ) [ Key to insert )
389    /// Offsets:    X              10                20
390    ///
391    /// Search key:                 | (10..10)
392    /// ```
393    ///
394    /// Searching with `10..10` has end=10, len=0. Because shorter extents sort earlier when
395    /// ends match, `10..10` sorts before any valid extent `X..10` (len > 0). This positions the
396    /// iterator at `X..10` so it can be merged with `10..20`.
397    pub fn lower_bound_for_merge_into(&self) -> AllocatorKey {
398        AllocatorKey { device_range: self.device_range.key_for_merge_into() }
399    }
400}
401
402impl LayerKey for AllocatorKey {
403    fn merge_type(&self) -> MergeType {
404        MergeType::OptimizedMerge
405    }
406
407    fn next_key(&self) -> Option<Self> {
408        Some(Self { device_range: Extent::search_key_from_offset(self.device_range.end) })
409    }
410
411    fn search_key(&self) -> Option<Self> {
412        Some(Self { device_range: self.device_range.search_key() })
413    }
414
415    fn is_search_key(&self) -> bool {
416        self.device_range.is_search_key()
417    }
418
419    fn overlaps(&self, other: &Self) -> bool {
420        self.device_range.overlaps(&other.device_range)
421    }
422}
423
424impl OrdUpperBound for AllocatorKey {
425    fn cmp_upper_bound(&self, other: &AllocatorKey) -> std::cmp::Ordering {
426        // Defer to cmp_upper_bound ordering provided by Extent type which
427        // uses (end, len) to order ranges.
428        self.device_range.cmp_upper_bound(&other.device_range)
429    }
430}
431
432impl OrdLowerBound for AllocatorKey {
433    fn cmp_lower_bound(&self, other: &AllocatorKey) -> std::cmp::Ordering {
434        // The ordering over range.end is significant here as it is used in
435        // the heap ordering that feeds into our merge function and
436        // a total ordering over range lets us remove a symmetry case from
437        // the allocator merge function.
438        self.device_range.cmp(&other.device_range)
439    }
440}
441
442/// Allocations are "owned" by a single ObjectStore and are reference counted
443/// (for future snapshot/clone support).
444pub type AllocatorValue = AllocatorValueV32;
445impl Value for AllocatorValue {
446    const DELETED_MARKER: Self = Self::None;
447}
448
449#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize, TypeFingerprint, Versioned)]
450#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
451pub enum AllocatorValueV32 {
452    // Tombstone variant indicating an extent is no longer allocated.
453    None,
454    // Used when we know there are no possible allocations below us in the stack.
455    // This is currently all the time. We used to have a related Delta type but
456    // it has been removed due to correctness issues (https://fxbug.dev/42179428).
457    Abs { count: u64, owner_object_id: u64 },
458}
459
460pub type AllocatorItem = Item<AllocatorKey, AllocatorValue>;
461
462/// Serialized information about the allocator.
463pub type AllocatorInfo = AllocatorInfoV32;
464
465#[derive(Debug, Default, Clone, Deserialize, Serialize, TypeFingerprint, Versioned)]
466pub struct AllocatorInfoV32 {
467    /// Holds the set of layer file object_id for the LSM tree (newest first).
468    pub layers: Vec<u64>,
469    /// Maps from owner_object_id to bytes allocated.
470    pub allocated_bytes: BTreeMap<u64, u64>,
471    /// Set of owner_object_id that we should ignore if found in layer files.  For now, this should
472    /// always be empty on-disk because we always do full compactions.
473    pub marked_for_deletion: HashSet<u64>,
474    /// The limit for the number of allocates bytes per `owner_object_id` whereas the value. If
475    /// there is no limit present here for an `owner_object_id` assume it is max u64.
476    pub limit_bytes: BTreeMap<u64, u64>,
477}
478
479const MAX_ALLOCATOR_INFO_SERIALIZED_SIZE: usize = 131_072;
480
481/// Computes the target maximum extent size based on the block size of the allocator.
482pub fn max_extent_size_for_block_size(block_size: BlockSize) -> u64 {
483    // Each block in an extent contains an 8-byte checksum (which due to varint encoding is 9
484    // bytes), and a given extent record must be no larger DEFAULT_MAX_SERIALIZED_RECORD_SIZE.  We
485    // also need to leave a bit of room (arbitrarily, 64 bytes) for the rest of the extent's
486    // metadata.
487    block_size * ((DEFAULT_MAX_SERIALIZED_RECORD_SIZE - 64) / 9)
488}
489
490#[derive(Default)]
491struct AllocatorCounters {
492    num_flushes: u64,
493    last_flush_time: Option<std::time::SystemTime>,
494}
495
496pub struct Allocator {
497    filesystem: Weak<FxFilesystem>,
498    block_size: BlockSize,
499    device_size: u64,
500    object_id: u64,
501    max_extent_size_bytes: u64,
502    tree: LSMTree<AllocatorKey, AllocatorValue>,
503    // A list of allocations which are temporary; i.e. they are not actually stored in the LSM tree,
504    // but we still don't want to be able to allocate over them.  This layer is merged into the
505    // allocator's LSM tree whilst reading it, so the allocations appear to exist in the LSM tree.
506    // This is used in a few places, for example to hold back allocations that have been added to a
507    // transaction but are not yet committed yet, or to prevent the allocation of a deleted extent
508    // until the device is flushed.
509    temporary_allocations: Arc<SkipListLayer<AllocatorKey, AllocatorValue>>,
510    inner: Mutex<Inner>,
511    allocation_mutex: futures::lock::Mutex<()>,
512    counters: Mutex<AllocatorCounters>,
513    maximum_offset: AtomicU64,
514    allocations_allowed: AtomicBool,
515}
516
517/// Tracks the different stages of byte allocations for an individual owner.
518#[derive(Debug, Default, PartialEq)]
519struct ByteTracking {
520    /// This value is the up-to-date count of the number of allocated bytes per owner_object_id
521    /// whereas the value in `Info::allocated_bytes` is the value as it was when we last flushed.
522    /// This field should be regarded as *untrusted*; it can be invalid due to filesystem
523    /// inconsistencies, and this is why it is `Saturating<u64>` rather than just `u64`.  Any
524    /// computations that use this value should typically propagate `Saturating` so that it is clear
525    /// the value is *untrusted*.
526    allocated_bytes: Saturating<u64>,
527
528    /// This value is the number of bytes allocated to uncommitted allocations.
529    /// (Bytes allocated, but not yet persisted to disk)
530    uncommitted_allocated_bytes: u64,
531
532    /// This value is the number of bytes allocated to reservations.
533    reserved_bytes: u64,
534
535    /// Committed deallocations that we cannot use until they are flushed to the device.  Each entry
536    /// in this list is the log file offset at which it was committed and an array of deallocations
537    /// that occurred at that time.
538    committed_deallocated_bytes: u64,
539}
540
541impl ByteTracking {
542    // Returns the total number of bytes that are taken either from reservations, allocations or
543    // uncommitted allocations.
544    fn used_bytes(&self) -> Saturating<u64> {
545        self.allocated_bytes + Saturating(self.uncommitted_allocated_bytes + self.reserved_bytes)
546    }
547
548    // Returns the amount that is not available to be allocated, which includes actually allocated
549    // bytes, bytes that have been allocated for a transaction but the transaction hasn't committed
550    // yet, and bytes that have been deallocated, but the device hasn't been flushed yet so we can't
551    // reuse those bytes yet.
552    fn unavailable_bytes(&self) -> Saturating<u64> {
553        self.allocated_bytes
554            + Saturating(self.uncommitted_allocated_bytes)
555            + Saturating(self.committed_deallocated_bytes)
556    }
557
558    // Like `unavailable_bytes`, but treats as available the bytes which have been deallocated and
559    // require a device flush to be reused.
560    fn unavailable_after_sync_bytes(&self) -> Saturating<u64> {
561        self.allocated_bytes + Saturating(self.uncommitted_allocated_bytes)
562    }
563}
564
565#[derive(Debug)]
566struct CommittedDeallocation {
567    // The offset at which this deallocation was committed.
568    log_file_offset: u64,
569    // The device range being deallocated.
570    range: Range<u64>,
571    // The owning object id which originally allocated it.
572    owner_object_id: u64,
573}
574
575struct Inner {
576    info: AllocatorInfo,
577
578    /// The allocator can only be opened if there have been no allocations and it has not already
579    /// been opened or initialized.
580    opened: bool,
581
582    /// When we allocate a range from RAM, we add it to Allocator::temporary_allocations.
583    /// This layer is added to layer_set in rebuild_strategy and take_for trim so we don't assume
584    /// the range is available.
585    /// If we apply the mutation, we move it from `temporary_allocations` to the LSMTree.
586    /// If we drop the mutation, we just delete it from `temporary_allocations` and call `free`.
587    ///
588    /// We need to be very careful about races when we move the allocation into the LSMTree.
589    /// A layer_set is a snapshot in time of a set of layers. Say:
590    ///  1. `rebuild_strategy` takes a layer_set that includes the mutable layer.
591    ///  2. A compaction operation seals that mutable layer and creates a new mutable layer.
592    ///     Note that the layer_set in rebuild_strategy doesn't have this new mutable layer.
593    ///  3. An apply_mutation operation adds the allocation to the new mutable layer.
594    ///     It then removes the allocation from `temporary_allocations`.
595    ///  4. `rebuild_strategy` iterates the layer_set, missing both the temporary mutation AND
596    ///     the record in the new mutable layer, making the allocated range available for double
597    ///     allocation and future filesystem corruption.
598    ///
599    /// We don't have (or want) a means to lock compctions during rebuild_strategy operations so
600    /// to avoid this scenario, we can't remove from temporary_allocations when we apply a mutation.
601    /// Instead we add them to dropped_temporary_allocations and make the actual removals from
602    /// `temporary_allocations` while holding the allocator lock. Because rebuild_strategy only
603    /// runs with this lock held, this guarantees that we don't remove entries from temporary
604    /// mutations while also iterating over it.
605    ///
606    /// A related issue happens when we deallocate a range. We use temporary_allocations this time
607    /// to prevent reuse until after the deallocation has been successfully flushed.
608    /// In this scenario we don't want to call 'free()' until the range is ready to use again.
609    dropped_temporary_allocations: Vec<Range<u64>>,
610
611    /// The per-owner counters for bytes at various stages of the data life-cycle. From initial
612    /// reservation through until the bytes are unallocated and eventually uncommitted.
613    owner_bytes: BTreeMap<u64, ByteTracking>,
614
615    /// This value is the number of bytes allocated to reservations but not tracked as part of a
616    /// particular volume.
617    unattributed_reserved_bytes: u64,
618
619    /// Committed deallocations that we cannot use until they are flushed to the device.
620    committed_deallocated: VecDeque<CommittedDeallocation>,
621
622    /// Bytes which are currently being trimmed.  These can still be allocated from, but the
623    /// allocation will block until the current batch of trimming is finished.
624    trim_reserved_bytes: u64,
625
626    /// While a trim is being performed, this listener is set.  When the current batch of extents
627    /// being trimmed have been released (and trim_reserved_bytes is 0), this is signaled.
628    /// This should only be listened to while the allocation_mutex is held.
629    trim_listener: Option<EventListener>,
630
631    /// This controls how we allocate our free space to manage fragmentation.
632    strategy: strategy::BestFit,
633
634    /// Tracks the number of allocations of size 1,2,...63,>=64.
635    allocation_size_histogram: [u64; 64],
636    /// Tracks which size bucket triggers rebuild_strategy.
637    rebuild_strategy_trigger_histogram: [u64; 64],
638
639    /// This is distinct from the set contained in `info`.  New entries are inserted *after* a
640    /// device has been flushed (it is not safe to reuse the space taken by a deleted volume prior
641    /// to this) and are removed after a major compaction.
642    marked_for_deletion: HashSet<u64>,
643
644    /// The set of volumes deleted that are waiting for a sync.
645    volumes_deleted_pending_sync: HashSet<u64>,
646}
647
648impl Inner {
649    fn allocated_bytes(&self) -> Saturating<u64> {
650        let mut total = Saturating(0);
651        for (_, bytes) in &self.owner_bytes {
652            total += bytes.allocated_bytes;
653        }
654        total
655    }
656
657    fn uncommitted_allocated_bytes(&self) -> u64 {
658        self.owner_bytes.values().map(|x| &x.uncommitted_allocated_bytes).sum()
659    }
660
661    fn reserved_bytes(&self) -> u64 {
662        self.owner_bytes.values().map(|x| &x.reserved_bytes).sum::<u64>()
663            + self.unattributed_reserved_bytes
664    }
665
666    fn owner_id_limit_bytes(&self, owner_object_id: u64) -> u64 {
667        match self.info.limit_bytes.get(&owner_object_id) {
668            Some(v) => *v,
669            None => u64::MAX,
670        }
671    }
672
673    fn owner_id_bytes_left(&self, owner_object_id: u64) -> u64 {
674        let limit = self.owner_id_limit_bytes(owner_object_id);
675        let used = self.owner_bytes.get(&owner_object_id).map_or(Saturating(0), |b| b.used_bytes());
676        (Saturating(limit) - used).0
677    }
678
679    // Returns the amount that is not available to be allocated, which includes actually allocated
680    // bytes, bytes that have been allocated for a transaction but the transaction hasn't committed
681    // yet, and bytes that have been deallocated, but the device hasn't been flushed yet so we can't
682    // reuse those bytes yet.
683    fn unavailable_bytes(&self) -> Saturating<u64> {
684        let mut total = Saturating(0);
685        for (_, bytes) in &self.owner_bytes {
686            total += bytes.unavailable_bytes();
687        }
688        total
689    }
690
691    // Returns the total number of bytes that are taken either from reservations, allocations or
692    // uncommitted allocations.
693    fn used_bytes(&self) -> Saturating<u64> {
694        let mut total = Saturating(0);
695        for (_, bytes) in &self.owner_bytes {
696            total += bytes.used_bytes();
697        }
698        total + Saturating(self.unattributed_reserved_bytes)
699    }
700
701    // Like `unavailable_bytes`, but treats as available the bytes which have been deallocated and
702    // require a device flush to be reused.
703    fn unavailable_after_sync_bytes(&self) -> Saturating<u64> {
704        let mut total = Saturating(0);
705        for (_, bytes) in &self.owner_bytes {
706            total += bytes.unavailable_after_sync_bytes();
707        }
708        total
709    }
710
711    // Returns the number of bytes which will be available after the current batch of trimming
712    // completes.
713    fn bytes_available_not_being_trimmed(&self, device_size: u64) -> Result<u64, Error> {
714        device_size
715            .checked_sub(
716                (self.unavailable_after_sync_bytes() + Saturating(self.trim_reserved_bytes)).0,
717            )
718            .ok_or_else(|| anyhow!(FxfsError::Inconsistent))
719    }
720
721    fn add_reservation(&mut self, owner_object_id: Option<u64>, amount: u64) {
722        match owner_object_id {
723            Some(owner) => self.owner_bytes.entry(owner).or_default().reserved_bytes += amount,
724            None => self.unattributed_reserved_bytes += amount,
725        };
726    }
727
728    fn remove_reservation(&mut self, owner_object_id: Option<u64>, amount: u64) {
729        match owner_object_id {
730            Some(owner) => {
731                let owner_entry = self.owner_bytes.entry(owner).or_default();
732                assert!(
733                    owner_entry.reserved_bytes >= amount,
734                    "{} >= {}",
735                    owner_entry.reserved_bytes,
736                    amount
737                );
738                owner_entry.reserved_bytes -= amount;
739            }
740            None => {
741                assert!(
742                    self.unattributed_reserved_bytes >= amount,
743                    "{} >= {}",
744                    self.unattributed_reserved_bytes,
745                    amount
746                );
747                self.unattributed_reserved_bytes -= amount
748            }
749        };
750    }
751}
752
753/// A container for a set of extents which are known to be free and can be trimmed.  Returned by
754/// `take_for_trimming`.
755pub struct TrimmableExtents<'a> {
756    allocator: &'a Allocator,
757    extents: Vec<Range<u64>>,
758    // The allocator can subscribe to this event to wait until these extents are dropped.  This way,
759    // we don't fail an allocation attempt if blocks are tied up for trimming; rather, we just wait
760    // until the batch is finished with and then proceed.
761    _drop_event: DropEvent,
762}
763
764impl<'a> TrimmableExtents<'a> {
765    pub fn extents(&self) -> &Vec<Range<u64>> {
766        &self.extents
767    }
768
769    // Also returns an EventListener which is signaled when this is dropped.
770    fn new(allocator: &'a Allocator) -> (Self, EventListener) {
771        let drop_event = DropEvent::new();
772        let listener = drop_event.listen();
773        (Self { allocator, extents: vec![], _drop_event: drop_event }, listener)
774    }
775
776    fn add_extent(&mut self, extent: Range<u64>) {
777        self.extents.push(extent);
778    }
779}
780
781impl<'a> Drop for TrimmableExtents<'a> {
782    fn drop(&mut self) {
783        let mut inner = self.allocator.inner.lock();
784        for device_range in std::mem::take(&mut self.extents) {
785            inner.strategy.free(device_range.clone()).expect("drop trim extent");
786            self.allocator
787                .temporary_allocations
788                .erase(&AllocatorKey { device_range: Extent(device_range) });
789        }
790        inner.trim_reserved_bytes = 0;
791    }
792}
793
794impl Allocator {
795    pub fn new(filesystem: Arc<FxFilesystem>, object_id: u64) -> Allocator {
796        let block_size = filesystem.block_size();
797        // We expect device size to be a multiple of block size. Throw away any tail.
798        let device_size = block_size.align_down(filesystem.device().size());
799        if device_size != filesystem.device().size() {
800            warn!("Device size is not block aligned. Rounding down.");
801        }
802        let max_extent_size_bytes = max_extent_size_for_block_size(block_size);
803        let mut strategy = strategy::BestFit::default();
804        strategy.free(0..device_size).expect("new fs");
805        Allocator {
806            filesystem: Arc::downgrade(&filesystem),
807            block_size,
808            device_size,
809            object_id,
810            max_extent_size_bytes,
811            tree: LSMTree::new(merge, None),
812            temporary_allocations: SkipListLayer::new(1024),
813            inner: Mutex::new(Inner {
814                info: AllocatorInfo::default(),
815                opened: false,
816                dropped_temporary_allocations: Vec::new(),
817                owner_bytes: BTreeMap::new(),
818                unattributed_reserved_bytes: 0,
819                committed_deallocated: VecDeque::new(),
820                trim_reserved_bytes: 0,
821                trim_listener: None,
822                strategy,
823                allocation_size_histogram: [0; 64],
824                rebuild_strategy_trigger_histogram: [0; 64],
825                marked_for_deletion: HashSet::new(),
826                volumes_deleted_pending_sync: HashSet::new(),
827            }),
828            allocation_mutex: futures::lock::Mutex::new(()),
829            counters: Mutex::new(AllocatorCounters::default()),
830            maximum_offset: AtomicU64::new(0),
831            allocations_allowed: AtomicBool::new(filesystem.options().image_builder_mode.is_none()),
832        }
833    }
834
835    pub fn tree(&self) -> &LSMTree<AllocatorKey, AllocatorValue> {
836        &self.tree
837    }
838
839    /// Enables allocations.
840    /// This is only valid to call if the allocator was created in image builder mode.
841    pub fn enable_allocations(&self) {
842        self.allocations_allowed.store(true, Ordering::SeqCst);
843    }
844
845    /// Returns an iterator that yields all allocations, filtering out tombstones and any
846    /// owner_object_id that have been marked as deleted.  If `committed_marked_for_deletion` is
847    /// true, then filter using the committed volumes marked for deletion rather than the in-memory
848    /// copy which excludes volumes that have been deleted but there hasn't been a sync yet.
849    pub async fn filter<'a>(
850        &self,
851        iter: impl LayerIterator<AllocatorKey, AllocatorValue> + 'a,
852        committed_marked_for_deletion: bool,
853    ) -> Result<impl LayerIterator<AllocatorKey, AllocatorValue> + 'a, Error> {
854        let marked_for_deletion = {
855            let inner = self.inner.lock();
856            if committed_marked_for_deletion {
857                &inner.info.marked_for_deletion
858            } else {
859                &inner.marked_for_deletion
860            }
861            .clone()
862        };
863        let iter =
864            filter_marked_for_deletion(filter_tombstones(iter).await?, marked_for_deletion).await?;
865        Ok(iter)
866    }
867
868    /// A histogram of allocation request sizes.
869    /// The index into the array is 'number of blocks'.
870    /// The last bucket is a catch-all for larger allocation requests.
871    pub fn allocation_size_histogram(&self) -> [u64; 64] {
872        self.inner.lock().allocation_size_histogram
873    }
874
875    /// Creates a new (empty) allocator.
876    pub async fn create(&self, transaction: &mut Transaction<'_>) -> Result<(), Error> {
877        // Mark the allocator as opened before creating the file because creating a new
878        // transaction requires a reservation.
879        assert_eq!(std::mem::replace(&mut self.inner.lock().opened, true), false);
880
881        let filesystem = self.filesystem.upgrade().unwrap();
882        let root_store = filesystem.root_store();
883        ObjectStore::create_object_with_id(
884            &root_store,
885            transaction,
886            ReservedId::new(&root_store, NonZero::new(self.object_id()).unwrap()),
887            HandleOptions::default(),
888            None,
889        )?;
890        root_store.update_last_object_id(self.object_id());
891        Ok(())
892    }
893
894    // Opens the allocator.  This is not thread-safe; this should be called prior to
895    // replaying the allocator's mutations.
896    pub async fn open(self: &Arc<Self>) -> Result<(), Error> {
897        let filesystem = self.filesystem.upgrade().unwrap();
898        let root_store = filesystem.root_store();
899
900        self.inner.lock().strategy = strategy::BestFit::default();
901
902        let handle =
903            ObjectStore::open_object(&root_store, self.object_id, HandleOptions::default(), None)
904                .await
905                .context("Failed to open allocator object")?;
906
907        if handle.get_size() > 0 {
908            let serialized_info = handle
909                .contents(MAX_ALLOCATOR_INFO_SERIALIZED_SIZE)
910                .await
911                .context("Failed to read AllocatorInfo")?;
912            let mut cursor = std::io::Cursor::new(serialized_info);
913            let (info, _version) = AllocatorInfo::deserialize_with_version(&mut cursor)
914                .context("Failed to deserialize AllocatorInfo")?;
915
916            let (layers, total_size) = open_layers(&root_store, info.layers.iter().cloned(), None)
917                .await
918                .context("Failed to open allocator layer file")?;
919
920            {
921                let mut inner = self.inner.lock();
922
923                // Check allocated_bytes fits within the device.
924                let mut device_bytes = self.device_size;
925                for (&owner_object_id, &bytes) in &info.allocated_bytes {
926                    ensure!(
927                        bytes <= device_bytes,
928                        anyhow!(FxfsError::Inconsistent).context(format!(
929                            "Allocated bytes exceeds device size: {:?}",
930                            info.allocated_bytes
931                        ))
932                    );
933                    device_bytes -= bytes;
934
935                    inner.owner_bytes.entry(owner_object_id).or_default().allocated_bytes =
936                        Saturating(bytes);
937                }
938
939                inner.info = info;
940            }
941
942            self.tree.append_open_layers(layers);
943            self.filesystem.upgrade().unwrap().object_manager().update_reservation(
944                self.object_id,
945                tree::reservation_amount_from_layer_size(total_size),
946            );
947        }
948
949        Ok(())
950    }
951
952    pub async fn on_replay_complete(self: &Arc<Self>) -> Result<(), Error> {
953        // We can assume the device has been flushed.
954        {
955            let mut inner = self.inner.lock();
956            inner.volumes_deleted_pending_sync.clear();
957            inner.marked_for_deletion = inner.info.marked_for_deletion.clone();
958        }
959
960        // Build free extent structure from disk. For now, we expect disks to have *some* free
961        // space at all times. This may change if we ever support mounting of read-only
962        // redistributable filesystem images.
963        if !self.rebuild_strategy().await.context("Build free extents")? {
964            if self.filesystem.upgrade().unwrap().options().read_only {
965                info!("Device contains no free space (read-only mode).");
966            } else {
967                info!("Device contains no free space.");
968                return Err(FxfsError::Inconsistent)
969                    .context("Device appears to contain no free space");
970            }
971        }
972
973        assert_eq!(std::mem::replace(&mut self.inner.lock().opened, true), false);
974        Ok(())
975    }
976
977    /// Walk all allocations to generate the set of free regions between allocations.
978    /// It is safe to re-run this on live filesystems but it should not be called concurrently
979    /// with allocations, trims or other rebuild_strategy invocations -- use allocation_mutex.
980    /// Returns true if the rebuild made changes to either the free ranges or overflow markers.
981    async fn rebuild_strategy(self: &Arc<Self>) -> Result<bool, Error> {
982        let mut changed = false;
983        let mut layer_set = self.tree.empty_layer_set();
984        layer_set.layers.push((self.temporary_allocations.clone() as Arc<dyn Layer<_, _>>).into());
985        self.tree.add_all_layers_to_layer_set(&mut layer_set);
986
987        let overflow_markers = self.inner.lock().strategy.overflow_markers();
988        self.inner.lock().strategy.reset_overflow_markers();
989
990        let mut to_add = Vec::new();
991        let mut merger = layer_set.merger();
992        let mut iter = self.filter(merger.query(Query::FullScan).await?, false).await?;
993        let mut last_offset = 0;
994        while last_offset < self.device_size {
995            let next_range = match iter.get() {
996                None => {
997                    assert!(last_offset <= self.device_size);
998                    let range = last_offset..self.device_size;
999                    last_offset = self.device_size;
1000                    iter.advance().await?;
1001                    range
1002                }
1003                Some(ItemRef { key, .. }) => {
1004                    let device_range = &key.device_range;
1005                    if device_range.end < last_offset {
1006                        iter.advance().await?;
1007                        continue;
1008                    }
1009                    if device_range.start <= last_offset {
1010                        last_offset = device_range.end;
1011                        iter.advance().await?;
1012                        continue;
1013                    }
1014                    let range = last_offset..device_range.start;
1015                    last_offset = device_range.end;
1016                    iter.advance().await?;
1017                    range
1018                }
1019            };
1020            to_add.push(next_range);
1021            // Avoid taking a lock on inner for every free range.
1022            if to_add.len() > 100 {
1023                let mut inner = self.inner.lock();
1024                for range in to_add.drain(..) {
1025                    changed |= inner.strategy.force_free(range)?;
1026                }
1027            }
1028        }
1029        let mut inner = self.inner.lock();
1030        for range in to_add {
1031            changed |= inner.strategy.force_free(range)?;
1032        }
1033        if overflow_markers != inner.strategy.overflow_markers() {
1034            changed = true;
1035        }
1036        Ok(changed)
1037    }
1038
1039    /// Collects up to `extents_per_batch` free extents of size up to `max_extent_size` from
1040    /// `offset`.  The extents will be reserved.
1041    /// Note that only one `FreeExtents` can exist for the allocator at any time.
1042    pub async fn take_for_trimming(
1043        &self,
1044        offset: u64,
1045        max_extent_size: usize,
1046        extents_per_batch: usize,
1047    ) -> Result<TrimmableExtents<'_>, Error> {
1048        let _guard = self.allocation_mutex.lock().await;
1049
1050        let (mut result, listener) = TrimmableExtents::new(self);
1051        if offset >= self.device_size {
1052            return Ok(result);
1053        }
1054        let mut bytes = 0;
1055
1056        // We can't just use self.strategy here because it doesn't necessarily hold all
1057        // free extent metadata in RAM.
1058        let mut layer_set = self.tree.empty_layer_set();
1059        layer_set.layers.push((self.temporary_allocations.clone() as Arc<dyn Layer<_, _>>).into());
1060        self.tree.add_all_layers_to_layer_set(&mut layer_set);
1061        let mut merger = layer_set.merger();
1062        let mut iter = self
1063            .filter(
1064                merger
1065                    .query(Query::FullRange(&AllocatorKey {
1066                        device_range: Extent::search_key_from_offset(offset),
1067                    }))
1068                    .await?,
1069                false,
1070            )
1071            .await?;
1072        let mut last_offset = offset;
1073        'outer: while last_offset < self.device_size {
1074            let mut range = match iter.get() {
1075                None => {
1076                    assert!(last_offset <= self.device_size);
1077                    let range = last_offset..self.device_size;
1078                    last_offset = self.device_size;
1079                    iter.advance().await?;
1080                    range
1081                }
1082                Some(ItemRef { key: AllocatorKey { device_range, .. }, .. }) => {
1083                    if device_range.end <= last_offset {
1084                        iter.advance().await?;
1085                        continue;
1086                    }
1087                    if device_range.start <= last_offset {
1088                        last_offset = device_range.end;
1089                        iter.advance().await?;
1090                        continue;
1091                    }
1092                    let range = last_offset..device_range.start;
1093                    last_offset = device_range.end;
1094                    iter.advance().await?;
1095                    range
1096                }
1097            };
1098            if range.start < offset {
1099                continue;
1100            }
1101
1102            // 'range' is based on the on-disk LSM tree. We need to check against any uncommitted
1103            // allocations and remove temporarily allocate the range for ourselves.
1104            let mut inner = self.inner.lock();
1105
1106            while range.start < range.end {
1107                let prefix =
1108                    range.start..std::cmp::min(range.start + max_extent_size as u64, range.end);
1109                range = prefix.end..range.end;
1110                bytes += prefix.length()?;
1111                // We can assume the following are safe because we've checked both LSM tree and
1112                // temporary_allocations while holding the lock.
1113                inner.strategy.remove(prefix.clone());
1114                self.temporary_allocations.insert(AllocatorItem {
1115                    key: AllocatorKey { device_range: Extent(prefix.clone()) },
1116                    value: AllocatorValue::Abs { owner_object_id: INVALID_OBJECT_ID, count: 1 },
1117                })?;
1118                result.add_extent(prefix);
1119                if result.extents.len() == extents_per_batch {
1120                    break 'outer;
1121                }
1122            }
1123            if result.extents.len() == extents_per_batch {
1124                break 'outer;
1125            }
1126        }
1127        {
1128            let mut inner = self.inner.lock();
1129
1130            ensure!(inner.trim_reserved_bytes == 0, FxfsError::AlreadyBound);
1131            inner.trim_listener = Some(listener);
1132            inner.trim_reserved_bytes = bytes;
1133            debug_assert!(
1134                (Saturating(inner.trim_reserved_bytes) + inner.unavailable_bytes()).0
1135                    <= self.device_size
1136            );
1137        }
1138        Ok(result)
1139    }
1140
1141    /// Returns all objects that exist in the parent store that pertain to this allocator.
1142    pub fn parent_objects(&self) -> Vec<u64> {
1143        // The allocator tree needs to store a file for each of the layers in the tree, so we return
1144        // those, since nothing else references them.
1145        self.inner.lock().info.layers.clone()
1146    }
1147
1148    /// Returns all the current owner byte limits (in pairs of `(owner_object_id, bytes)`).
1149    pub fn owner_byte_limits(&self) -> Vec<(u64, u64)> {
1150        self.inner.lock().info.limit_bytes.iter().map(|(k, v)| (*k, *v)).collect()
1151    }
1152
1153    /// Returns (allocated_bytes, byte_limit) for the given owner.
1154    pub fn owner_allocation_info(&self, owner_object_id: u64) -> (u64, Option<u64>) {
1155        let inner = self.inner.lock();
1156        (
1157            inner.owner_bytes.get(&owner_object_id).map(|b| b.used_bytes().0).unwrap_or(0u64),
1158            inner.info.limit_bytes.get(&owner_object_id).copied(),
1159        )
1160    }
1161
1162    /// Returns owner bytes debug information.
1163    pub fn owner_bytes_debug(&self) -> String {
1164        format!("{:?}", self.inner.lock().owner_bytes)
1165    }
1166
1167    fn needs_sync(&self) -> bool {
1168        // TODO(https://fxbug.dev/42178048): This will only trigger if *all* free space is taken up with
1169        // committed deallocated bytes, but we might want to trigger a sync if we're low and there
1170        // happens to be a lot of deallocated bytes as that might mean we can fully satisfy
1171        // allocation requests.
1172        let inner = self.inner.lock();
1173        inner.unavailable_bytes().0 >= self.device_size
1174    }
1175
1176    fn is_system_store(&self, owner_object_id: u64) -> bool {
1177        let fs = self.filesystem.upgrade().unwrap();
1178        owner_object_id == fs.object_manager().root_store_object_id()
1179            || owner_object_id == fs.object_manager().root_parent_store_object_id()
1180    }
1181
1182    /// Updates the accounting to track that a byte reservation has been moved out of an owner to
1183    /// the unattributed pool.
1184    pub fn disown_reservation(&self, old_owner_object_id: Option<u64>, amount: u64) {
1185        if old_owner_object_id.is_none() || amount == 0 {
1186            return;
1187        }
1188        // These 2 mutations should behave as though they're a single atomic mutation.
1189        let mut inner = self.inner.lock();
1190        inner.remove_reservation(old_owner_object_id, amount);
1191        inner.add_reservation(None, amount);
1192    }
1193
1194    /// Creates a lazy inspect node named `str` under `parent` which will yield statistics for the
1195    /// allocator when queried.
1196    pub fn track_statistics(self: &Arc<Self>, parent: &fuchsia_inspect::Node, name: &str) {
1197        let this = Arc::downgrade(self);
1198        parent.record_lazy_child(name, move || {
1199            let this_clone = this.clone();
1200            async move {
1201                let inspector = fuchsia_inspect::Inspector::default();
1202                if let Some(this) = this_clone.upgrade() {
1203                    let counters = this.counters.lock();
1204                    let root = inspector.root();
1205                    root.record_uint("max_extent_size_bytes", this.max_extent_size_bytes);
1206                    root.record_uint("bytes_total", this.device_size);
1207                    let (allocated, reserved, used, unavailable) = {
1208                        // TODO(https://fxbug.dev/42069513): Push-back or rate-limit to prevent DoS.
1209                        let inner = this.inner.lock();
1210                        (
1211                            inner.allocated_bytes().0,
1212                            inner.reserved_bytes(),
1213                            inner.used_bytes().0,
1214                            inner.unavailable_bytes().0,
1215                        )
1216                    };
1217                    root.record_uint("bytes_allocated", allocated);
1218                    root.record_uint("bytes_reserved", reserved);
1219                    root.record_uint("bytes_used", used);
1220                    root.record_uint("bytes_unavailable", unavailable);
1221
1222                    // TODO(https://fxbug.dev/42068224): Post-compute rather than manually computing
1223                    // metrics.
1224                    if let Some(x) = round_div(100 * allocated, this.device_size) {
1225                        root.record_uint("bytes_allocated_percent", x);
1226                    }
1227                    if let Some(x) = round_div(100 * reserved, this.device_size) {
1228                        root.record_uint("bytes_reserved_percent", x);
1229                    }
1230                    if let Some(x) = round_div(100 * used, this.device_size) {
1231                        root.record_uint("bytes_used_percent", x);
1232                    }
1233                    if let Some(x) = round_div(100 * unavailable, this.device_size) {
1234                        root.record_uint("bytes_unavailable_percent", x);
1235                    }
1236
1237                    root.record_uint("num_flushes", counters.num_flushes);
1238                    if let Some(last_flush_time) = counters.last_flush_time.as_ref() {
1239                        root.record_uint(
1240                            "last_flush_time_ms",
1241                            last_flush_time
1242                                .duration_since(std::time::UNIX_EPOCH)
1243                                .unwrap_or(std::time::Duration::ZERO)
1244                                .as_millis()
1245                                .try_into()
1246                                .unwrap_or(0u64),
1247                        );
1248                    }
1249
1250                    let data = this.allocation_size_histogram();
1251                    let alloc_sizes = root.create_uint_linear_histogram(
1252                        "allocation_size_histogram",
1253                        fuchsia_inspect::LinearHistogramParams {
1254                            floor: 1,
1255                            step_size: 1,
1256                            buckets: 64,
1257                        },
1258                    );
1259                    for (i, count) in data.iter().enumerate() {
1260                        if i != 0 {
1261                            alloc_sizes.insert_multiple(i as u64, *count as usize);
1262                        }
1263                    }
1264                    root.record(alloc_sizes);
1265
1266                    let data = this.inner.lock().rebuild_strategy_trigger_histogram;
1267                    let triggers = root.create_uint_linear_histogram(
1268                        "rebuild_strategy_triggers",
1269                        fuchsia_inspect::LinearHistogramParams {
1270                            floor: 1,
1271                            step_size: 1,
1272                            buckets: 64,
1273                        },
1274                    );
1275                    for (i, count) in data.iter().enumerate() {
1276                        if i != 0 {
1277                            triggers.insert_multiple(i as u64, *count as usize);
1278                        }
1279                    }
1280                    root.record(triggers);
1281
1282                    let allocator_ref = this.clone();
1283                    root.record_child("lsm_tree", move |node| {
1284                        allocator_ref.tree.record_inspect_data(node);
1285                    });
1286                }
1287                Ok(inspector)
1288            }
1289            .boxed()
1290        });
1291    }
1292
1293    /// Returns the offset of the first byte which has not been used by the allocator since its
1294    /// creation.
1295    /// NB: This does *not* take into account existing allocations.  This is only reliable when the
1296    /// allocator was created from scratch, without any pre-existing allocations.
1297    pub fn maximum_offset(&self) -> u64 {
1298        self.maximum_offset.load(Ordering::Relaxed)
1299    }
1300}
1301
1302impl Drop for Allocator {
1303    fn drop(&mut self) {
1304        let inner = self.inner.lock();
1305        // Uncommitted and reserved should be released back using RAII, so they should be zero.
1306        assert_eq!(inner.uncommitted_allocated_bytes(), 0);
1307        assert_eq!(inner.reserved_bytes(), 0);
1308    }
1309}
1310
1311#[fxfs_trace::trace]
1312impl Allocator {
1313    /// Returns the object ID for the allocator.
1314    pub fn object_id(&self) -> u64 {
1315        self.object_id
1316    }
1317
1318    /// Returns information about the allocator such as the layer files storing persisted
1319    /// allocations.
1320    pub fn info(&self) -> AllocatorInfo {
1321        self.inner.lock().info.clone()
1322    }
1323
1324    fn reservation_description(&self, reservation: &Reservation) -> String {
1325        let metadata_info = self.filesystem.upgrade().and_then(|fs| {
1326            let om = fs.object_manager();
1327            if om.is_metadata_reservation(reservation) {
1328                Some((om.borrowed_metadata_space(), om.max_store_reservation()))
1329            } else {
1330                None
1331            }
1332        });
1333        match metadata_info {
1334            Some((borrowed, max_store_reservation)) => {
1335                format!(
1336                    "reservation with {} bytes available (borrowed metadata space: {}, max store \
1337                    reservation: {})",
1338                    reservation.amount(),
1339                    borrowed,
1340                    max_store_reservation
1341                )
1342            }
1343            None => {
1344                format!("reservation with {} bytes available", reservation.amount())
1345            }
1346        }
1347    }
1348
1349    /// Tries to allocate enough space for |object_range| in the specified object and returns the
1350    /// device range allocated.
1351    /// The allocated range may be short (e.g. due to fragmentation), in which case the caller can
1352    /// simply call allocate again until they have enough blocks.
1353    ///
1354    /// We also store the object store ID of the store that the allocation should be assigned to so
1355    /// that we have a means to delete encrypted stores without needing the encryption key.
1356    #[trace]
1357    pub async fn allocate(
1358        self: &Arc<Self>,
1359        transaction: &mut Transaction<'_>,
1360        owner_object_id: u64,
1361        mut len: u64,
1362    ) -> Result<Range<u64>, Error> {
1363        ensure!(self.allocations_allowed.load(Ordering::SeqCst), FxfsError::Unavailable);
1364        assert!(self.block_size.is_aligned(len));
1365        len = std::cmp::min(len, self.max_extent_size_bytes);
1366        debug_assert_ne!(owner_object_id, INVALID_OBJECT_ID);
1367
1368        let requested_len = len;
1369
1370        // Make sure we have space reserved before we try and find the space.
1371        let reservation = if let Some(reservation) = transaction.allocator_reservation() {
1372            match reservation.owner_object_id {
1373                // If there is no owner, this must be a system store that we're allocating for.
1374                None => assert!(self.is_system_store(owner_object_id)),
1375                // If it has an owner, it should not be different than the allocating owner.
1376                Some(res_owner_object_id) => assert_eq!(owner_object_id, res_owner_object_id),
1377            };
1378            // Reservation limits aren't necessarily a multiple of the block size.
1379            let r = reservation
1380                .reserve_with(|limit| std::cmp::min(len, self.block_size.align_down(limit)));
1381            len = r.amount();
1382            Left(r)
1383        } else {
1384            let mut inner = self.inner.lock();
1385            assert!(inner.opened);
1386            // Do not exceed the limit for the owner or the device.
1387            let device_used = inner.used_bytes();
1388            let owner_bytes_left = inner.owner_id_bytes_left(owner_object_id);
1389            // We must take care not to use up space that might be reserved.
1390            let limit =
1391                std::cmp::min(owner_bytes_left, (Saturating(self.device_size) - device_used).0);
1392            len = self.block_size.align_down(std::cmp::min(len, limit));
1393            let owner_entry = inner.owner_bytes.entry(owner_object_id).or_default();
1394            owner_entry.reserved_bytes += len;
1395            Right(ReservationImpl::<_, Self>::new(&**self, Some(owner_object_id), len))
1396        };
1397
1398        if len == 0 {
1399            if let Some(reservation) = transaction.allocator_reservation() {
1400                bail!(anyhow!(FxfsError::NoSpace).context(format!(
1401                    "Failed to allocate {} bytes for owner {} from {}",
1402                    requested_len,
1403                    owner_object_id,
1404                    self.reservation_description(reservation),
1405                )));
1406            } else {
1407                let inner = self.inner.lock();
1408                bail!(anyhow!(FxfsError::NoSpace).context(format!(
1409                    "Failed to allocate {} bytes for owner {} (owner_bytes_left: {}, device_used: \
1410                    {}, device_size: {})",
1411                    requested_len,
1412                    owner_object_id,
1413                    inner.owner_id_bytes_left(owner_object_id),
1414                    inner.used_bytes(),
1415                    self.device_size
1416                )));
1417            }
1418        }
1419
1420        // If volumes have been deleted, flush the device so that we can use any of the freed space.
1421        let volumes_deleted = {
1422            let inner = self.inner.lock();
1423            (!inner.volumes_deleted_pending_sync.is_empty())
1424                .then(|| inner.volumes_deleted_pending_sync.clone())
1425        };
1426
1427        if let Some(volumes_deleted) = volumes_deleted {
1428            // No locks are held here, so in theory, there could be unnecessary syncs, but it
1429            // should be sufficiently rare that it won't matter.
1430            self.filesystem
1431                .upgrade()
1432                .unwrap()
1433                .sync(SyncOptions {
1434                    flush_device: true,
1435                    precondition: Some(Box::new(|| {
1436                        !self.inner.lock().volumes_deleted_pending_sync.is_empty()
1437                    })),
1438                    ..Default::default()
1439                })
1440                .await?;
1441
1442            {
1443                let mut inner = self.inner.lock();
1444                for owner_id in volumes_deleted {
1445                    inner.volumes_deleted_pending_sync.remove(&owner_id);
1446                    inner.marked_for_deletion.insert(owner_id);
1447                }
1448            }
1449
1450            let _guard = self.allocation_mutex.lock().await;
1451            self.rebuild_strategy().await?;
1452        }
1453
1454        #[allow(clippy::never_loop)] // Loop used as a for {} else {}.
1455        let _guard = 'sync: loop {
1456            // Cap number of sync attempts before giving up on finding free space.
1457            for _ in 0..10 {
1458                {
1459                    let guard = self.allocation_mutex.lock().await;
1460
1461                    if !self.needs_sync() {
1462                        break 'sync guard;
1463                    }
1464                }
1465
1466                // All the free space is currently tied up with deallocations, so we need to sync
1467                // and flush the device to free that up.
1468                //
1469                // We can't hold the allocation lock whilst we sync here because the allocation lock
1470                // is also taken in apply_mutations, which is called when journal locks are held,
1471                // and we call sync here which takes those same locks, so it would have the
1472                // potential to result in a deadlock.  Sync holds its own lock to guard against
1473                // multiple syncs occurring at the same time, and we can supply a precondition that
1474                // is evaluated under that lock to ensure we don't sync twice if we don't need to.
1475                self.filesystem
1476                    .upgrade()
1477                    .unwrap()
1478                    .sync(SyncOptions {
1479                        flush_device: true,
1480                        precondition: Some(Box::new(|| self.needs_sync())),
1481                        ..Default::default()
1482                    })
1483                    .await?;
1484            }
1485            bail!(
1486                anyhow!(FxfsError::NoSpace).context("Sync failed to yield sufficient free space.")
1487            );
1488        };
1489
1490        let mut trim_listener = None;
1491        {
1492            let mut inner = self.inner.lock();
1493            inner.allocation_size_histogram[std::cmp::min(63, len / self.block_size) as usize] += 1;
1494
1495            // If trimming would be the reason that this allocation gets cut short, wait for
1496            // trimming to complete before proceeding.
1497            let avail = self
1498                .device_size
1499                .checked_sub(inner.unavailable_bytes().0)
1500                .ok_or(FxfsError::Inconsistent)?;
1501            let free_and_not_being_trimmed =
1502                inner.bytes_available_not_being_trimmed(self.device_size)?;
1503            if free_and_not_being_trimmed < std::cmp::min(len, avail) {
1504                debug_assert!(inner.trim_reserved_bytes > 0);
1505                trim_listener = std::mem::take(&mut inner.trim_listener);
1506            }
1507        }
1508
1509        if let Some(listener) = trim_listener {
1510            listener.await;
1511        }
1512
1513        let result = loop {
1514            {
1515                let mut inner = self.inner.lock();
1516
1517                // While we know rebuild_strategy and take_for_trim are not running (we hold guard),
1518                // apply temporary_allocation removals.
1519                for device_range in inner.dropped_temporary_allocations.drain(..) {
1520                    self.temporary_allocations
1521                        .erase(&AllocatorKey { device_range: Extent(device_range) });
1522                }
1523
1524                match inner.strategy.allocate(len) {
1525                    Err(FxfsError::NotFound) => {
1526                        // Overflow. Fall through and rebuild
1527                        inner.rebuild_strategy_trigger_histogram
1528                            [std::cmp::min(63, (len / self.block_size) as usize)] += 1;
1529                    }
1530                    Err(err) => {
1531                        error!(err:%; "Likely filesystem corruption.");
1532                        return Err(err.into());
1533                    }
1534                    Ok(x) => {
1535                        break x;
1536                    }
1537                }
1538            }
1539            // We've run out of extents of the requested length in RAM but there
1540            // exists more of this size on device. Rescan device and circle back.
1541            // We already hold the allocation_mutex, so exclusive access is guaranteed.
1542            if !self.rebuild_strategy().await? {
1543                error!("Cannot find additional free space. Corruption?");
1544                return Err(FxfsError::Inconsistent.into());
1545            }
1546        };
1547
1548        debug!(device_range:? = result; "allocate");
1549
1550        let len = result.length().unwrap();
1551        let reservation_owner = reservation.either(
1552            // Left means we got an outside reservation.
1553            |l| {
1554                l.forget_some(len);
1555                l.owner_object_id()
1556            },
1557            |r| {
1558                r.forget_some(len);
1559                r.owner_object_id()
1560            },
1561        );
1562
1563        {
1564            let mut inner = self.inner.lock();
1565            let owner_entry = inner.owner_bytes.entry(owner_object_id).or_default();
1566            owner_entry.uncommitted_allocated_bytes += len;
1567            // If the reservation has an owner, ensure they are the same.
1568            assert_eq!(owner_object_id, reservation_owner.unwrap_or(owner_object_id));
1569            inner.remove_reservation(reservation_owner, len);
1570            self.temporary_allocations.insert(AllocatorItem {
1571                key: AllocatorKey { device_range: Extent(result.clone()) },
1572                value: AllocatorValue::Abs { owner_object_id, count: 1 },
1573            })?;
1574        }
1575
1576        let mutation =
1577            AllocatorMutation::Allocate { device_range: result.clone().into(), owner_object_id };
1578        assert!(transaction.add(self.object_id(), Mutation::Allocator(mutation)).is_none());
1579
1580        Ok(result)
1581    }
1582
1583    /// Marks the given device range as allocated.  The main use case for this at this time is for
1584    /// the super-block which needs to be at a fixed location on the device.
1585    #[trace]
1586    pub fn mark_allocated(
1587        &self,
1588        transaction: &mut Transaction<'_>,
1589        owner_object_id: u64,
1590        device_range: Range<u64>,
1591    ) -> Result<(), Error> {
1592        debug_assert_ne!(owner_object_id, INVALID_OBJECT_ID);
1593        {
1594            let len = device_range.length().map_err(|_| FxfsError::InvalidArgs)?;
1595
1596            let mut inner = self.inner.lock();
1597            let device_used = inner.used_bytes();
1598            let owner_id_bytes_left = inner.owner_id_bytes_left(owner_object_id);
1599            let owner_entry = inner.owner_bytes.entry(owner_object_id).or_default();
1600            ensure!(
1601                device_range.end <= self.device_size
1602                    && (Saturating(self.device_size) - device_used).0 >= len
1603                    && owner_id_bytes_left >= len,
1604                anyhow!(FxfsError::NoSpace).context(format!(
1605                    "Failed to allocate {} bytes at {:?} for owner {} (owner_bytes_left: {}, \
1606                    device_used: {}, device_size: {})",
1607                    len,
1608                    device_range,
1609                    owner_object_id,
1610                    owner_id_bytes_left,
1611                    device_used,
1612                    self.device_size
1613                ))
1614            );
1615            if let Some(reservation) = transaction.allocator_reservation() {
1616                // The transaction takes ownership of this hold.
1617                reservation
1618                    .reserve(len)
1619                    .ok_or_else(|| {
1620                        anyhow!(FxfsError::NoSpace).context(format!(
1621                            "Failed to reserve {} bytes at {:?} from {}",
1622                            len,
1623                            device_range,
1624                            self.reservation_description(reservation),
1625                        ))
1626                    })?
1627                    .forget();
1628            }
1629            owner_entry.uncommitted_allocated_bytes += len;
1630            inner.strategy.remove(device_range.clone());
1631            self.temporary_allocations.insert(AllocatorItem {
1632                key: AllocatorKey { device_range: Extent(device_range.clone()) },
1633                value: AllocatorValue::Abs { owner_object_id, count: 1 },
1634            })?;
1635        }
1636        let mutation =
1637            AllocatorMutation::Allocate { device_range: device_range.into(), owner_object_id };
1638        transaction.add(self.object_id(), Mutation::Allocator(mutation));
1639        Ok(())
1640    }
1641
1642    /// Sets the limits for an owner object in terms of usage.
1643    pub fn set_bytes_limit(
1644        &self,
1645        transaction: &mut Transaction<'_>,
1646        owner_object_id: u64,
1647        bytes: u64,
1648    ) -> Result<(), Error> {
1649        // System stores cannot be given limits.
1650        assert!(!self.is_system_store(owner_object_id));
1651        transaction.add(
1652            self.object_id(),
1653            Mutation::Allocator(AllocatorMutation::SetLimit { owner_object_id, bytes }),
1654        );
1655        Ok(())
1656    }
1657
1658    /// Gets the bytes limit for an owner object.
1659    pub fn get_owner_bytes_limit(&self, owner_object_id: u64) -> Option<u64> {
1660        self.inner.lock().info.limit_bytes.get(&owner_object_id).copied()
1661    }
1662
1663    /// Gets the number of bytes currently used by allocations, reservations or uncommitted
1664    /// allocations.
1665    pub fn get_owner_bytes_used(&self, owner_object_id: u64) -> u64 {
1666        self.inner.lock().owner_bytes.get(&owner_object_id).map_or(0, |info| info.used_bytes().0)
1667    }
1668
1669    /// Deallocates the given device range for the specified object.
1670    #[trace]
1671    pub async fn deallocate(
1672        &self,
1673        transaction: &mut Transaction<'_>,
1674        owner_object_id: u64,
1675        dealloc_range: Range<u64>,
1676    ) -> Result<u64, Error> {
1677        debug!(device_range:? = dealloc_range; "deallocate");
1678        ensure!(dealloc_range.is_valid(), FxfsError::InvalidArgs);
1679        // We don't currently support sharing of allocations (value.count always equals 1), so as
1680        // long as we can assume the deallocated range is actually allocated, we can avoid device
1681        // access.
1682        let deallocated = dealloc_range.end - dealloc_range.start;
1683        let mutation = AllocatorMutation::Deallocate {
1684            device_range: dealloc_range.clone().into(),
1685            owner_object_id,
1686        };
1687        transaction.add(self.object_id(), Mutation::Allocator(mutation));
1688
1689        let _guard = self.allocation_mutex.lock().await;
1690
1691        // We use `dropped_temporary_allocations` to defer removals from `temporary_allocations` in
1692        // places where we can't execute async code or take locks.
1693        //
1694        // It's important we don't ever remove entries from `temporary_allocations` without
1695        // holding the `allocation_mutex` lock or else we may end up with an inconsistent view of
1696        // available disk space when we combine temporary_allocations with the LSMTree.
1697        // This is normally done in `allocate()` but we also need to apply these here because
1698        // `temporary_allocations` is also used to track deallocated space until it has been
1699        // flushed (see comment below). A user may allocate and then deallocate space before calling
1700        // allocate() a second time, so if we do not clean up here, we may end up with the same
1701        // range in temporary_allocations twice (once for allocate, once for deallocate).
1702        let mut inner = self.inner.lock();
1703        for device_range in inner.dropped_temporary_allocations.drain(..) {
1704            self.temporary_allocations.erase(&AllocatorKey { device_range: Extent(device_range) });
1705        }
1706
1707        // We can't reuse deallocated space immediately because failure to successfully flush will
1708        // mean that on next mount, we may find this space is still assigned to the deallocated
1709        // region. To avoid immediate reuse, we hold these regions in 'temporary_allocations' until
1710        // after a successful flush so we know the region is safe to reuse.
1711        self.temporary_allocations
1712            .insert(AllocatorItem {
1713                key: AllocatorKey { device_range: Extent(dealloc_range.clone()) },
1714                value: AllocatorValue::Abs { owner_object_id, count: 1 },
1715            })
1716            .context("tracking deallocated")?;
1717
1718        Ok(deallocated)
1719    }
1720
1721    /// Marks allocations associated with a given |owner_object_id| for deletion.
1722    ///
1723    /// This is used as part of deleting encrypted volumes (ObjectStore) without having the keys.
1724    ///
1725    /// MarkForDeletion mutations eventually manipulates allocator metadata (AllocatorInfo) instead
1726    /// of the mutable layer but we must be careful not to do this too early and risk premature
1727    /// reuse of extents.
1728    ///
1729    /// Replay is not guaranteed until the *device* gets flushed, so we cannot reuse the deleted
1730    /// extents until we've flushed the device.
1731    ///
1732    /// TODO(b/316827348): Consider removing the use of mark_for_deletion in AllocatorInfo and
1733    /// just compacting?
1734    ///
1735    /// After an allocator.flush() (i.e. a major compaction), we know that there is no data left
1736    /// in the layer files for this owner_object_id and we are able to clear `marked_for_deletion`.
1737    pub fn mark_for_deletion(&self, transaction: &mut Transaction<'_>, owner_object_id: u64) {
1738        // Note that because the actual time of deletion (the next major compaction) is undefined,
1739        // |owner_object_id| should not be reused after this call.
1740        transaction.add(
1741            self.object_id(),
1742            Mutation::Allocator(AllocatorMutation::MarkForDeletion(owner_object_id)),
1743        );
1744    }
1745
1746    /// Called when the device has been flush and indicates what the journal log offset was when
1747    /// that happened.
1748    pub fn did_flush_device(&self, flush_log_offset: u64) {
1749        // First take out the deallocations that we now know to be flushed.  The list is maintained
1750        // in order, so we can stop on the first entry that we find that should not be unreserved
1751        // yet.
1752        #[allow(clippy::never_loop)] // Loop used as a for {} else {}.
1753        let deallocs = 'deallocs_outer: loop {
1754            let mut inner = self.inner.lock();
1755            for (index, dealloc) in inner.committed_deallocated.iter().enumerate() {
1756                if dealloc.log_file_offset >= flush_log_offset {
1757                    let mut deallocs = inner.committed_deallocated.split_off(index);
1758                    // Swap because we want the opposite of what split_off does.
1759                    std::mem::swap(&mut inner.committed_deallocated, &mut deallocs);
1760                    break 'deallocs_outer deallocs;
1761                }
1762            }
1763            break std::mem::take(&mut inner.committed_deallocated);
1764        };
1765
1766        // Now we can free those elements.
1767        let mut inner = self.inner.lock();
1768        let mut totals = BTreeMap::<u64, u64>::new();
1769        for dealloc in deallocs {
1770            *(totals.entry(dealloc.owner_object_id).or_default()) +=
1771                dealloc.range.length().unwrap();
1772            inner.strategy.free(dealloc.range.clone()).expect("dealloced ranges");
1773            self.temporary_allocations.erase(&AllocatorKey { device_range: Extent(dealloc.range) });
1774        }
1775
1776        // This *must* come after we've removed the records from reserved reservations because the
1777        // allocator uses this value to decide whether or not a device-flush is required and it must
1778        // be possible to find free space if it thinks no device-flush is required.
1779        for (owner_object_id, total) in totals {
1780            match inner.owner_bytes.get_mut(&owner_object_id) {
1781                Some(counters) => counters.committed_deallocated_bytes -= total,
1782                None => {
1783                    // This should only be possible if the volume has been deleted and a sync is
1784                    // pending.
1785                    assert!(inner.volumes_deleted_pending_sync.contains(&owner_object_id));
1786                }
1787            }
1788        }
1789    }
1790
1791    /// Returns a reservation that can be used later, or None if there is insufficient space. The
1792    /// |owner_object_id| indicates which object in the root object store the reservation is for.
1793    pub fn reserve(
1794        self: Arc<Self>,
1795        owner_object_id: Option<u64>,
1796        amount: u64,
1797    ) -> Option<Reservation> {
1798        {
1799            let mut inner = self.inner.lock();
1800
1801            let device_free = (Saturating(self.device_size) - inner.used_bytes()).0;
1802
1803            let limit = match owner_object_id {
1804                Some(id) => std::cmp::min(inner.owner_id_bytes_left(id), device_free),
1805                None => device_free,
1806            };
1807            if limit < amount {
1808                return None;
1809            }
1810            inner.add_reservation(owner_object_id, amount);
1811        }
1812        Some(Reservation::new(self, owner_object_id, amount))
1813    }
1814
1815    /// Like reserve, but takes a callback is passed the |limit| and should return the amount,
1816    /// which can be zero.
1817    pub fn reserve_with(
1818        self: Arc<Self>,
1819        owner_object_id: Option<u64>,
1820        amount: impl FnOnce(u64) -> u64,
1821    ) -> Reservation {
1822        let amount = {
1823            let mut inner = self.inner.lock();
1824
1825            let device_free = (Saturating(self.device_size) - inner.used_bytes()).0;
1826
1827            let amount = amount(match owner_object_id {
1828                Some(id) => std::cmp::min(inner.owner_id_bytes_left(id), device_free),
1829                None => device_free,
1830            });
1831
1832            inner.add_reservation(owner_object_id, amount);
1833
1834            amount
1835        };
1836
1837        Reservation::new(self, owner_object_id, amount)
1838    }
1839
1840    /// Returns the total number of allocated bytes.
1841    pub fn get_allocated_bytes(&self) -> u64 {
1842        self.inner.lock().allocated_bytes().0
1843    }
1844
1845    /// Returns the size of bytes available to allocate.
1846    pub fn get_disk_bytes(&self) -> u64 {
1847        self.device_size
1848    }
1849
1850    /// Returns the total number of allocated bytes per owner_object_id.
1851    /// Note that this is quite an expensive operation as it copies the collection.
1852    /// This is intended for use in fsck() and friends, not general use code.
1853    pub fn get_owner_allocated_bytes(&self) -> BTreeMap<u64, u64> {
1854        self.inner.lock().owner_bytes.iter().map(|(k, v)| (*k, v.allocated_bytes.0)).collect()
1855    }
1856
1857    /// Returns the number of allocated and reserved bytes.
1858    pub fn get_used_bytes(&self) -> Saturating<u64> {
1859        let inner = self.inner.lock();
1860        inner.used_bytes()
1861    }
1862
1863    pub async fn flush(&self) -> Result<Version, Error> {
1864        let filesystem = self.filesystem.upgrade().unwrap();
1865        let object_manager = filesystem.object_manager();
1866        let earliest_version = self.tree.get_earliest_version();
1867        if !object_manager.needs_flush(self.object_id()) && earliest_version == LATEST_VERSION {
1868            // Early exit, but still return the earliest version used by a struct in the tree
1869            return Ok(earliest_version);
1870        }
1871
1872        let fs = self.filesystem.upgrade().unwrap();
1873        let mut flusher = Flusher::new(self, &fs).await;
1874        let (new_layer_file, info) = flusher.start().await?;
1875        flusher.finish(new_layer_file, info).await
1876    }
1877}
1878
1879impl ReservationOwner for Allocator {
1880    fn release_reservation(&self, owner_object_id: Option<u64>, amount: u64) {
1881        self.inner.lock().remove_reservation(owner_object_id, amount);
1882    }
1883}
1884
1885#[async_trait]
1886impl JournalingObject for Allocator {
1887    fn apply_mutation(
1888        &self,
1889        mutation: Mutation,
1890        context: &ApplyContext<'_, '_>,
1891        _assoc_obj: AssocObj<'_>,
1892    ) -> Result<(), Error> {
1893        match mutation {
1894            Mutation::Allocator(AllocatorMutation::MarkForDeletion(owner_object_id)) => {
1895                let mut inner = self.inner.lock();
1896                inner.owner_bytes.remove(&owner_object_id);
1897
1898                // We use `info.marked_for_deletion` to track the committed state and
1899                // `inner.marked_for_deletion` to track volumes marked for deletion *after* we have
1900                // flushed the device.  It is not safe to use extents belonging to deleted volumes
1901                // until after we have flushed the device.
1902                inner.info.marked_for_deletion.insert(owner_object_id);
1903                inner.volumes_deleted_pending_sync.insert(owner_object_id);
1904
1905                inner.info.limit_bytes.remove(&owner_object_id);
1906            }
1907            Mutation::Allocator(AllocatorMutation::Allocate { device_range, owner_object_id }) => {
1908                self.maximum_offset.fetch_max(device_range.end, Ordering::Relaxed);
1909                let item = AllocatorItem {
1910                    key: AllocatorKey { device_range: Extent(device_range.0.clone()) },
1911                    value: AllocatorValue::Abs { count: 1, owner_object_id },
1912                };
1913                let len = item.key.device_range.length().unwrap();
1914                let lower_bound = item.key.lower_bound_for_merge_into();
1915                self.tree.merge_into(item, &lower_bound);
1916                let mut inner = self.inner.lock();
1917                let entry = inner.owner_bytes.entry(owner_object_id).or_default();
1918                entry.allocated_bytes += len;
1919                if let ApplyMode::Live(transaction) = context.mode {
1920                    entry.uncommitted_allocated_bytes -= len;
1921                    // Note that we cannot drop entries from temporary_allocations without holding
1922                    // the allocation_mutex as it may introduce races. We instead add the range to
1923                    // a Vec that can be applied later when we hold the lock (See comment on
1924                    // `dropped_temporary_allocations` above).
1925                    inner.dropped_temporary_allocations.push(device_range.0);
1926                    if let Some(reservation) = transaction.allocator_reservation() {
1927                        reservation.commit(len);
1928                    }
1929                }
1930            }
1931            Mutation::Allocator(AllocatorMutation::Deallocate {
1932                device_range,
1933                owner_object_id,
1934            }) => {
1935                let item = AllocatorItem {
1936                    key: AllocatorKey { device_range: Extent(device_range.0) },
1937                    value: AllocatorValue::None,
1938                };
1939                let len = item.key.device_range.length().unwrap();
1940
1941                {
1942                    let mut inner = self.inner.lock();
1943                    {
1944                        let entry = inner.owner_bytes.entry(owner_object_id).or_default();
1945                        entry.allocated_bytes -= len;
1946                        if context.mode.is_live() {
1947                            entry.committed_deallocated_bytes += len;
1948                        }
1949                    }
1950                    if context.mode.is_live() {
1951                        inner.committed_deallocated.push_back(CommittedDeallocation {
1952                            log_file_offset: context.checkpoint.file_offset,
1953                            range: (*item.key.device_range).clone(),
1954                            owner_object_id,
1955                        });
1956                    }
1957                    if let ApplyMode::Live(transaction) = context.mode
1958                        && let Some(reservation) = transaction.allocator_reservation()
1959                    {
1960                        inner.add_reservation(reservation.owner_object_id(), len);
1961                        reservation.add(len);
1962                    }
1963                }
1964                let lower_bound = item.key.lower_bound_for_merge_into();
1965                self.tree.merge_into(item, &lower_bound);
1966            }
1967            Mutation::Allocator(AllocatorMutation::SetLimit { owner_object_id, bytes }) => {
1968                // Journal replay is ordered and each of these calls is idempotent. So the last one
1969                // will be respected, it doesn't matter if the value is already set, or gets changed
1970                // multiple times during replay. When it gets opened it will be merged in with the
1971                // snapshot.
1972                self.inner.lock().info.limit_bytes.insert(owner_object_id, bytes);
1973            }
1974            Mutation::BeginFlush => {
1975                self.tree.seal();
1976                // Transfer our running count for allocated_bytes so that it gets written to the new
1977                // info file when flush completes.
1978                let mut inner = self.inner.lock();
1979                let allocated_bytes =
1980                    inner.owner_bytes.iter().map(|(k, v)| (*k, v.allocated_bytes.0)).collect();
1981                inner.info.allocated_bytes = allocated_bytes;
1982            }
1983            Mutation::EndFlush => {}
1984            _ => bail!("unexpected mutation: {:?}", mutation),
1985        }
1986        Ok(())
1987    }
1988
1989    fn drop_mutation(&self, mutation: Mutation, transaction: &Transaction<'_>) {
1990        match mutation {
1991            Mutation::Allocator(AllocatorMutation::Allocate { device_range, owner_object_id }) => {
1992                let len = device_range.length().unwrap();
1993                let mut inner = self.inner.lock();
1994                inner
1995                    .owner_bytes
1996                    .entry(owner_object_id)
1997                    .or_default()
1998                    .uncommitted_allocated_bytes -= len;
1999                if let Some(reservation) = transaction.allocator_reservation() {
2000                    let res_owner = reservation.owner_object_id();
2001                    inner.add_reservation(res_owner, len);
2002                    reservation.release_reservation(res_owner, len);
2003                }
2004                inner.strategy.free(device_range.0.clone()).expect("drop mutaton");
2005                self.temporary_allocations
2006                    .erase(&AllocatorKey { device_range: Extent(device_range.0) });
2007            }
2008            Mutation::Allocator(AllocatorMutation::Deallocate { device_range, .. }) => {
2009                self.temporary_allocations
2010                    .erase(&AllocatorKey { device_range: Extent(device_range.0) });
2011            }
2012            _ => {}
2013        }
2014    }
2015
2016    async fn flush(&self, _reason: crate::filesystem::FlushReason) -> Result<Version, Error> {
2017        self.flush().await
2018    }
2019}
2020
2021// The merger is unable to merge extents that exist like the following:
2022//
2023//     |----- +1 -----|
2024//                    |----- -1 -----|
2025//                    |----- +2 -----|
2026//
2027// It cannot coalesce them because it has to emit the +1 record so that it can move on and merge the
2028// -1 and +2 records. To address this, we add another stage that applies after merging which
2029// coalesces records after they have been emitted.  This is a bit simpler than merging because the
2030// records cannot overlap, so it's just a question of merging adjacent records if they happen to
2031// have the same delta and object_id.
2032
2033pub struct CoalescingIterator<I> {
2034    iter: I,
2035    item: Option<AllocatorItem>,
2036}
2037
2038impl<I: LayerIterator<AllocatorKey, AllocatorValue>> CoalescingIterator<I> {
2039    pub async fn new(iter: I) -> Result<CoalescingIterator<I>, Error> {
2040        let mut iter = Self { iter, item: None };
2041        iter.advance().await?;
2042        Ok(iter)
2043    }
2044}
2045
2046impl<I: LayerIterator<AllocatorKey, AllocatorValue>> LayerIterator<AllocatorKey, AllocatorValue>
2047    for CoalescingIterator<I>
2048{
2049    async fn advance(&mut self) -> Result<(), Error> {
2050        self.item = self.iter.get().map(|x| x.cloned());
2051        if self.item.is_none() {
2052            return Ok(());
2053        }
2054        let left = self.item.as_mut().unwrap();
2055        loop {
2056            self.iter.advance().await?;
2057            match self.iter.get() {
2058                None => return Ok(()),
2059                Some(right) => {
2060                    // The two records cannot overlap.
2061                    ensure!(
2062                        left.key.device_range.end <= right.key.device_range.start,
2063                        FxfsError::Inconsistent
2064                    );
2065                    // We can only coalesce records if they are touching and have the same value.
2066                    if left.key.device_range.end < right.key.device_range.start
2067                        || left.value != *right.value
2068                    {
2069                        return Ok(());
2070                    }
2071                    left.key.device_range.end = right.key.device_range.end;
2072                }
2073            }
2074        }
2075    }
2076
2077    fn advance_dyn<'a>(&'a mut self) -> Result<Option<BoxFuture<'a, Result<(), Error>>>, Error> {
2078        Ok(Some(Box::pin(self.advance())))
2079    }
2080
2081    fn get(&self) -> Option<ItemRef<'_, AllocatorKey, AllocatorValue>> {
2082        self.item.as_ref().map(|x| x.as_item_ref())
2083    }
2084}
2085
2086struct Flusher<'a> {
2087    allocator: &'a Allocator,
2088    fs: &'a Arc<FxFilesystem>,
2089    _guard: WriteGuard<'a>,
2090}
2091
2092impl<'a> Flusher<'a> {
2093    async fn new(allocator: &'a Allocator, fs: &'a Arc<FxFilesystem>) -> Self {
2094        let keys = lock_keys![LockKey::flush(allocator.object_id())];
2095        Self { allocator, fs, _guard: fs.lock_manager().write_lock(keys).await }
2096    }
2097
2098    fn txn_options() -> Options<'static> {
2099        Options {
2100            skip_journal_checks: true,
2101            reservation: ReservationOptions::BorrowedMetadataAndData,
2102            ..Default::default()
2103        }
2104    }
2105
2106    async fn start(&mut self) -> Result<(DataObjectHandle<ObjectStore>, AllocatorInfo), Error> {
2107        let root_store = self.fs.root_store();
2108        let mut transaction = root_store.new_transaction(lock_keys![], Self::txn_options()).await?;
2109        let layer_object_handle = ObjectStore::create_object(
2110            &root_store,
2111            &mut transaction,
2112            HandleOptions { skip_journal_checks: true, ..Default::default() },
2113            None,
2114        )
2115        .await?;
2116        root_store.add_to_graveyard(&mut transaction, layer_object_handle.object_id());
2117        // It's important that this transaction does not include any allocations because we use
2118        // BeginFlush as a snapshot point for mutations to the tree: other allocator mutations
2119        // within this transaction might get applied before seal (which would be OK), but they could
2120        // equally get applied afterwards (since Transaction makes no guarantees about the order in
2121        // which mutations are applied whilst committing), in which case they'd get lost on replay
2122        // because the journal will only send mutations that follow this transaction.
2123        transaction.add(self.allocator.object_id(), Mutation::BeginFlush);
2124        let info = transaction
2125            .commit_with_callback(|_| {
2126                // We must capture `info` as it is when we start flushing. Subsequent transactions
2127                // can end up modifying `info` and we shouldn't capture those here.
2128                self.allocator.inner.lock().info.clone()
2129            })
2130            .await?;
2131        Ok((layer_object_handle, info))
2132    }
2133
2134    async fn finish(
2135        self,
2136        layer_object_handle: DataObjectHandle<ObjectStore>,
2137        mut info: AllocatorInfo,
2138    ) -> Result<Version, Error> {
2139        let txn_options = Self::txn_options();
2140
2141        let layer_set = self.allocator.tree.immutable_layer_set();
2142        let total_len = layer_set.sum_len();
2143        {
2144            let start_time = std::time::Instant::now();
2145            let merged_layer_count = layer_set.layers.len();
2146            let mut merger = layer_set.merger();
2147            let iter = self.allocator.filter(merger.query(Query::FullScan).await?, true).await?;
2148            let iter = CoalescingIterator::new(iter).await?;
2149            let bytes_written = compact_with_iterator(
2150                iter,
2151                total_len,
2152                DirectWriter::new(&layer_object_handle, txn_options).await,
2153                layer_object_handle.block_size(),
2154                Some(self.fs.journal().get_compaction_yielder()),
2155            )
2156            .await?;
2157
2158            self.allocator.tree.report_compaction_metrics(
2159                bytes_written,
2160                start_time.elapsed(),
2161                merged_layer_count,
2162            );
2163        }
2164
2165        let root_store = self.fs.root_store();
2166
2167        // Both of these forward-declared variables need to outlive the transaction.
2168        let object_handle;
2169        let reservation_update;
2170        let mut transaction = root_store
2171            .new_transaction(
2172                lock_keys![LockKey::object(
2173                    root_store.store_object_id(),
2174                    self.allocator.object_id()
2175                )],
2176                txn_options,
2177            )
2178            .await?;
2179        let mut serialized_info = Vec::new();
2180
2181        debug!(oid = layer_object_handle.object_id(); "new allocator layer file");
2182        object_handle = ObjectStore::open_object(
2183            &root_store,
2184            self.allocator.object_id(),
2185            HandleOptions::default(),
2186            None,
2187        )
2188        .await?;
2189
2190        // Move all the existing layers to the graveyard.
2191        for object_id in &info.layers {
2192            root_store.add_to_graveyard(&mut transaction, *object_id);
2193        }
2194
2195        // Write out updated info.
2196
2197        // After successfully flushing, all the stores that were marked for deletion at the time of
2198        // the BeginFlush transaction, no longer need to be marked for deletion.  There can be
2199        // stores that have been deleted since the BeginFlush transaction, but they will be covered
2200        // by a MarkForDeletion mutation.
2201        let marked_for_deletion = std::mem::take(&mut info.marked_for_deletion);
2202
2203        info.layers = vec![layer_object_handle.object_id()];
2204
2205        info.serialize_with_version(&mut serialized_info)?;
2206
2207        let mut buf = object_handle.allocate_buffer(serialized_info.len()).await;
2208        buf.copy_from_slice(&serialized_info);
2209        object_handle.txn_write(&mut transaction, 0u64, buf.as_ref()).await?;
2210
2211        reservation_update = ReservationUpdate::new(tree::reservation_amount_from_layer_size(
2212            layer_object_handle.get_size(),
2213        ));
2214
2215        // It's important that EndFlush is in the same transaction that we write AllocatorInfo,
2216        // because we use EndFlush to make the required adjustments to allocated_bytes.
2217        transaction.add_with_object(
2218            self.allocator.object_id(),
2219            Mutation::EndFlush,
2220            AssocObj::Borrowed(&reservation_update),
2221        );
2222        root_store.remove_from_graveyard(&mut transaction, layer_object_handle.object_id());
2223
2224        let layers = vec![layer_from_handle(layer_object_handle, None).await?];
2225        transaction
2226            .commit_with_callback(|_| {
2227                self.allocator.tree.set_layers(layers);
2228
2229                // At this point we've committed the new layers to disk so we can start using them.
2230                // This means we can also switch to the new AllocatorInfo which clears
2231                // marked_for_deletion.
2232                let mut inner = self.allocator.inner.lock();
2233                inner.info.layers = info.layers;
2234                for owner_id in marked_for_deletion {
2235                    inner.marked_for_deletion.remove(&owner_id);
2236                    inner.info.marked_for_deletion.remove(&owner_id);
2237                }
2238            })
2239            .await?;
2240
2241        // Now close the layers and purge them.
2242        for layer in layer_set.layers {
2243            let object_id = layer.handle().map(|h| h.object_id());
2244            layer.close_layer().await;
2245            if let Some(object_id) = object_id {
2246                root_store.tombstone_object(object_id, txn_options, None).await?;
2247            }
2248        }
2249
2250        let mut counters = self.allocator.counters.lock();
2251        counters.num_flushes += 1;
2252        counters.last_flush_time = Some(std::time::SystemTime::now());
2253        // Return the earliest version used by a struct in the tree
2254        Ok(self.allocator.tree.get_earliest_version())
2255    }
2256}
2257
2258#[cfg(test)]
2259mod tests {
2260    use crate::filesystem::{FxFilesystem, FxFilesystemBuilder, OpenFxFilesystem};
2261    use crate::fsck::fsck;
2262    use crate::lsm_tree::skip_list_layer::SkipListLayer;
2263    use crate::lsm_tree::types::{FuzzyHash as _, Item, ItemRef, LayerIterator, LayerKey as _};
2264    use crate::lsm_tree::{LSMTree, Query};
2265    use crate::object_handle::ObjectHandle;
2266    use crate::object_store::allocator::merge::merge;
2267    use crate::object_store::allocator::{
2268        Allocator, AllocatorKey, AllocatorValue, CoalescingIterator, EXTENT_HASH_BUCKET_SIZE,
2269    };
2270    use crate::object_store::extent::MIN_BLOCK_SIZE;
2271    use crate::object_store::transaction::{
2272        Options, ReservationOptions, TRANSACTION_METADATA_MAX_AMOUNT, lock_keys,
2273    };
2274    use crate::object_store::volume::root_volume;
2275    use crate::object_store::{Directory, FxfsError, LockKey, NewChildStoreOptions, ObjectStore};
2276    use crate::range::RangeExt;
2277    use crate::serialized_types::{LATEST_VERSION, Versioned};
2278    use crate::testing;
2279    use bincode::Options as _;
2280    use fuchsia_async as fasync;
2281    use fuchsia_sync::Mutex;
2282    use std::cmp::{max, min};
2283    use std::ops::{Bound, Range};
2284    use std::sync::Arc;
2285    use storage_device::DeviceHolder;
2286    use storage_device::fake_device::FakeDevice;
2287
2288    #[test]
2289    fn test_allocator_key_is_range_based() {
2290        // Make sure we disallow using allocator keys with point queries.
2291        assert!(AllocatorKey { device_range: (0..100).into() }.is_range_key());
2292    }
2293
2294    #[test]
2295    fn test_allocator_key_fuzzy_hash_len() {
2296        let key = AllocatorKey { device_range: (0..512).into() };
2297        let mut iter = key.fuzzy_hash();
2298        assert_eq!(iter.len(), 1);
2299        assert_eq!(iter.size_hint(), (1, Some(1)));
2300        assert!(iter.next().is_some());
2301        assert_eq!(iter.len(), 0);
2302        assert_eq!(iter.size_hint(), (0, Some(0)));
2303        assert_eq!(iter.next(), None);
2304
2305        let key = AllocatorKey { device_range: (0..3 * EXTENT_HASH_BUCKET_SIZE).into() };
2306        let mut iter = key.fuzzy_hash();
2307        assert_eq!(iter.len(), 3);
2308        assert_eq!(iter.size_hint(), (3, Some(3)));
2309        assert!(iter.next().is_some());
2310        assert_eq!(iter.len(), 2);
2311        assert_eq!(iter.size_hint(), (2, Some(2)));
2312        assert!(iter.next().is_some());
2313        assert_eq!(iter.len(), 1);
2314        assert_eq!(iter.size_hint(), (1, Some(1)));
2315        assert!(iter.next().is_some());
2316        assert_eq!(iter.len(), 0);
2317        assert_eq!(iter.size_hint(), (0, Some(0)));
2318        assert_eq!(iter.next(), None);
2319    }
2320
2321    #[test]
2322    fn test_allocator_key_search_key() {
2323        let key = AllocatorKey { device_range: (MIN_BLOCK_SIZE.get()..3 * MIN_BLOCK_SIZE).into() };
2324        assert!(!key.is_search_key());
2325        let search_key = key.search_key().unwrap();
2326        assert!(search_key.is_search_key());
2327        assert_eq!(search_key.device_range, (MIN_BLOCK_SIZE.get()..2 * MIN_BLOCK_SIZE).into());
2328    }
2329
2330    #[test]
2331    fn test_allocator_key_serialization_compatibility() {
2332        let range = 1024..2048;
2333        let key = AllocatorKey { device_range: range.clone().into() };
2334
2335        // 1. Serialize the range directly using bincode (which is what the manual
2336        //    Versioned implementation did)
2337        let options = bincode::DefaultOptions::new().allow_trailing_bytes();
2338        let expected_bytes = options.serialize(&range).unwrap();
2339
2340        // 2. Serialize AllocatorKey via our derived Versioned implementation
2341        let mut key_bytes = Vec::new();
2342        key.serialize_into(&mut key_bytes).unwrap();
2343
2344        // Verify they are byte-compatible
2345        assert_eq!(key_bytes, expected_bytes);
2346
2347        // 3. Deserialize a Range<u64> as AllocatorKey via our derived Versioned
2348        //    implementation
2349        let deserialized_key =
2350            AllocatorKey::deserialize_from(&mut &expected_bytes[..], LATEST_VERSION).unwrap();
2351        assert_eq!(deserialized_key, key);
2352    }
2353
2354    #[fuchsia::test]
2355    async fn test_coalescing_iterator() {
2356        let skip_list = SkipListLayer::new(100);
2357        let items = [
2358            Item::new(
2359                AllocatorKey { device_range: (0..100 * 512).into() },
2360                AllocatorValue::Abs { count: 1, owner_object_id: 99 },
2361            ),
2362            Item::new(
2363                AllocatorKey { device_range: (100 * 512..200 * 512).into() },
2364                AllocatorValue::Abs { count: 1, owner_object_id: 99 },
2365            ),
2366        ];
2367        skip_list.insert(items[1].clone()).expect("insert error");
2368        skip_list.insert(items[0].clone()).expect("insert error");
2369        let mut iter =
2370            CoalescingIterator::new(skip_list.seek(Bound::Unbounded)).await.expect("new failed");
2371        let ItemRef { key, value, .. } = iter.get().expect("get failed");
2372        assert_eq!(
2373            (key, value),
2374            (
2375                &AllocatorKey { device_range: (0..200 * 512).into() },
2376                &AllocatorValue::Abs { count: 1, owner_object_id: 99 }
2377            )
2378        );
2379        iter.advance().await.expect("advance failed");
2380        assert!(iter.get().is_none());
2381    }
2382
2383    #[fuchsia::test]
2384    async fn test_merge_and_coalesce_across_three_layers() {
2385        let lsm_tree = LSMTree::new(merge, None);
2386        lsm_tree
2387            .insert(Item::new(
2388                AllocatorKey { device_range: (100 * 512..200 * 512).into() },
2389                AllocatorValue::Abs { count: 1, owner_object_id: 99 },
2390            ))
2391            .expect("insert error");
2392        lsm_tree.seal();
2393        lsm_tree
2394            .insert(Item::new(
2395                AllocatorKey { device_range: (0..100 * 512).into() },
2396                AllocatorValue::Abs { count: 1, owner_object_id: 99 },
2397            ))
2398            .expect("insert error");
2399
2400        let layer_set = lsm_tree.layer_set();
2401        let mut merger = layer_set.merger();
2402        let mut iter =
2403            CoalescingIterator::new(merger.query(Query::FullScan).await.expect("seek failed"))
2404                .await
2405                .expect("new failed");
2406        let ItemRef { key, value, .. } = iter.get().expect("get failed");
2407        assert_eq!(
2408            (key, value),
2409            (
2410                &AllocatorKey { device_range: (0..200 * 512).into() },
2411                &AllocatorValue::Abs { count: 1, owner_object_id: 99 }
2412            )
2413        );
2414        iter.advance().await.expect("advance failed");
2415        assert!(iter.get().is_none());
2416    }
2417
2418    #[fuchsia::test]
2419    async fn test_merge_and_coalesce_wont_merge_across_object_id() {
2420        let lsm_tree = LSMTree::new(merge, None);
2421        lsm_tree
2422            .insert(Item::new(
2423                AllocatorKey { device_range: (100 * 512..200 * 512).into() },
2424                AllocatorValue::Abs { count: 1, owner_object_id: 99 },
2425            ))
2426            .expect("insert error");
2427        lsm_tree.seal();
2428        lsm_tree
2429            .insert(Item::new(
2430                AllocatorKey { device_range: (0..100 * 512).into() },
2431                AllocatorValue::Abs { count: 1, owner_object_id: 98 },
2432            ))
2433            .expect("insert error");
2434
2435        let layer_set = lsm_tree.layer_set();
2436        let mut merger = layer_set.merger();
2437        let mut iter =
2438            CoalescingIterator::new(merger.query(Query::FullScan).await.expect("seek failed"))
2439                .await
2440                .expect("new failed");
2441        let ItemRef { key, value, .. } = iter.get().expect("get failed");
2442        assert_eq!(
2443            (key, value),
2444            (
2445                &AllocatorKey { device_range: (0..100 * 512).into() },
2446                &AllocatorValue::Abs { count: 1, owner_object_id: 98 },
2447            )
2448        );
2449        iter.advance().await.expect("advance failed");
2450        let ItemRef { key, value, .. } = iter.get().expect("get failed");
2451        assert_eq!(
2452            (key, value),
2453            (
2454                &AllocatorKey { device_range: (100 * 512..200 * 512).into() },
2455                &AllocatorValue::Abs { count: 1, owner_object_id: 99 }
2456            )
2457        );
2458        iter.advance().await.expect("advance failed");
2459        assert!(iter.get().is_none());
2460    }
2461
2462    fn overlap(a: &Range<u64>, b: &Range<u64>) -> u64 {
2463        if a.end > b.start && a.start < b.end {
2464            min(a.end, b.end) - max(a.start, b.start)
2465        } else {
2466            0
2467        }
2468    }
2469
2470    async fn collect_allocations(allocator: &Allocator) -> Vec<Range<u64>> {
2471        let layer_set = allocator.tree.layer_set();
2472        let mut merger = layer_set.merger();
2473        let mut iter = allocator
2474            .filter(merger.query(Query::FullScan).await.expect("seek failed"), false)
2475            .await
2476            .expect("build iterator");
2477        let mut allocations: Vec<Range<u64>> = Vec::new();
2478        while let Some(ItemRef { key: AllocatorKey { device_range }, .. }) = iter.get() {
2479            if let Some(r) = allocations.last() {
2480                assert!(device_range.start >= r.end);
2481            }
2482            allocations.push(device_range.clone().into());
2483            iter.advance().await.expect("advance failed");
2484        }
2485        allocations
2486    }
2487
2488    async fn check_allocations(allocator: &Allocator, expected_allocations: &[Range<u64>]) {
2489        let layer_set = allocator.tree.layer_set();
2490        let mut merger = layer_set.merger();
2491        let mut iter = allocator
2492            .filter(merger.query(Query::FullScan).await.expect("seek failed"), false)
2493            .await
2494            .expect("build iterator");
2495        let mut found = 0;
2496        while let Some(ItemRef { key: AllocatorKey { device_range }, .. }) = iter.get() {
2497            let mut l = device_range.length().expect("Invalid range");
2498            found += l;
2499            // Make sure that the entire range we have found completely overlaps with all the
2500            // allocations we expect to find.
2501            for range in expected_allocations {
2502                l -= overlap(range, device_range);
2503                if l == 0 {
2504                    break;
2505                }
2506            }
2507            assert_eq!(l, 0, "range {:?} not covered by expectations", device_range);
2508            iter.advance().await.expect("advance failed");
2509        }
2510        // Make sure the total we found adds up to what we expect.
2511        assert_eq!(found, expected_allocations.iter().map(|r| r.length().unwrap()).sum::<u64>());
2512    }
2513
2514    async fn test_fs() -> (OpenFxFilesystem, Arc<Allocator>) {
2515        let device = DeviceHolder::new(FakeDevice::new(4096, 4096));
2516        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2517        let allocator = fs.allocator();
2518        (fs, allocator)
2519    }
2520
2521    #[fuchsia::test]
2522    async fn test_allocations() {
2523        const STORE_OBJECT_ID: u64 = 99;
2524        let (fs, allocator) = test_fs().await;
2525        let mut transaction = fs
2526            .root_store()
2527            .new_transaction(lock_keys![], Options::default())
2528            .await
2529            .expect("new failed");
2530        let mut device_ranges = collect_allocations(&allocator).await;
2531
2532        // Expected extents:
2533        let expected = vec![
2534            0..4096,        // Superblock A (4k)
2535            4096..139264,   // root_store layer files, StoreInfo.. (33x4k blocks)
2536            139264..204800, // Superblock A extension (64k)
2537            204800..335872, // Initial Journal (128k)
2538            335872..401408, // Superblock B extension (64k)
2539            524288..528384, // Superblock B (4k)
2540        ];
2541        assert_eq!(device_ranges, expected);
2542        device_ranges.push(
2543            allocator
2544                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
2545                .await
2546                .expect("allocate failed"),
2547        );
2548        assert_eq!(device_ranges.last().unwrap().length().expect("Invalid range"), fs.block_size());
2549        device_ranges.push(
2550            allocator
2551                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
2552                .await
2553                .expect("allocate failed"),
2554        );
2555        assert_eq!(device_ranges.last().unwrap().length().expect("Invalid range"), fs.block_size());
2556        assert_eq!(overlap(&device_ranges[0], &device_ranges[1]), 0);
2557        transaction.commit().await.expect("commit failed");
2558        let mut transaction = fs
2559            .root_store()
2560            .new_transaction(lock_keys![], Options::default())
2561            .await
2562            .expect("new failed");
2563        device_ranges.push(
2564            allocator
2565                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
2566                .await
2567                .expect("allocate failed"),
2568        );
2569        assert_eq!(device_ranges[7].length().expect("Invalid range"), fs.block_size());
2570        assert_eq!(overlap(&device_ranges[5], &device_ranges[7]), 0);
2571        assert_eq!(overlap(&device_ranges[6], &device_ranges[7]), 0);
2572        transaction.commit().await.expect("commit failed");
2573
2574        check_allocations(&allocator, &device_ranges).await;
2575    }
2576
2577    #[fuchsia::test]
2578    async fn test_allocate_more_than_max_size() {
2579        const STORE_OBJECT_ID: u64 = 99;
2580        let (fs, allocator) = test_fs().await;
2581        let mut transaction = fs
2582            .root_store()
2583            .new_transaction(lock_keys![], Options::default())
2584            .await
2585            .expect("new failed");
2586        let mut device_ranges = collect_allocations(&allocator).await;
2587        device_ranges.push(
2588            allocator
2589                .allocate(&mut transaction, STORE_OBJECT_ID, fs.device().size())
2590                .await
2591                .expect("allocate failed"),
2592        );
2593        assert_eq!(
2594            device_ranges.last().unwrap().length().expect("Invalid range"),
2595            allocator.max_extent_size_bytes
2596        );
2597        transaction.commit().await.expect("commit failed");
2598
2599        check_allocations(&allocator, &device_ranges).await;
2600    }
2601
2602    #[fuchsia::test]
2603    async fn test_deallocations() {
2604        const STORE_OBJECT_ID: u64 = 99;
2605        let (fs, allocator) = test_fs().await;
2606        let initial_allocations = collect_allocations(&allocator).await;
2607
2608        let mut transaction = fs
2609            .root_store()
2610            .new_transaction(lock_keys![], Options::default())
2611            .await
2612            .expect("new failed");
2613        let device_range1 = allocator
2614            .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
2615            .await
2616            .expect("allocate failed");
2617        assert_eq!(device_range1.length().expect("Invalid range"), fs.block_size());
2618        transaction.commit().await.expect("commit failed");
2619
2620        let mut transaction = fs
2621            .root_store()
2622            .new_transaction(lock_keys![], Options::default())
2623            .await
2624            .expect("new failed");
2625        allocator
2626            .deallocate(&mut transaction, STORE_OBJECT_ID, device_range1)
2627            .await
2628            .expect("deallocate failed");
2629        transaction.commit().await.expect("commit failed");
2630
2631        check_allocations(&allocator, &initial_allocations).await;
2632    }
2633
2634    #[fuchsia::test]
2635    async fn test_mark_allocated() {
2636        const STORE_OBJECT_ID: u64 = 99;
2637        let (fs, allocator) = test_fs().await;
2638        let mut device_ranges = collect_allocations(&allocator).await;
2639        let range = {
2640            let mut transaction = fs
2641                .root_store()
2642                .new_transaction(lock_keys![], Options::default())
2643                .await
2644                .expect("new failed");
2645            // First, allocate 2 blocks.
2646            allocator
2647                .allocate(&mut transaction, STORE_OBJECT_ID, 2 * fs.block_size())
2648                .await
2649                .expect("allocate failed")
2650            // Let the transaction drop.
2651        };
2652
2653        let mut transaction = fs
2654            .root_store()
2655            .new_transaction(lock_keys![], Options::default())
2656            .await
2657            .expect("new failed");
2658
2659        // If we allocate 1 block, the two blocks that were allocated earlier should be available,
2660        // and this should return the first of them.
2661        device_ranges.push(
2662            allocator
2663                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
2664                .await
2665                .expect("allocate failed"),
2666        );
2667
2668        assert_eq!(device_ranges.last().unwrap().start, range.start);
2669
2670        // Mark the second block as allocated.
2671        let mut range2 = range.clone();
2672        range2.start += fs.block_size();
2673        allocator
2674            .mark_allocated(&mut transaction, STORE_OBJECT_ID, range2.clone())
2675            .expect("mark_allocated failed");
2676        device_ranges.push(range2);
2677
2678        // This should avoid the range we marked as allocated.
2679        device_ranges.push(
2680            allocator
2681                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
2682                .await
2683                .expect("allocate failed"),
2684        );
2685        let last_range = device_ranges.last().unwrap();
2686        assert_eq!(last_range.length().expect("Invalid range"), fs.block_size());
2687        assert_eq!(overlap(last_range, &range), 0);
2688        transaction.commit().await.expect("commit failed");
2689
2690        check_allocations(&allocator, &device_ranges).await;
2691    }
2692
2693    #[fuchsia::test]
2694    async fn test_mark_for_deletion() {
2695        const STORE_OBJECT_ID: u64 = 99;
2696        let (fs, allocator) = test_fs().await;
2697
2698        // Allocate some stuff.
2699        let initial_allocated_bytes = allocator.get_allocated_bytes();
2700        let mut device_ranges = collect_allocations(&allocator).await;
2701        let mut transaction = fs
2702            .root_store()
2703            .new_transaction(lock_keys![], Options::default())
2704            .await
2705            .expect("new failed");
2706        // Note we have a cap on individual allocation length so we allocate over multiple mutation.
2707        for _ in 0..15 {
2708            device_ranges.push(
2709                allocator
2710                    .allocate(&mut transaction, STORE_OBJECT_ID, 100 * fs.block_size())
2711                    .await
2712                    .expect("allocate failed"),
2713            );
2714            device_ranges.push(
2715                allocator
2716                    .allocate(&mut transaction, STORE_OBJECT_ID, 100 * fs.block_size())
2717                    .await
2718                    .expect("allocate2 failed"),
2719            );
2720        }
2721        transaction.commit().await.expect("commit failed");
2722        check_allocations(&allocator, &device_ranges).await;
2723
2724        assert_eq!(
2725            allocator.get_allocated_bytes(),
2726            initial_allocated_bytes + fs.block_size() * 3000
2727        );
2728
2729        // Mark for deletion.
2730        let mut transaction = fs
2731            .root_store()
2732            .new_transaction(lock_keys![], Options::default())
2733            .await
2734            .expect("new failed");
2735        allocator.mark_for_deletion(&mut transaction, STORE_OBJECT_ID);
2736        transaction.commit().await.expect("commit failed");
2737
2738        // Expect that allocated bytes is updated immediately but device ranges are still allocated.
2739        assert_eq!(allocator.get_allocated_bytes(), initial_allocated_bytes);
2740        check_allocations(&allocator, &device_ranges).await;
2741
2742        // Allocate more space than we have until we deallocate the mark_for_deletion space.
2743        // This should force a flush on allocate(). (1500 * 3 > test_fs size of 4096 blocks).
2744        device_ranges.clear();
2745
2746        let mut transaction = fs
2747            .root_store()
2748            .new_transaction(lock_keys![], Options::default())
2749            .await
2750            .expect("new failed");
2751        let target_bytes = 1500 * fs.block_size();
2752        while device_ranges.iter().map(|x| x.length().unwrap()).sum::<u64>() != target_bytes {
2753            let len = std::cmp::min(
2754                target_bytes - device_ranges.iter().map(|x| x.length().unwrap()).sum::<u64>(),
2755                100 * fs.block_size(),
2756            );
2757            device_ranges.push(
2758                allocator.allocate(&mut transaction, 100, len).await.expect("allocate failed"),
2759            );
2760        }
2761        transaction.commit().await.expect("commit failed");
2762
2763        // Have the deleted ranges cleaned up.
2764        allocator.flush().await.expect("flush failed");
2765
2766        // The flush above seems to trigger an allocation for the allocator itself.
2767        // We will just check that we have the right size for the owner we care about.
2768
2769        assert!(!allocator.get_owner_allocated_bytes().contains_key(&STORE_OBJECT_ID));
2770        assert_eq!(*allocator.get_owner_allocated_bytes().get(&100).unwrap(), target_bytes,);
2771    }
2772
2773    async fn create_file(store: &Arc<ObjectStore>, size: usize) {
2774        let root_directory =
2775            Directory::open(store, store.root_directory_object_id()).await.expect("open failed");
2776
2777        let mut transaction = store
2778            .filesystem()
2779            .root_store()
2780            .new_transaction(
2781                lock_keys![LockKey::object(
2782                    store.store_object_id(),
2783                    store.root_directory_object_id()
2784                )],
2785                Options::default(),
2786            )
2787            .await
2788            .expect("new_transaction failed");
2789        let file = root_directory
2790            .create_child_file(&mut transaction, &format!("foo {}", size))
2791            .await
2792            .expect("create_child_file failed");
2793        transaction.commit().await.expect("commit failed");
2794
2795        let buffer = file.allocate_buffer(size).await;
2796
2797        // Append some data to it.
2798        const CHUNK_SIZE: usize = 1_048_576;
2799        for offset in (0..size).step_by(CHUNK_SIZE) {
2800            let len = std::cmp::min(CHUNK_SIZE, size - offset);
2801            let mut transaction = file.new_transaction().await.expect("new_transaction failed");
2802            file.txn_write(&mut transaction, offset as u64, buffer.subslice(offset..offset + len))
2803                .await
2804                .expect("txn_write failed");
2805            transaction.commit().await.expect("commit failed");
2806        }
2807    }
2808
2809    #[fuchsia::test]
2810    async fn test_replay_with_deleted_store_and_compaction() {
2811        let (fs, _) = test_fs().await;
2812
2813        const FILE_SIZE: usize = 10_000_000;
2814
2815        let mut store_id = {
2816            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2817            let store = root_vol
2818                .new_volume("vol", NewChildStoreOptions::default())
2819                .await
2820                .expect("new_volume failed");
2821
2822            create_file(&store, FILE_SIZE).await;
2823            store.store_object_id()
2824        };
2825
2826        fs.close().await.expect("close failed");
2827        let device = fs.take_device().await;
2828        device.reopen(false);
2829
2830        let mut fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
2831
2832        // Compact so that when we replay the transaction to delete the store won't find any
2833        // mutations.
2834        fs.journal().force_compact().await.expect("compact failed");
2835
2836        for _ in 0..2 {
2837            {
2838                let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2839
2840                let transaction = fs
2841                    .root_store()
2842                    .new_transaction(
2843                        lock_keys![
2844                            LockKey::object(
2845                                root_vol.volume_directory().store().store_object_id(),
2846                                root_vol.volume_directory().object_id(),
2847                            ),
2848                            LockKey::flush(store_id)
2849                        ],
2850                        Options {
2851                            reservation: ReservationOptions::BorrowedMetadata,
2852                            ..Default::default()
2853                        },
2854                    )
2855                    .await
2856                    .expect("new_transaction failed");
2857                root_vol
2858                    .delete_volume("vol", transaction, || {})
2859                    .await
2860                    .expect("delete_volume failed");
2861
2862                let store = root_vol
2863                    .new_volume("vol", NewChildStoreOptions::default())
2864                    .await
2865                    .expect("new_volume failed");
2866                create_file(&store, FILE_SIZE).await;
2867                store_id = store.store_object_id();
2868            }
2869
2870            fs.close().await.expect("close failed");
2871            let device = fs.take_device().await;
2872            device.reopen(false);
2873
2874            fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
2875        }
2876
2877        fsck(fs.clone()).await.expect("fsck failed");
2878        fs.close().await.expect("close failed");
2879    }
2880
2881    #[fuchsia::test(threads = 4)]
2882    async fn test_compaction_delete_race() {
2883        let (fs, _allocator) = test_fs().await;
2884
2885        {
2886            const FILE_SIZE: usize = 10_000_000;
2887
2888            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2889            let store = root_vol
2890                .new_volume("vol", NewChildStoreOptions::default())
2891                .await
2892                .expect("new_volume failed");
2893
2894            create_file(&store, FILE_SIZE).await;
2895
2896            // Race compaction with deleting a store.
2897            let fs_clone = fs.clone();
2898
2899            // Even though the executor has 4 threads, it's hard to get it to run with
2900            // multiple threads.
2901            let executor_tasks = testing::force_executor_threads_to_run(4).await;
2902
2903            let task = fasync::Task::spawn(async move {
2904                fs_clone.journal().force_compact().await.expect("compact failed");
2905            });
2906
2907            // We don't need the executor tasks any more.
2908            drop(executor_tasks);
2909
2910            // This range is chosen such that it caused this test to fail after quite a low number
2911            // of iterations for the bug that this test was introduced for.
2912            let sleep = rand::random_range(3000..6000);
2913            std::thread::sleep(std::time::Duration::from_micros(sleep));
2914            log::info!("sleep {sleep}us");
2915
2916            let transaction = fs
2917                .root_store()
2918                .new_transaction(
2919                    lock_keys![
2920                        LockKey::object(
2921                            root_vol.volume_directory().store().store_object_id(),
2922                            root_vol.volume_directory().object_id(),
2923                        ),
2924                        LockKey::flush(store.store_object_id())
2925                    ],
2926                    Options {
2927                        reservation: ReservationOptions::BorrowedMetadata,
2928                        ..Default::default()
2929                    },
2930                )
2931                .await
2932                .expect("new_transaction failed");
2933            root_vol.delete_volume("vol", transaction, || {}).await.expect("delete_volume failed");
2934
2935            task.await;
2936        }
2937
2938        fs.journal().force_compact().await.expect("compact failed");
2939        fs.close().await.expect("close failed");
2940
2941        let device = fs.take_device().await;
2942        device.reopen(false);
2943
2944        let fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
2945        fsck(fs.clone()).await.expect("fsck failed");
2946        fs.close().await.expect("close failed");
2947    }
2948
2949    #[fuchsia::test]
2950    async fn test_delete_multiple_volumes() {
2951        let (mut fs, _) = test_fs().await;
2952
2953        for _ in 0..50 {
2954            {
2955                let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2956                let store = root_vol
2957                    .new_volume("vol", NewChildStoreOptions::default())
2958                    .await
2959                    .expect("new_volume failed");
2960
2961                create_file(&store, 1_000_000).await;
2962
2963                let transaction = fs
2964                    .root_store()
2965                    .new_transaction(
2966                        lock_keys![
2967                            LockKey::object(
2968                                root_vol.volume_directory().store().store_object_id(),
2969                                root_vol.volume_directory().object_id(),
2970                            ),
2971                            LockKey::flush(store.store_object_id())
2972                        ],
2973                        Options {
2974                            reservation: ReservationOptions::BorrowedMetadata,
2975                            ..Default::default()
2976                        },
2977                    )
2978                    .await
2979                    .expect("new_transaction failed");
2980                root_vol
2981                    .delete_volume("vol", transaction, || {})
2982                    .await
2983                    .expect("delete_volume failed");
2984
2985                fs.allocator().flush().await.expect("flush failed");
2986            }
2987
2988            fs.close().await.expect("close failed");
2989            let device = fs.take_device().await;
2990            device.reopen(false);
2991
2992            fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
2993        }
2994
2995        fsck(fs.clone()).await.expect("fsck failed");
2996        fs.close().await.expect("close failed");
2997    }
2998
2999    #[fuchsia::test]
3000    async fn test_allocate_free_reallocate() {
3001        const STORE_OBJECT_ID: u64 = 99;
3002        let (fs, allocator) = test_fs().await;
3003
3004        // Allocate some stuff.
3005        let mut device_ranges = Vec::new();
3006        let mut transaction = fs
3007            .root_store()
3008            .new_transaction(lock_keys![], Options::default())
3009            .await
3010            .expect("new failed");
3011        for _ in 0..30 {
3012            device_ranges.push(
3013                allocator
3014                    .allocate(&mut transaction, STORE_OBJECT_ID, 100 * fs.block_size())
3015                    .await
3016                    .expect("allocate failed"),
3017            );
3018        }
3019        transaction.commit().await.expect("commit failed");
3020
3021        assert_eq!(
3022            fs.block_size() * 3000,
3023            *allocator.get_owner_allocated_bytes().entry(STORE_OBJECT_ID).or_default()
3024        );
3025
3026        // Delete it all.
3027        let mut transaction = fs
3028            .root_store()
3029            .new_transaction(lock_keys![], Options::default())
3030            .await
3031            .expect("new failed");
3032        for range in std::mem::replace(&mut device_ranges, Vec::new()) {
3033            allocator.deallocate(&mut transaction, STORE_OBJECT_ID, range).await.expect("dealloc");
3034        }
3035        transaction.commit().await.expect("commit failed");
3036
3037        assert_eq!(0, *allocator.get_owner_allocated_bytes().entry(STORE_OBJECT_ID).or_default());
3038
3039        // Allocate some more stuff. Due to storage pressure, this requires us to flush device
3040        // before reusing the above space
3041        let mut transaction = fs
3042            .root_store()
3043            .new_transaction(lock_keys![], Options::default())
3044            .await
3045            .expect("new failed");
3046        let target_len = 1500 * fs.block_size();
3047        while device_ranges.iter().map(|i| i.length().unwrap()).sum::<u64>() != target_len {
3048            let len = target_len - device_ranges.iter().map(|i| i.length().unwrap()).sum::<u64>();
3049            device_ranges.push(
3050                allocator
3051                    .allocate(&mut transaction, STORE_OBJECT_ID, len)
3052                    .await
3053                    .expect("allocate failed"),
3054            );
3055        }
3056        transaction.commit().await.expect("commit failed");
3057
3058        assert_eq!(
3059            fs.block_size() * 1500,
3060            *allocator.get_owner_allocated_bytes().entry(STORE_OBJECT_ID).or_default()
3061        );
3062    }
3063
3064    #[fuchsia::test]
3065    async fn test_flush() {
3066        const STORE_OBJECT_ID: u64 = 99;
3067
3068        let mut device_ranges = Vec::new();
3069        let device = {
3070            let (fs, allocator) = test_fs().await;
3071            let mut transaction = fs
3072                .root_store()
3073                .new_transaction(lock_keys![], Options::default())
3074                .await
3075                .expect("new failed");
3076            device_ranges.push(
3077                allocator
3078                    .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3079                    .await
3080                    .expect("allocate failed"),
3081            );
3082            device_ranges.push(
3083                allocator
3084                    .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3085                    .await
3086                    .expect("allocate failed"),
3087            );
3088            device_ranges.push(
3089                allocator
3090                    .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3091                    .await
3092                    .expect("allocate failed"),
3093            );
3094            transaction.commit().await.expect("commit failed");
3095
3096            allocator.flush().await.expect("flush failed");
3097
3098            fs.close().await.expect("close failed");
3099            fs.take_device().await
3100        };
3101
3102        device.reopen(false);
3103        let fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
3104        let allocator = fs.allocator();
3105
3106        let allocated = collect_allocations(&allocator).await;
3107
3108        // Make sure the ranges we allocated earlier are still allocated.
3109        for i in &device_ranges {
3110            let mut overlapping = 0;
3111            for j in &allocated {
3112                overlapping += overlap(i, j);
3113            }
3114            assert_eq!(overlapping, i.length().unwrap(), "Range {i:?} not allocated");
3115        }
3116
3117        let mut transaction = fs
3118            .root_store()
3119            .new_transaction(lock_keys![], Options::default())
3120            .await
3121            .expect("new failed");
3122        let range = allocator
3123            .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3124            .await
3125            .expect("allocate failed");
3126
3127        // Make sure the range just allocated doesn't overlap any other allocated ranges.
3128        for r in &allocated {
3129            assert_eq!(overlap(r, &range), 0);
3130        }
3131        transaction.commit().await.expect("commit failed");
3132    }
3133
3134    #[fuchsia::test]
3135    async fn test_dropped_transaction() {
3136        const STORE_OBJECT_ID: u64 = 99;
3137        let (fs, allocator) = test_fs().await;
3138        let allocated_range = {
3139            let mut transaction = fs
3140                .root_store()
3141                .new_transaction(lock_keys![], Options::default())
3142                .await
3143                .expect("new_transaction failed");
3144            allocator
3145                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3146                .await
3147                .expect("allocate failed")
3148        };
3149        // After dropping the transaction and attempting to allocate again, we should end up with
3150        // the same range because the reservation should have been released.
3151        let mut transaction = fs
3152            .root_store()
3153            .new_transaction(lock_keys![], Options::default())
3154            .await
3155            .expect("new_transaction failed");
3156        assert_eq!(
3157            allocator
3158                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3159                .await
3160                .expect("allocate failed"),
3161            allocated_range
3162        );
3163    }
3164
3165    #[fuchsia::test]
3166    async fn test_cleanup_removed_owner() {
3167        const STORE_OBJECT_ID: u64 = 99;
3168        let device = {
3169            let (fs, allocator) = test_fs().await;
3170
3171            assert!(!allocator.get_owner_allocated_bytes().contains_key(&STORE_OBJECT_ID));
3172            {
3173                let mut transaction = fs
3174                    .root_store()
3175                    .new_transaction(lock_keys![], Options::default())
3176                    .await
3177                    .unwrap();
3178                allocator
3179                    .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3180                    .await
3181                    .expect("Allocating");
3182                transaction.commit().await.expect("Committing.");
3183            }
3184            allocator.flush().await.expect("Flushing");
3185            assert!(allocator.get_owner_allocated_bytes().contains_key(&STORE_OBJECT_ID));
3186            {
3187                let mut transaction = fs
3188                    .root_store()
3189                    .new_transaction(lock_keys![], Options::default())
3190                    .await
3191                    .unwrap();
3192                allocator.mark_for_deletion(&mut transaction, STORE_OBJECT_ID);
3193                transaction.commit().await.expect("Committing.");
3194            }
3195            assert!(!allocator.get_owner_allocated_bytes().contains_key(&STORE_OBJECT_ID));
3196            fs.close().await.expect("Closing");
3197            fs.take_device().await
3198        };
3199
3200        device.reopen(false);
3201        let fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
3202        let allocator = fs.allocator();
3203        assert!(!allocator.get_owner_allocated_bytes().contains_key(&STORE_OBJECT_ID));
3204    }
3205
3206    #[fuchsia::test]
3207    async fn test_allocated_bytes() {
3208        const STORE_OBJECT_ID: u64 = 99;
3209        let (fs, allocator) = test_fs().await;
3210
3211        let initial_allocated_bytes = allocator.get_allocated_bytes();
3212
3213        // Verify allocated_bytes reflects allocation changes.
3214        let allocated_bytes = initial_allocated_bytes + fs.block_size();
3215        let allocated_range = {
3216            let mut transaction = fs
3217                .root_store()
3218                .new_transaction(lock_keys![], Options::default())
3219                .await
3220                .expect("new_transaction failed");
3221            let range = allocator
3222                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3223                .await
3224                .expect("allocate failed");
3225            transaction.commit().await.expect("commit failed");
3226            assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3227            range
3228        };
3229
3230        {
3231            let mut transaction = fs
3232                .root_store()
3233                .new_transaction(lock_keys![], Options::default())
3234                .await
3235                .expect("new_transaction failed");
3236            allocator
3237                .allocate(&mut transaction, STORE_OBJECT_ID, fs.block_size().get())
3238                .await
3239                .expect("allocate failed");
3240
3241            // Prior to committing, the count of allocated bytes shouldn't change.
3242            assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3243        }
3244
3245        // After dropping the prior transaction, the allocated bytes still shouldn't have changed.
3246        assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3247
3248        // Verify allocated_bytes reflects deallocations.
3249        let deallocate_range = allocated_range.start + 20..allocated_range.end - 20;
3250        let mut transaction = fs
3251            .root_store()
3252            .new_transaction(lock_keys![], Options::default())
3253            .await
3254            .expect("new failed");
3255        allocator
3256            .deallocate(&mut transaction, STORE_OBJECT_ID, deallocate_range)
3257            .await
3258            .expect("deallocate failed");
3259
3260        // Before committing, there should be no change.
3261        assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3262
3263        transaction.commit().await.expect("commit failed");
3264
3265        // After committing, all but 40 bytes should remain allocated.
3266        assert_eq!(allocator.get_allocated_bytes(), initial_allocated_bytes + 40);
3267    }
3268
3269    #[fuchsia::test]
3270    async fn test_persist_bytes_limit() {
3271        const LIMIT: u64 = 12345;
3272        const OWNER_ID: u64 = 12;
3273
3274        let (fs, allocator) = test_fs().await;
3275        {
3276            let mut transaction = fs
3277                .root_store()
3278                .new_transaction(lock_keys![], Options::default())
3279                .await
3280                .expect("new_transaction failed");
3281            allocator
3282                .set_bytes_limit(&mut transaction, OWNER_ID, LIMIT)
3283                .expect("Failed to set limit.");
3284            assert!(allocator.inner.lock().info.limit_bytes.get(&OWNER_ID).is_none());
3285            transaction.commit().await.expect("Failed to commit transaction");
3286            let bytes: u64 = *allocator
3287                .inner
3288                .lock()
3289                .info
3290                .limit_bytes
3291                .get(&OWNER_ID)
3292                .expect("Failed to find limit");
3293            assert_eq!(LIMIT, bytes);
3294        }
3295    }
3296
3297    /// Given a sorted list of non-overlapping ranges, this will coalesce adjacent ranges.
3298    /// This allows comparison of equivalent sets of ranges which may occur due to differences
3299    /// across allocator strategies.
3300    fn coalesce_ranges(ranges: Vec<Range<u64>>) -> Vec<Range<u64>> {
3301        let mut coalesced = Vec::new();
3302        let mut prev: Option<Range<u64>> = None;
3303        for range in ranges {
3304            if let Some(prev_range) = &mut prev {
3305                if range.start == prev_range.end {
3306                    prev_range.end = range.end;
3307                } else {
3308                    coalesced.push(prev_range.clone());
3309                    prev = Some(range);
3310                }
3311            } else {
3312                prev = Some(range);
3313            }
3314        }
3315        if let Some(prev_range) = prev {
3316            coalesced.push(prev_range);
3317        }
3318        coalesced
3319    }
3320
3321    #[fuchsia::test]
3322    async fn test_take_for_trimming() {
3323        const STORE_OBJECT_ID: u64 = 99;
3324
3325        // Allocate a large chunk, then free a few bits of it, so we have free chunks interleaved
3326        // with allocated chunks.
3327        let allocated_range;
3328        let expected_free_ranges;
3329        let device = {
3330            let (fs, allocator) = test_fs().await;
3331            let bs = fs.block_size();
3332            let mut transaction = fs
3333                .root_store()
3334                .new_transaction(lock_keys![], Options::default())
3335                .await
3336                .expect("new failed");
3337            allocated_range = allocator
3338                .allocate(&mut transaction, STORE_OBJECT_ID, 32 * bs)
3339                .await
3340                .expect("allocate failed");
3341            transaction.commit().await.expect("commit failed");
3342
3343            let mut transaction = fs
3344                .root_store()
3345                .new_transaction(lock_keys![], Options::default())
3346                .await
3347                .expect("new failed");
3348            let base = allocated_range.start;
3349            expected_free_ranges = vec![
3350                base..(base + (bs * 1)),
3351                (base + (bs * 2))..(base + (bs * 3)),
3352                // Note that the next three ranges are adjacent and will be treated as one free
3353                // range once applied.  We separate them here to exercise the handling of "large"
3354                // free ranges.
3355                (base + (bs * 4))..(base + (bs * 8)),
3356                (base + (bs * 8))..(base + (bs * 12)),
3357                (base + (bs * 12))..(base + (bs * 13)),
3358                (base + (bs * 29))..(base + (bs * 30)),
3359            ];
3360            for range in &expected_free_ranges {
3361                allocator
3362                    .deallocate(&mut transaction, STORE_OBJECT_ID, range.clone())
3363                    .await
3364                    .expect("deallocate failed");
3365            }
3366            transaction.commit().await.expect("commit failed");
3367
3368            allocator.flush().await.expect("flush failed");
3369
3370            fs.close().await.expect("close failed");
3371            fs.take_device().await
3372        };
3373
3374        device.reopen(false);
3375        let fs = FxFilesystemBuilder::new().open(device).await.expect("open failed");
3376        let allocator = fs.allocator();
3377
3378        // These values were picked so that each of them would be the reason why
3379        // collect_free_extents finished, and so we would return after partially processing one of
3380        // the free extents.
3381        let max_extent_size = (fs.block_size() * 4) as usize;
3382        const EXTENTS_PER_BATCH: usize = 2;
3383        let mut free_ranges = vec![];
3384        let mut offset = allocated_range.start;
3385        while offset < allocated_range.end {
3386            let free = allocator
3387                .take_for_trimming(offset, max_extent_size, EXTENTS_PER_BATCH)
3388                .await
3389                .expect("take_for_trimming failed");
3390            free_ranges.extend(
3391                free.extents().iter().filter(|range| range.end <= allocated_range.end).cloned(),
3392            );
3393            offset = free.extents().last().expect("Unexpectedly hit the end of free extents").end;
3394        }
3395        // Coalesce adjacent free ranges because the allocator may return smaller aligned chunks
3396        // but the overall range should always be equivalent.
3397        let coalesced_free_ranges = coalesce_ranges(free_ranges);
3398        let coalesced_expected_free_ranges = coalesce_ranges(expected_free_ranges);
3399
3400        assert_eq!(coalesced_free_ranges, coalesced_expected_free_ranges);
3401    }
3402
3403    #[fuchsia::test]
3404    async fn test_allocations_wait_for_free_extents() {
3405        const STORE_OBJECT_ID: u64 = 99;
3406        let (fs, allocator) = test_fs().await;
3407        let allocator_clone = allocator.clone();
3408
3409        let mut transaction = fs
3410            .root_store()
3411            .new_transaction(lock_keys![], Options::default())
3412            .await
3413            .expect("new failed");
3414
3415        // Tie up all of the free extents on the device, and make sure allocations block.
3416        let max_extent_size = fs.device().size() as usize;
3417        const EXTENTS_PER_BATCH: usize = usize::MAX;
3418
3419        // HACK: Treat `trimmable_extents` as being locked by `trim_done` (i.e. it should only be
3420        // accessed whilst `trim_done` is locked). We can't combine them into the same mutex,
3421        // because the inner type would be "poisoned" by the lifetime parameter of
3422        // `trimmable_extents` (which is in the lifetime of `allocator`), and then we can't move it
3423        // into `alloc_task` which would require a `'static` lifetime.
3424        let trim_done = Arc::new(Mutex::new(false));
3425        let trimmable_extents = allocator
3426            .take_for_trimming(0, max_extent_size, EXTENTS_PER_BATCH)
3427            .await
3428            .expect("take_for_trimming failed");
3429
3430        let trim_done_clone = trim_done.clone();
3431        let bs = fs.block_size();
3432        let alloc_task = fasync::Task::spawn(async move {
3433            allocator_clone
3434                .allocate(&mut transaction, STORE_OBJECT_ID, bs.get())
3435                .await
3436                .expect("allocate should fail");
3437            {
3438                assert!(*trim_done_clone.lock(), "Allocation finished before trim completed");
3439            }
3440            transaction.commit().await.expect("commit failed");
3441        });
3442
3443        // Add a small delay to simulate the trim taking some nonzero amount of time.  Otherwise,
3444        // this will almost certainly always beat the allocation attempt.
3445        fasync::Timer::new(std::time::Duration::from_millis(100)).await;
3446
3447        // Once the free extents are released, the task should unblock.
3448        {
3449            let mut trim_done = trim_done.lock();
3450            std::mem::drop(trimmable_extents);
3451            *trim_done = true;
3452        }
3453
3454        alloc_task.await;
3455    }
3456
3457    #[fuchsia::test]
3458    async fn test_allocation_with_reservation_not_multiple_of_block_size() {
3459        const STORE_OBJECT_ID: u64 = 99;
3460        let (fs, allocator) = test_fs().await;
3461
3462        // Reserve an amount that isn't a multiple of the block size.
3463        const RESERVATION_AMOUNT: u64 = TRANSACTION_METADATA_MAX_AMOUNT + 5000;
3464        let reservation =
3465            allocator.clone().reserve(Some(STORE_OBJECT_ID), RESERVATION_AMOUNT).unwrap();
3466
3467        let mut transaction = fs
3468            .root_store()
3469            .new_transaction(
3470                lock_keys![],
3471                Options {
3472                    reservation: ReservationOptions::Hold(&reservation),
3473                    ..Options::default()
3474                },
3475            )
3476            .await
3477            .expect("new failed");
3478
3479        let range = allocator
3480            .allocate(
3481                &mut transaction,
3482                STORE_OBJECT_ID,
3483                fs.block_size().align_up(RESERVATION_AMOUNT).unwrap(),
3484            )
3485            .await
3486            .expect("allocate faiiled");
3487        assert!(fs.block_size().is_aligned(range.end - range.start));
3488
3489        println!("{}", range.end - range.start);
3490    }
3491
3492    #[fuchsia::test]
3493    async fn test_concurrent_take_for_trimming_returns_error() {
3494        let (fs, allocator) = test_fs().await;
3495        let max_extent_size = fs.device().size() as usize;
3496        const EXTENTS_PER_BATCH: usize = usize::MAX;
3497
3498        {
3499            let _trimmable_extents = allocator
3500                .take_for_trimming(0, max_extent_size, EXTENTS_PER_BATCH)
3501                .await
3502                .expect("take_for_trimming failed");
3503
3504            let res = allocator.take_for_trimming(0, max_extent_size, EXTENTS_PER_BATCH).await;
3505            assert!(matches!(res, Err(e) if FxfsError::AlreadyBound.matches(&e)));
3506        }
3507
3508        fs.close().await.expect("close failed");
3509    }
3510}