1pub 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
138pub trait ReservationOwner: Send + Sync {
140 fn release_reservation(&self, owner_object_id: Option<u64>, amount: u64);
144}
145
146pub 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: u64,
164
165 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 pub fn amount(&self) -> u64 {
195 self.inner.lock().amount
196 }
197
198 pub fn add(&self, amount: u64) {
200 self.inner.lock().amount += amount;
201 }
202
203 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 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 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 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 pub fn commit(&self, amount: u64) {
244 let mut inner = self.inner.lock();
245 inner.reserved -= amount;
246 inner.amount -= amount;
247 }
248
249 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 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 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
304pub 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 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 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 self.device_range.cmp(&other.device_range)
439 }
440}
441
442pub 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 None,
454 Abs { count: u64, owner_object_id: u64 },
458}
459
460pub type AllocatorItem = Item<AllocatorKey, AllocatorValue>;
461
462pub type AllocatorInfo = AllocatorInfoV32;
464
465#[derive(Debug, Default, Clone, Deserialize, Serialize, TypeFingerprint, Versioned)]
466pub struct AllocatorInfoV32 {
467 pub layers: Vec<u64>,
469 pub allocated_bytes: BTreeMap<u64, u64>,
471 pub marked_for_deletion: HashSet<u64>,
474 pub limit_bytes: BTreeMap<u64, u64>,
477}
478
479const MAX_ALLOCATOR_INFO_SERIALIZED_SIZE: usize = 131_072;
480
481pub fn max_extent_size_for_block_size(block_size: BlockSize) -> u64 {
483 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 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#[derive(Debug, Default, PartialEq)]
519struct ByteTracking {
520 allocated_bytes: Saturating<u64>,
527
528 uncommitted_allocated_bytes: u64,
531
532 reserved_bytes: u64,
534
535 committed_deallocated_bytes: u64,
539}
540
541impl ByteTracking {
542 fn used_bytes(&self) -> Saturating<u64> {
545 self.allocated_bytes + Saturating(self.uncommitted_allocated_bytes + self.reserved_bytes)
546 }
547
548 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 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 log_file_offset: u64,
569 range: Range<u64>,
571 owner_object_id: u64,
573}
574
575struct Inner {
576 info: AllocatorInfo,
577
578 opened: bool,
581
582 dropped_temporary_allocations: Vec<Range<u64>>,
610
611 owner_bytes: BTreeMap<u64, ByteTracking>,
614
615 unattributed_reserved_bytes: u64,
618
619 committed_deallocated: VecDeque<CommittedDeallocation>,
621
622 trim_reserved_bytes: u64,
625
626 trim_listener: Option<EventListener>,
630
631 strategy: strategy::BestFit,
633
634 allocation_size_histogram: [u64; 64],
636 rebuild_strategy_trigger_histogram: [u64; 64],
638
639 marked_for_deletion: HashSet<u64>,
643
644 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 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 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 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 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
753pub struct TrimmableExtents<'a> {
756 allocator: &'a Allocator,
757 extents: Vec<Range<u64>>,
758 _drop_event: DropEvent,
762}
763
764impl<'a> TrimmableExtents<'a> {
765 pub fn extents(&self) -> &Vec<Range<u64>> {
766 &self.extents
767 }
768
769 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 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 pub fn enable_allocations(&self) {
842 self.allocations_allowed.store(true, Ordering::SeqCst);
843 }
844
845 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 pub fn allocation_size_histogram(&self) -> [u64; 64] {
872 self.inner.lock().allocation_size_histogram
873 }
874
875 pub async fn create(&self, transaction: &mut Transaction<'_>) -> Result<(), Error> {
877 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 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 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 {
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 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 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 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 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 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 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 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 pub fn parent_objects(&self) -> Vec<u64> {
1143 self.inner.lock().info.layers.clone()
1146 }
1147
1148 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 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 pub fn owner_bytes_debug(&self) -> String {
1164 format!("{:?}", self.inner.lock().owner_bytes)
1165 }
1166
1167 fn needs_sync(&self) -> bool {
1168 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 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 let mut inner = self.inner.lock();
1190 inner.remove_reservation(old_owner_object_id, amount);
1191 inner.add_reservation(None, amount);
1192 }
1193
1194 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 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 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 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 assert_eq!(inner.uncommitted_allocated_bytes(), 0);
1307 assert_eq!(inner.reserved_bytes(), 0);
1308 }
1309}
1310
1311#[fxfs_trace::trace]
1312impl Allocator {
1313 pub fn object_id(&self) -> u64 {
1315 self.object_id
1316 }
1317
1318 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 #[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 let reservation = if let Some(reservation) = transaction.allocator_reservation() {
1372 match reservation.owner_object_id {
1373 None => assert!(self.is_system_store(owner_object_id)),
1375 Some(res_owner_object_id) => assert_eq!(owner_object_id, res_owner_object_id),
1377 };
1378 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 let device_used = inner.used_bytes();
1388 let owner_bytes_left = inner.owner_id_bytes_left(owner_object_id);
1389 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 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 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)] let _guard = 'sync: loop {
1456 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 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 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 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 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 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 |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 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 #[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 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 pub fn set_bytes_limit(
1644 &self,
1645 transaction: &mut Transaction<'_>,
1646 owner_object_id: u64,
1647 bytes: u64,
1648 ) -> Result<(), Error> {
1649 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 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 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 #[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 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 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 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 pub fn mark_for_deletion(&self, transaction: &mut Transaction<'_>, owner_object_id: u64) {
1738 transaction.add(
1741 self.object_id(),
1742 Mutation::Allocator(AllocatorMutation::MarkForDeletion(owner_object_id)),
1743 );
1744 }
1745
1746 pub fn did_flush_device(&self, flush_log_offset: u64) {
1749 #[allow(clippy::never_loop)] 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 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 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 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 assert!(inner.volumes_deleted_pending_sync.contains(&owner_object_id));
1786 }
1787 }
1788 }
1789 }
1790
1791 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 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 pub fn get_allocated_bytes(&self) -> u64 {
1842 self.inner.lock().allocated_bytes().0
1843 }
1844
1845 pub fn get_disk_bytes(&self) -> u64 {
1847 self.device_size
1848 }
1849
1850 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 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 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 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 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 self.inner.lock().info.limit_bytes.insert(owner_object_id, bytes);
1973 }
1974 Mutation::BeginFlush => {
1975 self.tree.seal();
1976 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
2021pub 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 ensure!(
2062 left.key.device_range.end <= right.key.device_range.start,
2063 FxfsError::Inconsistent
2064 );
2065 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 transaction.add(self.allocator.object_id(), Mutation::BeginFlush);
2124 let info = transaction
2125 .commit_with_callback(|_| {
2126 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 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 for object_id in &info.layers {
2192 root_store.add_to_graveyard(&mut transaction, *object_id);
2193 }
2194
2195 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 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 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 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 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 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 let options = bincode::DefaultOptions::new().allow_trailing_bytes();
2338 let expected_bytes = options.serialize(&range).unwrap();
2339
2340 let mut key_bytes = Vec::new();
2342 key.serialize_into(&mut key_bytes).unwrap();
2343
2344 assert_eq!(key_bytes, expected_bytes);
2346
2347 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 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 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 let expected = vec![
2534 0..4096, 4096..139264, 139264..204800, 204800..335872, 335872..401408, 524288..528384, ];
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 allocator
2647 .allocate(&mut transaction, STORE_OBJECT_ID, 2 * fs.block_size())
2648 .await
2649 .expect("allocate failed")
2650 };
2652
2653 let mut transaction = fs
2654 .root_store()
2655 .new_transaction(lock_keys![], Options::default())
2656 .await
2657 .expect("new failed");
2658
2659 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 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 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 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 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 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 assert_eq!(allocator.get_allocated_bytes(), initial_allocated_bytes);
2740 check_allocations(&allocator, &device_ranges).await;
2741
2742 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 allocator.flush().await.expect("flush failed");
2765
2766 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 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 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 let fs_clone = fs.clone();
2898
2899 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 drop(executor_tasks);
2909
2910 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 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 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 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 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 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 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 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 assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3243 }
3244
3245 assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3247
3248 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 assert_eq!(allocator.get_allocated_bytes(), allocated_bytes);
3262
3263 transaction.commit().await.expect("commit failed");
3264
3265 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 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 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 (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 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 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 let max_extent_size = fs.device().size() as usize;
3417 const EXTENTS_PER_BATCH: usize = usize::MAX;
3418
3419 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 fasync::Timer::new(std::time::Duration::from_millis(100)).await;
3446
3447 {
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 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}