1use crate::errors::FxfsError;
6use crate::fsck::{FsckOptions, fsck_volume_with_options, fsck_with_options};
7use crate::hooks::HooksHandle;
8use crate::log::*;
9use crate::metrics;
10use crate::object_handle::LayerObject;
11use crate::object_store::allocator::Allocator;
12use crate::object_store::directory::Directory;
13use crate::object_store::graveyard::Graveyard;
14use crate::object_store::journal::super_block::{SuperBlockHeader, SuperBlockInstance};
15use crate::object_store::journal::{self, Journal, JournalCheckpoint, JournalOptions};
16use crate::object_store::object_manager::ObjectManager;
17use crate::object_store::transaction::{
18 self, AssocObj, LockKey, LockManager, MetadataReservation, Mutation, ObjectMutationIterator,
19 TRANSACTION_METADATA_MAX_AMOUNT, Transaction, WriteGuard, lock_keys,
20};
21use crate::object_store::volume::{VOLUMES_DIRECTORY, root_volume};
22use crate::object_store::{
23 AttributeId, DataObjectHandle, NewChildStoreOptions, ObjectStore, StoreOptions,
24};
25use crate::range::RangeExt;
26use crate::serialized_types::{LATEST_VERSION, Version};
27use anyhow::{Context, Error, anyhow, bail};
28use async_trait::async_trait;
29use event_listener::Event;
30use fuchsia_async as fasync;
31use fuchsia_async::condition::Condition;
32use fuchsia_inspect::{Inspector, LazyNode, NumericProperty as _, UintProperty};
33use fuchsia_sync::Mutex;
34use futures::future::BoxFuture;
35use futures::stream::BoxStream;
36use futures::{FutureExt, Stream};
37use fxfs_crypto::{Crypt, UnwrappedKey};
38use fxfs_trace::{TraceFutureExt, trace_future_args};
39use static_assertions::const_assert;
40use std::pin::pin;
41use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
42use std::sync::{Arc, OnceLock, Weak};
43use std::task::Poll;
44use std::time::{Duration, Instant};
45use storage_device::{Device, DeviceHolder};
46use storage_units::BlockSize;
47
48pub const MIN_BLOCK_SIZE: BlockSize = BlockSize::SIZE_4KIB;
49pub const MAX_BLOCK_SIZE: BlockSize = BlockSize::SIZE_64KIB;
50
51pub const MAX_FILE_SIZE: u64 = i64::MAX as u64 - 4095;
56const_assert!(9223372036854771712 == MAX_FILE_SIZE);
57
58use futures::stream::StreamExt;
59
60pub(crate) const MAX_IN_FLIGHT_TRANSACTIONS: u64 = 4;
62
63const TRIM_AFTER_BOOT_TIMER: Duration = Duration::from_secs(60 * 60);
67
68const TRIM_INTERVAL_TIMER: Duration = Duration::from_secs(60 * 60 * 24);
70
71const CLEAN_TRANSFER_BUFFER_INTERVAL: Duration = Duration::from_secs(60);
74
75#[derive(Debug)]
78pub struct WakeLease {
79 #[cfg(target_os = "fuchsia")]
80 _lease: zx::Handle,
81}
82
83#[cfg(target_os = "fuchsia")]
84impl WakeLease {
85 pub fn new<H: TryInto<zx::Handle>>(lease: H) -> Self
86 where
87 H::Error: std::fmt::Debug,
88 {
89 Self { _lease: lease.try_into().expect("lease handle must be valid") }
90 }
91}
92
93pub trait PowerManager: Send + Sync {
94 fn watch_battery(self: Arc<Self>) -> BoxStream<'static, (bool, Option<WakeLease>)>;
96}
97
98#[async_trait]
100pub trait LayerPager: Send + Sync + 'static {
101 async fn open_layer(
102 &self,
103 handle: DataObjectHandle<ObjectStore>,
104 unwrapped_key: Option<UnwrappedKey>,
105 ) -> Result<Arc<dyn LayerObject>, Error>;
106}
107
108pub struct Info {
110 pub total_bytes: u64,
111 pub used_bytes: u64,
112}
113
114pub type PostCommitHook = Option<Box<dyn Fn() -> BoxFuture<'static, ()> + Send + Sync>>;
115
116pub struct Options {
117 pub read_only: bool,
119
120 pub hooks: Arc<HooksHandle>,
122
123 pub post_commit_hook: PostCommitHook,
126
127 pub skip_initial_reap: bool,
130
131 pub trim_config: Option<(Duration, Duration)>,
136
137 pub image_builder_mode: Option<SuperBlockInstance>,
140
141 pub inline_crypto_enabled: bool,
148
149 pub barriers_enabled: bool,
154
155 pub power_manager: Option<Arc<dyn PowerManager>>,
157
158 pub trim_charger_wait: Duration,
160
161 pub allow_type3_blobs: bool,
164}
165
166impl Default for Options {
167 fn default() -> Self {
168 Options {
169 read_only: false,
170 hooks: Arc::<HooksHandle>::default(),
171 post_commit_hook: None,
172 skip_initial_reap: false,
173 trim_config: Some((TRIM_AFTER_BOOT_TIMER, TRIM_INTERVAL_TIMER)),
174 image_builder_mode: None,
175 inline_crypto_enabled: false,
176 barriers_enabled: false,
177 power_manager: None,
178 trim_charger_wait: Duration::from_secs(10),
179 allow_type3_blobs: false,
180 }
181 }
182}
183
184pub struct ApplyContext<'a, 'b> {
186 pub mode: ApplyMode<'a, 'b>,
188
189 pub checkpoint: JournalCheckpoint,
191}
192
193pub enum ApplyMode<'a, 'b> {
196 Replay,
197 Live(&'a Transaction<'b>),
198}
199
200impl ApplyMode<'_, '_> {
201 pub fn is_replay(&self) -> bool {
202 matches!(self, ApplyMode::Replay)
203 }
204
205 pub fn is_live(&self) -> bool {
206 matches!(self, ApplyMode::Live(_))
207 }
208}
209
210#[derive(Copy, Clone, Debug, PartialEq, Eq)]
211pub enum ForceMajor {
212 True,
213 False,
214}
215
216#[derive(Copy, Clone, Debug, PartialEq, Eq)]
217pub enum FlushReason {
218 Journal(ForceMajor),
220
221 EncryptedMutations,
223
224 UpgradeVersion,
226}
227
228#[async_trait]
231pub trait JournalingObject: Send + Sync {
232 fn apply_mutation(
236 &self,
237 mutation: Mutation,
238 context: &ApplyContext<'_, '_>,
239 assoc_obj: AssocObj<'_>,
240 ) -> Result<(), Error>;
241
242 fn drop_mutation(&self, mutation: Mutation, transaction: &Transaction<'_>);
244
245 async fn prepare_commit<'a>(
248 &self,
249 _filesystem: &'a FxFilesystem,
250 _transaction: &Transaction<'_>,
251 ) -> Result<Option<WriteGuard<'a>>, Error> {
252 Ok(None)
253 }
254
255 async fn flush(&self, reason: FlushReason) -> Result<Version, Error>;
259
260 fn write_mutations(
263 &self,
264 mutations: ObjectMutationIterator<'_, '_>,
265 mut writer: journal::Writer<'_>,
266 ) {
267 for mutation in mutations {
268 writer.write(mutation.clone());
269 }
270 }
271}
272
273#[derive(Default)]
274pub struct SyncOptions<'a> {
275 pub flush_device: bool,
282
283 pub precondition: Option<Box<dyn FnOnce() -> bool + 'a + Send>>,
286}
287
288pub struct OpenFxFilesystem(Arc<FxFilesystem>);
289
290impl OpenFxFilesystem {
291 pub async fn take_device(self) -> DeviceHolder {
294 let fut = self.device.take_when_dropped();
295 std::mem::drop(self);
296 debug_assert_not_too_long!(fut)
297 }
298}
299
300impl From<Arc<FxFilesystem>> for OpenFxFilesystem {
301 fn from(fs: Arc<FxFilesystem>) -> Self {
302 Self(fs)
303 }
304}
305
306impl Drop for OpenFxFilesystem {
307 fn drop(&mut self) {
308 if self.options.image_builder_mode.is_some()
309 && self.journal().image_builder_mode().is_some()
310 {
311 error!("OpenFxFilesystem in image_builder_mode dropped without calling close().");
312 }
313 if !self.options.read_only && !self.closed.load(Ordering::SeqCst) {
314 error!("OpenFxFilesystem dropped without first being closed. Data loss may occur.");
315 }
316 }
317}
318
319impl std::ops::Deref for OpenFxFilesystem {
320 type Target = Arc<FxFilesystem>;
321
322 fn deref(&self) -> &Self::Target {
323 &self.0
324 }
325}
326
327pub struct FxFilesystemBuilder {
328 format: bool,
329 trace: bool,
330 options: Options,
331 journal_options: JournalOptions,
332 on_new_allocator: Option<Box<dyn Fn(Arc<Allocator>) + Send + Sync>>,
333 on_new_store: Option<Box<dyn Fn(&ObjectStore) + Send + Sync>>,
334 fsck_after_every_transaction: bool,
335 layer_pager: Option<Arc<dyn LayerPager>>,
336}
337
338impl FxFilesystemBuilder {
339 pub fn new() -> Self {
340 Self {
341 format: false,
342 trace: false,
343 options: Options::default(),
344 journal_options: JournalOptions::default(),
345 on_new_allocator: None,
346 on_new_store: None,
347 fsck_after_every_transaction: false,
348 layer_pager: None,
349 }
350 }
351
352 pub fn format(mut self, format: bool) -> Self {
354 self.format = format;
355 self
356 }
357
358 pub fn trace(mut self, trace: bool) -> Self {
360 self.trace = trace;
361 self
362 }
363
364 pub fn read_only(mut self, read_only: bool) -> Self {
367 self.options.read_only = read_only;
368 self
369 }
370
371 pub fn allow_type3_blobs(mut self, allow: bool) -> Self {
373 self.options.allow_type3_blobs = allow;
374 self
375 }
376
377 pub fn image_builder_mode(mut self, mode: Option<SuperBlockInstance>) -> Self {
384 self.options.image_builder_mode = mode;
385 self
386 }
387
388 pub fn post_commit_hook(
391 mut self,
392 hook: impl Fn() -> futures::future::BoxFuture<'static, ()> + Send + Sync + 'static,
393 ) -> Self {
394 self.options.post_commit_hook = Some(Box::new(hook));
395 self
396 }
397
398 pub fn skip_initial_reap(mut self, skip_initial_reap: bool) -> Self {
401 self.options.skip_initial_reap = skip_initial_reap;
402 self
403 }
404
405 pub fn journal_options(mut self, journal_options: JournalOptions) -> Self {
407 self.journal_options = journal_options;
408 self
409 }
410
411 pub fn on_new_allocator(
413 mut self,
414 on_new_allocator: impl Fn(Arc<Allocator>) + Send + Sync + 'static,
415 ) -> Self {
416 self.on_new_allocator = Some(Box::new(on_new_allocator));
417 self
418 }
419
420 pub fn on_new_store(
422 mut self,
423 on_new_store: impl Fn(&ObjectStore) + Send + Sync + 'static,
424 ) -> Self {
425 self.on_new_store = Some(Box::new(on_new_store));
426 self
427 }
428
429 pub fn fsck_after_every_transaction(mut self, fsck_after_every_transaction: bool) -> Self {
431 self.fsck_after_every_transaction = fsck_after_every_transaction;
432 self
433 }
434
435 pub fn trim_config(mut self, delay_and_interval: Option<(Duration, Duration)>) -> Self {
436 self.options.trim_config = delay_and_interval;
437 self
438 }
439
440 pub fn power_manager(mut self, power_manager: Arc<dyn PowerManager>) -> Self {
441 self.options.power_manager = Some(power_manager);
442 self
443 }
444
445 pub fn trim_charger_wait(mut self, wait: Duration) -> Self {
446 self.options.trim_charger_wait = wait;
447 self
448 }
449
450 pub fn inline_crypto_enabled(mut self, inline_crypto_enabled: bool) -> Self {
452 self.options.inline_crypto_enabled = inline_crypto_enabled;
453 self
454 }
455
456 pub fn barriers_enabled(mut self, barriers_enabled: bool) -> Self {
459 self.options.barriers_enabled = barriers_enabled;
460 self.journal_options.barriers_enabled = barriers_enabled;
461 self
462 }
463
464 pub fn hooks(mut self, hooks: Arc<crate::hooks::HooksHandle>) -> Self {
465 self.options.hooks = hooks;
466 self
467 }
468
469 pub fn layer_pager(mut self, layer_pager: Option<Arc<dyn LayerPager>>) -> Self {
470 self.layer_pager = layer_pager;
471 self
472 }
473
474 pub async fn open(self, device: DeviceHolder) -> Result<OpenFxFilesystem, Error> {
476 let read_only = self.options.read_only;
477 if self.format && read_only {
478 bail!("Cannot initialize a filesystem as read-only");
479 }
480
481 if self.options.inline_crypto_enabled && !self.options.barriers_enabled {
483 bail!("A filesystem using inline encryption requires barriers");
484 }
485
486 let objects = Arc::new(ObjectManager::new(self.on_new_store));
487 let journal = Arc::new(Journal::new(
488 objects.clone(),
489 self.journal_options,
490 self.options.hooks.clone(),
491 ));
492
493 let image_builder_mode = self.options.image_builder_mode;
494
495 let device_block_size =
496 BlockSize::new(device.block_size()).expect("Device block size is not a power of 2");
497 let block_size = std::cmp::max(device_block_size, MIN_BLOCK_SIZE);
498 assert!(block_size <= MAX_BLOCK_SIZE, "Max supported block size is 64KiB");
499
500 let mut fsck_after_every_transaction = None;
501 let mut filesystem_options = self.options;
502 if self.fsck_after_every_transaction {
503 let instance =
504 FsckAfterEveryTransaction::new(filesystem_options.post_commit_hook.take());
505 fsck_after_every_transaction = Some(instance.clone());
506 filesystem_options.post_commit_hook =
507 Some(Box::new(move || Box::pin(instance.clone().run())));
508 }
509
510 if !read_only && !self.format {
511 device.flush().await.context("Device flush failed")?;
514 }
515
516 let filesystem = Arc::new_cyclic(|weak: &Weak<FxFilesystem>| {
517 let weak = weak.clone();
518 FxFilesystem {
519 device,
520 block_size,
521 objects: objects.clone(),
522 journal,
523 commit_mutex: futures::lock::Mutex::new(()),
524 lock_manager: LockManager::new(),
525 flush_task: Mutex::new(None),
526 background_tasks: fasync::Scope::new(),
527 closed: AtomicBool::new(true),
528 trace: self.trace,
529 graveyard: Graveyard::new(weak.clone()),
530 completed_transactions: metrics::detail().create_uint("completed_transactions", 0),
531 options: filesystem_options,
532 in_flight_transactions: AtomicU64::new(0),
533 transaction_limit_event: Event::new(),
534 layer_pager: self.layer_pager,
535 _stores_node: metrics::register_fs(move || {
536 let weak = weak.clone();
537 Box::pin(async move {
538 if let Some(fs) = weak.upgrade() {
539 fs.populate_stores_node().await
540 } else {
541 Err(anyhow!("Filesystem has been dropped"))
542 }
543 })
544 }),
545 }
546 });
547
548 filesystem.journal().set_image_builder_mode(image_builder_mode);
549
550 filesystem.journal.set_trace(self.trace);
551 if self.format {
552 filesystem.journal.init_empty(filesystem.clone()).await?;
553 if image_builder_mode.is_none() {
554 filesystem.journal.init_superblocks().await?;
557
558 filesystem.graveyard.clone().reap_async();
560 }
561
562 let root_store = filesystem.root_store();
564 root_store.set_trace(self.trace);
565 let root_directory =
566 Directory::open(&root_store, root_store.root_directory_object_id())
567 .await
568 .context("Unable to open root volume directory")?;
569 let mut transaction = root_store
570 .new_transaction(
571 lock_keys![LockKey::object(
572 root_store.store_object_id(),
573 root_directory.object_id()
574 )],
575 transaction::Options::default(),
576 )
577 .await?;
578 let volume_directory =
579 root_directory.create_child_dir(&mut transaction, VOLUMES_DIRECTORY).await?;
580 transaction.commit().await?;
581 objects.set_volume_directory(volume_directory);
582 } else {
583 filesystem
584 .journal
585 .replay(filesystem.clone(), self.on_new_allocator)
586 .await
587 .context("Journal replay failed")?;
588 filesystem.root_store().set_trace(self.trace);
589
590 if !read_only {
591 for store in objects.unlocked_stores() {
595 filesystem.graveyard.initial_reap(&store).await?;
596 }
597 }
598 }
599
600 if let Some(fsck_after_every_transaction) = fsck_after_every_transaction {
602 fsck_after_every_transaction
603 .fs
604 .set(Arc::downgrade(&filesystem))
605 .unwrap_or_else(|_| unreachable!());
606 }
607
608 filesystem.closed.store(false, Ordering::SeqCst);
609
610 if !read_only && image_builder_mode.is_none() {
611 filesystem.graveyard.clone().reap_async();
613
614 if filesystem.options.trim_config.is_some() {
615 filesystem.start_trim_task();
616 }
617 filesystem.start_clean_transfer_buffer_task();
618 }
619
620 Ok(filesystem.into())
621 }
622}
623
624pub struct FxFilesystem {
625 block_size: BlockSize,
626 objects: Arc<ObjectManager>,
627 journal: Arc<Journal>,
628 commit_mutex: futures::lock::Mutex<()>,
629 lock_manager: LockManager,
630 flush_task: Mutex<Option<fasync::Task<()>>>,
631 background_tasks: fasync::Scope,
632 closed: AtomicBool,
633 trace: bool,
635 graveyard: Arc<Graveyard>,
636 completed_transactions: UintProperty,
637 options: Options,
638
639 in_flight_transactions: AtomicU64,
641
642 transaction_limit_event: Event,
645
646 _stores_node: LazyNode,
648
649 layer_pager: Option<Arc<dyn LayerPager>>,
650
651 device: DeviceHolder,
654}
655
656#[fxfs_trace::trace]
657impl FxFilesystem {
658 pub async fn new_empty(device: DeviceHolder) -> Result<OpenFxFilesystem, Error> {
659 FxFilesystemBuilder::new().format(true).open(device).await
660 }
661
662 pub async fn open(device: DeviceHolder) -> Result<OpenFxFilesystem, Error> {
663 FxFilesystemBuilder::new().open(device).await
664 }
665
666 pub fn root_parent_store(&self) -> Arc<ObjectStore> {
667 self.objects.root_parent_store()
668 }
669
670 pub fn layer_pager(&self) -> Option<&Arc<dyn LayerPager>> {
671 self.layer_pager.as_ref()
672 }
673
674 pub async fn close(&self) -> Result<(), Error> {
675 if self.journal().image_builder_mode().is_some() {
676 self.journal().allocate_journal().await?;
677 self.journal().set_image_builder_mode(None);
678 self.journal().force_compact().await?;
679 }
680 assert_eq!(self.closed.swap(true, Ordering::SeqCst), false);
681 debug_assert_not_too_long!(self.graveyard.wait_for_reap());
682 debug_assert_not_too_long!(self.background_tasks.clone().cancel());
683 self.journal.stop_compactions().await;
684 let sync_status =
685 if self.journal().image_builder_mode().is_some() || self.options().read_only {
686 Ok(None)
687 } else {
688 self.journal.sync(SyncOptions { flush_device: true, ..Default::default() }).await
689 };
690 match &sync_status {
691 Ok(None) => {}
692 Ok(checkpoint) => info!(
693 "Filesystem closed (checkpoint={}, metadata_reservation={:?}, \
694 reservation_required={}, borrowed={})",
695 checkpoint.as_ref().unwrap().0.file_offset,
696 self.object_manager().metadata_reservation(),
697 self.object_manager().required_reservation(),
698 self.object_manager().borrowed_metadata_space(),
699 ),
700 Err(e) => error!(error:? = e; "Failed to sync filesystem; data may be lost"),
701 }
702 self.journal.terminate();
703 let flush_task = self.flush_task.lock().take();
704 if let Some(task) = flush_task {
705 debug_assert_not_too_long!(task);
706 }
707 self.device().close().await.context("Failed to close device")?;
710 sync_status.map(|_| ())
711 }
712
713 pub fn device(&self) -> Arc<dyn Device> {
714 Arc::clone(&self.device)
715 }
716
717 pub fn root_store(&self) -> Arc<ObjectStore> {
718 self.objects.root_store()
719 }
720
721 pub fn allocator(&self) -> Arc<Allocator> {
722 self.objects.allocator()
723 }
724
725 pub fn enable_allocations(&self) {
729 self.allocator().enable_allocations();
730 }
731
732 pub fn object_manager(&self) -> &Arc<ObjectManager> {
733 &self.objects
734 }
735
736 pub fn journal(&self) -> &Arc<Journal> {
737 &self.journal
738 }
739
740 pub async fn sync(&self, options: SyncOptions<'_>) -> Result<(), Error> {
741 self.journal.sync(options).await.map(|_| ())
742 }
743
744 pub fn block_size(&self) -> BlockSize {
745 self.block_size
746 }
747
748 pub fn get_info(&self) -> Info {
749 Info {
750 total_bytes: self.device.size(),
751 used_bytes: self.object_manager().allocator().get_used_bytes().0,
752 }
753 }
754
755 pub fn super_block_header(&self) -> SuperBlockHeader {
756 self.journal.super_block_header()
757 }
758
759 pub fn graveyard(&self) -> &Arc<Graveyard> {
760 &self.graveyard
761 }
762
763 pub fn trace(&self) -> bool {
764 self.trace
765 }
766
767 pub fn options(&self) -> &Options {
768 &self.options
769 }
770
771 pub fn scope(&self) -> &fasync::Scope {
772 &self.background_tasks
773 }
774
775 pub async fn lock_commits(&self) -> futures::lock::MutexGuard<'_, ()> {
782 self.commit_mutex.lock().await
783 }
784
785 #[trace]
786 pub async fn commit_transaction<R: Send>(
787 &self,
788 transaction: &mut Transaction<'_>,
789 callback: impl FnOnce(u64) -> R + Send,
790 ) -> Result<R, Error> {
791 self.hooks().on_pre_commit(transaction)?;
792 debug_assert_not_too_long!(self.lock_manager.commit_prepare(&transaction));
793
794 let mut guards = Vec::new();
799 let mut last_object_id = 0;
800 for mutation in transaction.mutations() {
801 let object_id = mutation.object_id;
802
803 if matches!(mutation.mutation, Mutation::BeginFlush | Mutation::EndFlush) {
805 continue;
806 }
807
808 if object_id == last_object_id {
809 continue;
810 }
811 assert!(object_id > last_object_id);
812 last_object_id = object_id;
813
814 if let Some(obj) = self.object_manager().journaling_object(object_id) {
815 if let Some(guard) = obj.prepare_commit(self, transaction).await? {
816 guards.push(guard);
817 }
818 }
819 }
820
821 self.maybe_start_flush_task();
822
823 self.hooks().on_before_commit();
824
825 let _guard = debug_assert_not_too_long!(self.commit_mutex.lock());
826 let journal_offset = if self.journal().image_builder_mode().is_some() {
827 let journal_checkpoint =
828 JournalCheckpoint { file_offset: 0, checksum: 0, version: LATEST_VERSION };
829 let maybe_mutation = self
830 .object_manager()
831 .apply_transaction(transaction, &journal_checkpoint)
832 .expect("Transactions must not fail in image_builder_mode");
833 if let Some(mutation) = maybe_mutation {
834 assert!(matches!(mutation, Mutation::UpdateBorrowed(_)));
835 }
839 self.object_manager().did_commit_transaction(transaction, &journal_checkpoint, 0);
840 0
841 } else {
842 self.journal.commit(transaction).await?
843 };
844
845 std::mem::drop(guards);
846 self.completed_transactions.add(1);
847
848 let result = callback(journal_offset);
853
854 if let Some(hook) = self.options.post_commit_hook.as_ref() {
855 hook().await;
856 }
857
858 Ok(result)
859 }
860
861 pub fn lock_manager(&self) -> &LockManager {
862 &self.lock_manager
863 }
864
865 pub fn hooks(&self) -> &Arc<HooksHandle> {
866 &self.options.hooks
867 }
868
869 pub(crate) fn drop_transaction(&self, transaction: &mut Transaction<'_>) {
870 self.objects.drop_transaction(transaction);
871 self.lock_manager.drop_transaction(transaction);
872 }
873
874 fn maybe_start_flush_task(&self) {
875 if self.journal.image_builder_mode().is_some() {
876 return;
877 }
878 let mut flush_task = self.flush_task.lock();
879 if flush_task.is_none() {
880 let journal = self.journal.clone();
881 *flush_task = Some(fasync::Task::spawn(
882 journal.flush_task().trace(trace_future_args!("Journal::flush_task")),
883 ));
884 }
885 }
886
887 fn start_trim_task(self: &Arc<Self>) {
888 if !self.device.supports_trim() {
889 info!("Device does not support trim; not scheduling trimming");
890 return;
891 }
892 let this = self.clone();
893 self.background_tasks
894 .spawn(this.trim_task().trace(trace_future_args!("Filesystem::trim_task")));
895 }
896
897 async fn trim_task(self: Arc<Self>) {
898 let Some((mut next_timer, _)) = self.options.trim_config else { return };
900 loop {
901 fasync::Timer::new(next_timer.clone()).await;
902
903 let start = Instant::now();
906 let result = if let Some(pm) = &self.options.power_manager {
907 let mut watcher = pm.clone().watch_battery();
908
909 let pauser = Pauser::new(self.options.trim_charger_wait);
911
912 let mut pause_future = pin!(
913 async {
914 let mut wake_lease: Option<WakeLease> = None;
915 loop {
916 let Some((using_battery, new_lease)) = watcher.next_latest().await
917 else {
918 pauser.set_pause(false);
921 drop(wake_lease); return;
923 };
924
925 pauser.set_pause(using_battery);
927
928 if using_battery {
931 wake_lease = None;
932 } else if new_lease.is_some() {
933 wake_lease = new_lease;
934 }
935 }
936 }
937 .fuse()
938 );
939
940 let mut do_trim = pin!(self.do_trim(Some(&pauser)).fuse());
941
942 loop {
943 futures::select! {
944 _ = pause_future => {}
945 result = do_trim => break result,
946 }
947 }
948
949 } else {
952 self.do_trim(None).await
953 };
954
955 let duration = start.elapsed();
956 match result {
957 Ok(bytes_trimmed) => info!(
958 "Trimmed {bytes_trimmed} bytes in {duration:?}. Next trim in \
959 {next_timer:?}",
960 ),
961 Err(error) => error!(error:?; "Failed to trim"),
962 }
963
964 let Some((_, interval)) = self.options.trim_config else { return };
965 next_timer = interval;
966 if next_timer.is_zero() {
967 fasync::yield_now().await;
968 }
969 }
970 }
971
972 async fn do_trim(&self, pauser: Option<&Pauser>) -> Result<usize, Error> {
974 const MAX_EXTENTS_PER_BATCH: usize = 8;
975 const MAX_EXTENT_SIZE: usize = 256 * 1024;
976 let mut offset = 0;
977 let mut bytes_trimmed = 0;
978 loop {
979 let allocator = self.allocator();
980 if let Some(pauser) = pauser {
981 pauser.maybe_pause().await;
982 }
983 let trimmable_extents =
984 allocator.take_for_trimming(offset, MAX_EXTENT_SIZE, MAX_EXTENTS_PER_BATCH).await?;
985 for device_range in trimmable_extents.extents() {
986 self.device.trim(device_range.clone()).await?;
987 bytes_trimmed += device_range.length()? as usize;
988 }
989 if let Some(device_range) = trimmable_extents.extents().last() {
990 offset = device_range.end;
991 } else {
992 break;
993 }
994 }
995 Ok(bytes_trimmed)
996 }
997
998 fn start_clean_transfer_buffer_task(self: &Arc<Self>) {
999 let this = self.clone();
1000 self.background_tasks.spawn(
1001 async move {
1002 loop {
1003 fasync::Timer::new(CLEAN_TRANSFER_BUFFER_INTERVAL).await;
1004 this.device().clean_transfer_buffer();
1005 }
1006 }
1007 .trace(trace_future_args!("Filesystem::clean_transfer_buffer_task")),
1008 );
1009 }
1010
1011 pub(crate) async fn add_transaction_reservation(
1018 self: &Arc<Self>,
1019 transaction: &mut Transaction<'_>,
1020 ) -> Result<(), Error> {
1021 if self.options.image_builder_mode.is_some() {
1022 return Ok(());
1025 }
1026 if !transaction.skip_journal_checks {
1027 self.maybe_start_flush_task();
1028 self.journal.check_journal_space().await?;
1029 }
1030
1031 match &mut transaction.metadata_reservation {
1032 MetadataReservation::BorrowedMetadata
1033 | MetadataReservation::BorrowedMetadataAndData => {}
1034 MetadataReservation::Hold(hold) => {
1035 let amount = TRANSACTION_METADATA_MAX_AMOUNT.saturating_sub(hold.amount());
1036 hold.add(hold.owner().reserve(amount).ok_or(FxfsError::NoSpace)?.forget());
1037 }
1038 MetadataReservation::Reservation(txn_reservation) => {
1039 txn_reservation.add(
1040 self.allocator()
1041 .reserve(
1042 txn_reservation.owner_object_id(),
1043 TRANSACTION_METADATA_MAX_AMOUNT
1044 .saturating_sub(txn_reservation.amount()),
1045 )
1046 .ok_or(FxfsError::NoSpace)?
1047 .forget(),
1048 );
1049 }
1050 }
1051 Ok(())
1052 }
1053
1054 pub(crate) async fn add_transaction(&self, skip_journal_checks: bool) {
1055 if skip_journal_checks {
1056 self.in_flight_transactions.fetch_add(1, Ordering::Relaxed);
1057 } else {
1058 let inc = || {
1059 let mut in_flights = self.in_flight_transactions.load(Ordering::Relaxed);
1060 while in_flights < MAX_IN_FLIGHT_TRANSACTIONS {
1061 match self.in_flight_transactions.compare_exchange_weak(
1062 in_flights,
1063 in_flights + 1,
1064 Ordering::Relaxed,
1065 Ordering::Relaxed,
1066 ) {
1067 Ok(_) => return true,
1068 Err(x) => in_flights = x,
1069 }
1070 }
1071 return false;
1072 };
1073 while !inc() {
1074 let listener = self.transaction_limit_event.listen();
1075 if inc() {
1076 break;
1077 }
1078 debug_assert_not_too_long!(listener);
1079 }
1080 }
1081 }
1082
1083 pub(crate) fn sub_transaction(&self) {
1084 let old = self.in_flight_transactions.fetch_sub(1, Ordering::Relaxed);
1085 assert!(old != 0);
1086 if old <= MAX_IN_FLIGHT_TRANSACTIONS {
1087 self.transaction_limit_event.notify(usize::MAX);
1088 }
1089 }
1090
1091 pub async fn truncate_guard(&self, store_id: u64, object_id: u64) -> TruncateGuard<'_> {
1092 let keys = lock_keys![LockKey::truncate(store_id, object_id,)];
1093 TruncateGuard(self.lock_manager().write_lock(keys).await)
1094 }
1095
1096 pub async fn truncate_guard_owned(
1097 self: &Arc<Self>,
1098 store_id: u64,
1099 object_id: u64,
1100 ) -> TruncateGuard<'static> {
1101 self.truncate_guard(store_id, object_id).await.into_owned(self.clone())
1102 }
1103
1104 pub async fn tombstone_object(
1114 &self,
1115 store_id: u64,
1116 object_id: u64,
1117 truncate_guard: Option<&TruncateGuard<'_>>,
1118 ) -> Result<(), Error> {
1119 let store = self
1120 .objects
1121 .store(store_id)
1122 .with_context(|| format!("Failed to get store {}", store_id))?;
1123 let options = if store_id == self.objects.root_parent_store_object_id()
1127 || store_id == self.objects.root_store_object_id()
1128 {
1129 transaction::Options {
1130 reservation: transaction::ReservationOptions::BorrowedMetadataAndData,
1131 ..Default::default()
1132 }
1133 } else {
1134 transaction::Options {
1135 reservation: transaction::ReservationOptions::BorrowedMetadata,
1136 ..Default::default()
1137 }
1138 };
1139 store.tombstone_object(object_id, options, truncate_guard).await
1140 }
1141
1142 pub async fn tombstone_attribute(
1149 &self,
1150 store_id: u64,
1151 object_id: u64,
1152 attribute_id: AttributeId,
1153 ) -> Result<(), Error> {
1154 let store = self
1155 .objects
1156 .store(store_id)
1157 .with_context(|| format!("Failed to get store {}", store_id))?;
1158 let options = if store_id == self.objects.root_parent_store_object_id()
1162 || store_id == self.objects.root_store_object_id()
1163 {
1164 transaction::Options {
1165 reservation: transaction::ReservationOptions::BorrowedMetadataAndData,
1166 ..Default::default()
1167 }
1168 } else {
1169 transaction::Options {
1170 reservation: transaction::ReservationOptions::BorrowedMetadata,
1171 ..Default::default()
1172 }
1173 };
1174 store.tombstone_attribute(object_id, attribute_id, options).await
1175 }
1176
1177 async fn populate_stores_node(&self) -> Result<Inspector, Error> {
1178 let inspector = fuchsia_inspect::Inspector::default();
1179 let root = inspector.root();
1180 root.record_child("__root", |n| self.root_store().record_data(n));
1181 root.record_child("__root_parent", |n| self.root_parent_store().record_data(n));
1182 let object_manager = self.object_manager();
1183 let volume_directory = object_manager.volume_directory();
1184 let layer_set = volume_directory.store().tree().layer_set();
1185 let mut merger = layer_set.merger();
1186 let mut iter = volume_directory.iter(&mut merger).await?;
1187 while let Some((name, id, _)) = iter.get() {
1188 if let Some(store) = object_manager.store(id) {
1189 root.record_child(name.to_string(), |n| store.record_data(n));
1190 }
1191 iter.advance().await?;
1192 }
1193 Ok(inspector)
1194 }
1195}
1196
1197#[allow(dead_code)]
1199pub struct TruncateGuard<'a>(WriteGuard<'a>);
1200
1201impl<'a> TruncateGuard<'a> {
1202 pub fn into_owned(self, fs: Arc<FxFilesystem>) -> TruncateGuard<'static> {
1203 TruncateGuard(self.0.into_owned(fs))
1204 }
1205}
1206
1207pub async fn mkfs(device: DeviceHolder) -> Result<DeviceHolder, Error> {
1209 let fs = FxFilesystem::new_empty(device).await?;
1210 fs.close().await?;
1211 Ok(fs.take_device().await)
1212}
1213
1214pub async fn mkfs_with_volume(
1218 device: DeviceHolder,
1219 volume_name: &str,
1220 crypt: Option<Arc<dyn Crypt>>,
1221) -> Result<DeviceHolder, Error> {
1222 let fs = FxFilesystem::new_empty(device).await?;
1223 {
1224 let root_volume = root_volume(fs.clone()).await.expect("Open root_volume failed");
1227 root_volume
1228 .new_volume(
1229 volume_name,
1230 NewChildStoreOptions {
1231 options: StoreOptions { crypt, ..StoreOptions::default() },
1232 ..Default::default()
1233 },
1234 )
1235 .await
1236 .expect("Create volume failed");
1237 }
1238 fs.close().await?;
1239 Ok(fs.take_device().await)
1240}
1241
1242struct FsckAfterEveryTransaction {
1243 fs: OnceLock<Weak<FxFilesystem>>,
1244 old_hook: PostCommitHook,
1245}
1246
1247impl FsckAfterEveryTransaction {
1248 fn new(old_hook: PostCommitHook) -> Arc<Self> {
1249 Arc::new(Self { fs: OnceLock::new(), old_hook })
1250 }
1251
1252 async fn run(self: Arc<Self>) {
1253 if let Some(fs) = self.fs.get().and_then(Weak::upgrade) {
1254 let options = FsckOptions {
1255 fail_on_warning: true,
1256 no_lock: true,
1257 quiet: true,
1258 ..Default::default()
1259 };
1260 fsck_with_options(fs.clone(), &options).await.expect("fsck failed");
1261 let object_manager = fs.object_manager();
1262 for store in object_manager.unlocked_stores() {
1263 let store_id = store.store_object_id();
1264 if !object_manager.is_system_store(store_id) {
1265 fsck_volume_with_options(fs.as_ref(), &options, store_id, None)
1266 .await
1267 .expect("fsck_volume_with_options failed");
1268 }
1269 }
1270 }
1271 if let Some(old_hook) = self.old_hook.as_ref() {
1272 old_hook().await;
1273 }
1274 }
1275}
1276
1277struct Pauser {
1278 pause: Condition<bool>,
1279 bounce_delay: Duration,
1280}
1281
1282impl Pauser {
1283 fn new(bounce_delay: Duration) -> Self {
1285 Self { pause: Condition::new(true), bounce_delay }
1286 }
1287
1288 async fn maybe_pause(&self) {
1289 loop {
1290 if !*self.pause.lock() {
1291 return;
1292 }
1293 self.pause.when(|p| if **p { Poll::Pending } else { Poll::Ready(()) }).await;
1294 fasync::Timer::new(self.bounce_delay).await;
1295 }
1296 }
1297
1298 fn set_pause(&self, v: bool) {
1299 let mut guard = self.pause.lock();
1300 if *guard == v {
1301 return;
1302 }
1303 *guard = v;
1304 for waker in guard.drain_wakers() {
1305 waker.wake();
1306 }
1307 }
1308}
1309
1310trait NextLatest: Stream + Unpin {
1311 async fn next_latest(&mut self) -> Option<Self::Item> {
1313 let Some(mut next) = self.next().await else { return None };
1314
1315 loop {
1317 match self.next().now_or_never() {
1318 None => return Some(next),
1319 Some(None) => return None,
1320 Some(Some(n)) => next = n,
1321 }
1322 }
1323 }
1324}
1325
1326impl<T: ?Sized + Unpin> NextLatest for T where T: Stream {}
1327
1328#[cfg(test)]
1329mod tests {
1330 use super::{FxFilesystem, FxFilesystemBuilder, FxfsError, SyncOptions};
1331 use crate::fsck::{fsck, fsck_volume};
1332 use crate::log::*;
1333 use crate::lsm_tree::Operation;
1334 use crate::lsm_tree::types::Item;
1335 use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, WriteObjectHandle};
1336 use crate::object_store::directory::{Directory, replace_child};
1337 use crate::object_store::journal::JournalOptions;
1338 use crate::object_store::journal::super_block::SuperBlockInstance;
1339 use crate::object_store::transaction::{LockKey, Options, lock_keys};
1340 use crate::object_store::volume::root_volume;
1341 use crate::object_store::{
1342 HandleOptions, NewChildStoreOptions, ObjectDescriptor, ObjectStore, StoreOptions,
1343 };
1344 use crate::range::RangeExt;
1345 use fuchsia_async as fasync;
1346 use fuchsia_sync::Mutex;
1347 use futures::future::join_all;
1348 use futures::stream::{FuturesUnordered, TryStreamExt};
1349 use fxfs_insecure_crypto::new_insecure_crypt;
1350 use rustc_hash::FxHashMap as HashMap;
1351 use std::ops::Range;
1352 use std::sync::Arc;
1353 use std::sync::atomic::{AtomicU32, Ordering};
1354 use std::time::Duration;
1355 use storage_device::DeviceHolder;
1356 use storage_device::fake_device::{self, FakeDevice};
1357 use test_case::test_case;
1358
1359 const TEST_DEVICE_BLOCK_SIZE: u32 = 512;
1360
1361 #[fuchsia::test(threads = 10)]
1362 async fn test_compaction() {
1363 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1364
1365 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1367 let root_store = fs.root_store();
1368 let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
1369 .await
1370 .expect("open failed");
1371
1372 let mut tasks = Vec::new();
1373 for i in 0..2 {
1374 let mut transaction = fs
1375 .root_store()
1376 .new_transaction(
1377 lock_keys![LockKey::object(
1378 root_store.store_object_id(),
1379 root_directory.object_id()
1380 )],
1381 Options::default(),
1382 )
1383 .await
1384 .expect("new_transaction failed");
1385 let handle = root_directory
1386 .create_child_file(&mut transaction, &format!("{}", i))
1387 .await
1388 .expect("create_child_file failed");
1389 transaction.commit().await.expect("commit failed");
1390 tasks.push(fasync::Task::spawn(async move {
1391 const TEST_DATA: &[u8] = b"hello";
1392 let mut buf = handle.allocate_buffer(TEST_DATA.len()).await;
1393 buf.copy_from_slice(TEST_DATA);
1394 for _ in 0..1500 {
1395 handle.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
1396 }
1397 }));
1398 }
1399 join_all(tasks).await;
1400 fs.sync(SyncOptions::default()).await.expect("sync failed");
1401
1402 fsck(fs.clone()).await.expect("fsck failed");
1403 fs.close().await.expect("Close failed");
1404 }
1405
1406 #[fuchsia::test]
1407 async fn test_enable_allocations() {
1408 {
1410 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1411 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1412 fs.enable_allocations();
1413 let root_store = fs.root_store();
1414 let root_directory =
1415 Directory::open(&root_store, root_store.root_directory_object_id())
1416 .await
1417 .expect("open failed");
1418 let mut transaction = fs
1419 .root_store()
1420 .new_transaction(
1421 lock_keys![LockKey::object(
1422 root_store.store_object_id(),
1423 root_directory.object_id()
1424 )],
1425 Options::default(),
1426 )
1427 .await
1428 .expect("new_transaction failed");
1429 root_directory
1430 .create_child_file(&mut transaction, "test")
1431 .await
1432 .expect("create_child_file failed");
1433 transaction.commit().await.expect("commit failed");
1434 fs.close().await.expect("close failed");
1435 }
1436
1437 {
1439 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1440 let fs = FxFilesystemBuilder::new()
1441 .format(true)
1442 .image_builder_mode(Some(SuperBlockInstance::A))
1443 .open(device)
1444 .await
1445 .expect("open failed");
1446 let root_store = fs.root_store();
1447 let root_directory =
1448 Directory::open(&root_store, root_store.root_directory_object_id())
1449 .await
1450 .expect("open failed");
1451
1452 let mut transaction = fs
1453 .root_store()
1454 .new_transaction(
1455 lock_keys![LockKey::object(
1456 root_store.store_object_id(),
1457 root_directory.object_id()
1458 )],
1459 Options::default(),
1460 )
1461 .await
1462 .expect("new_transaction failed");
1463 let handle = root_directory
1464 .create_child_file(&mut transaction, "test_fail")
1465 .await
1466 .expect("create_child_file failed");
1467 transaction.commit().await.expect("commit failed");
1468
1469 assert!(
1471 FxfsError::Unavailable
1472 .matches(&handle.allocate(0..4096).await.expect_err("allocate should fail"))
1473 );
1474
1475 fs.enable_allocations();
1477 handle.allocate(0..4096).await.expect("allocate should work after enable_allocations");
1478
1479 fs.close().await.expect("close failed");
1483 }
1484 }
1488
1489 #[fuchsia::test(threads = 10)]
1490 async fn test_replay_is_identical() {
1491 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1492 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1493
1494 fs.close().await.expect("close failed");
1497 let device = fs.take_device().await;
1498 device.reopen(false);
1499
1500 struct Mutations<K, V>(Mutex<Vec<(Operation, Item<K, V>)>>);
1501
1502 impl<K: Clone, V: Clone> Mutations<K, V> {
1503 fn new() -> Self {
1504 Mutations(Mutex::new(Vec::new()))
1505 }
1506
1507 fn push(&self, operation: Operation, item: &Item<K, V>) {
1508 self.0.lock().push((operation, item.clone()));
1509 }
1510 }
1511
1512 let open_fs = |device,
1513 object_mutations: Arc<Mutex<HashMap<_, _>>>,
1514 allocator_mutations: Arc<Mutations<_, _>>| async {
1515 FxFilesystemBuilder::new()
1516 .journal_options(JournalOptions { reclaim_size: u64::MAX, ..Default::default() })
1517 .on_new_allocator(move |allocator| {
1518 let allocator_mutations = allocator_mutations.clone();
1519 allocator.tree().set_mutation_callback(Some(Box::new(move |op, item| {
1520 allocator_mutations.push(op, item)
1521 })));
1522 })
1523 .on_new_store(move |store| {
1524 let mutations = Arc::new(Mutations::new());
1525 object_mutations.lock().insert(store.store_object_id(), mutations.clone());
1526 store.tree().set_mutation_callback(Some(Box::new(move |op, item| {
1527 mutations.push(op, item)
1528 })));
1529 })
1530 .open(device)
1531 .await
1532 .expect("open failed")
1533 };
1534
1535 let allocator_mutations = Arc::new(Mutations::new());
1536 let object_mutations = Arc::new(Mutex::new(HashMap::default()));
1537 let fs = open_fs(device, object_mutations.clone(), allocator_mutations.clone()).await;
1538
1539 let root_store = fs.root_store();
1540 let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
1541 .await
1542 .expect("open failed");
1543
1544 let mut transaction = fs
1545 .root_store()
1546 .new_transaction(
1547 lock_keys![LockKey::object(
1548 root_store.store_object_id(),
1549 root_directory.object_id()
1550 )],
1551 Options::default(),
1552 )
1553 .await
1554 .expect("new_transaction failed");
1555 let object = root_directory
1556 .create_child_file(&mut transaction, "test")
1557 .await
1558 .expect("create_child_file failed");
1559 transaction.commit().await.expect("commit failed");
1560
1561 let buf = object.allocate_buffer(10000).await;
1563 object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
1564
1565 object.write_or_append(Some(5000), buf.as_ref()).await.expect("write failed");
1567
1568 object.truncate(3000).await.expect("truncate failed");
1570
1571 let mut transaction = fs
1573 .root_store()
1574 .new_transaction(
1575 lock_keys![
1576 LockKey::object(root_store.store_object_id(), root_directory.object_id()),
1577 LockKey::object(root_store.store_object_id(), object.object_id()),
1578 ],
1579 Options::default(),
1580 )
1581 .await
1582 .expect("new_transaction failed");
1583
1584 replace_child(&mut transaction, None, (&root_directory, "test"))
1585 .await
1586 .expect("replace_child failed");
1587
1588 transaction.commit().await.expect("commit failed");
1589
1590 root_store
1592 .tombstone_object(object.object_id(), Options::default(), None)
1593 .await
1594 .expect("tombstone failed");
1595
1596 fs.close().await.expect("close failed");
1598
1599 let metadata_reservation_amount = fs.object_manager().metadata_reservation().amount();
1600
1601 let device = fs.take_device().await;
1602 device.reopen(false);
1603
1604 let replayed_object_mutations = Arc::new(Mutex::new(HashMap::default()));
1605 let replayed_allocator_mutations = Arc::new(Mutations::new());
1606 let fs = open_fs(
1607 device,
1608 replayed_object_mutations.clone(),
1609 replayed_allocator_mutations.clone(),
1610 )
1611 .await;
1612
1613 let m1 = object_mutations.lock();
1614 let m2 = replayed_object_mutations.lock();
1615 assert_eq!(m1.len(), m2.len());
1616 for (store_id, mutations) in &*m1 {
1617 let mutations = mutations.0.lock();
1618 let replayed = m2.get(&store_id).expect("Found unexpected store").0.lock();
1619 assert_eq!(mutations.len(), replayed.len());
1620 for ((op1, i1), (op2, i2)) in mutations.iter().zip(replayed.iter()) {
1621 assert_eq!(op1, op2);
1622 assert_eq!(i1.key, i2.key);
1623 assert_eq!(i1.value, i2.value);
1624 }
1625 }
1626
1627 let a1 = allocator_mutations.0.lock();
1628 let a2 = replayed_allocator_mutations.0.lock();
1629 assert_eq!(a1.len(), a2.len());
1630 for ((op1, i1), (op2, i2)) in a1.iter().zip(a2.iter()) {
1631 assert_eq!(op1, op2);
1632 assert_eq!(i1.key, i2.key);
1633 assert_eq!(i1.value, i2.value);
1634 }
1635
1636 assert_eq!(
1637 fs.object_manager().metadata_reservation().amount(),
1638 metadata_reservation_amount
1639 );
1640 }
1641
1642 #[fuchsia::test]
1643 async fn test_max_in_flight_transactions() {
1644 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1645 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1646
1647 let store = fs.root_store();
1648 let transactions = FuturesUnordered::new();
1649 for _ in 0..super::MAX_IN_FLIGHT_TRANSACTIONS {
1650 transactions.push(store.new_transaction(lock_keys![], Options::default()));
1651 }
1652 let mut transactions: Vec<_> = transactions.try_collect().await.unwrap();
1653
1654 let mut fut = std::pin::pin!(store.new_transaction(lock_keys![], Options::default()));
1656 assert!(futures::poll!(&mut fut).is_pending());
1657
1658 transactions.pop();
1660
1661 assert!(futures::poll!(&mut fut).is_ready());
1662 }
1663
1664 #[fuchsia::test(threads = 10)]
1666 async fn test_continuously_trim() {
1667 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1668 let fs = FxFilesystemBuilder::new()
1669 .trim_config(Some((Duration::ZERO, Duration::ZERO)))
1670 .format(true)
1671 .open(device)
1672 .await
1673 .expect("open failed");
1674 fasync::Timer::new(Duration::from_millis(10)).await;
1676
1677 let root_store = fs.root_store();
1680 let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
1681 .await
1682 .expect("open failed");
1683 for _ in 0..100 {
1684 let mut transaction = fs
1685 .root_store()
1686 .new_transaction(
1687 lock_keys![LockKey::object(
1688 root_store.store_object_id(),
1689 root_directory.object_id()
1690 )],
1691 Options::default(),
1692 )
1693 .await
1694 .expect("new_transaction failed");
1695 let object = root_directory
1696 .create_child_file(&mut transaction, "test")
1697 .await
1698 .expect("create_child_file failed");
1699 transaction.commit().await.expect("commit failed");
1700
1701 {
1702 let buf = object.allocate_buffer(1024).await;
1703 object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
1704 }
1705 std::mem::drop(object);
1706
1707 let mut transaction = root_directory
1708 .acquire_context_for_replace(None, "test", true)
1709 .await
1710 .expect("acquire_context_for_replace failed")
1711 .transaction;
1712 replace_child(&mut transaction, None, (&root_directory, "test"))
1713 .await
1714 .expect("replace_child failed");
1715 transaction.commit().await.expect("commit failed");
1716 }
1717 fs.close().await.expect("close failed");
1718 }
1719
1720 #[test_case(true; "test power fail with barriers")]
1721 #[test_case(false; "test power fail with checksums")]
1722 #[fuchsia::test]
1723 async fn test_power_fail(barriers_enabled: bool) {
1724 for _ in 0..10 {
1727 let (store_id, device, test_file_object_id) = {
1728 let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
1729 let fs = if barriers_enabled {
1730 FxFilesystemBuilder::new()
1731 .barriers_enabled(true)
1732 .format(true)
1733 .open(device)
1734 .await
1735 .expect("new filesystem failed")
1736 } else {
1737 FxFilesystem::new_empty(device).await.expect("new_empty failed")
1738 };
1739 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
1740
1741 fs.sync(SyncOptions { flush_device: true, ..SyncOptions::default() })
1742 .await
1743 .expect("sync failed");
1744
1745 let store = root_volume
1746 .new_volume(
1747 "test",
1748 NewChildStoreOptions {
1749 options: StoreOptions {
1750 crypt: Some(Arc::new(new_insecure_crypt())),
1751 ..StoreOptions::default()
1752 },
1753 ..Default::default()
1754 },
1755 )
1756 .await
1757 .expect("new_volume failed");
1758 let root_directory = Directory::open(&store, store.root_directory_object_id())
1759 .await
1760 .expect("open failed");
1761
1762 async fn create_files(store: &Arc<ObjectStore>, prefix: &str) {
1764 let fs = store.filesystem();
1765 let root_directory = Directory::open(store, store.root_directory_object_id())
1766 .await
1767 .expect("open failed");
1768 for i in 0..100 {
1769 let mut transaction = fs
1770 .root_store()
1771 .new_transaction(
1772 lock_keys![LockKey::object(
1773 store.store_object_id(),
1774 store.root_directory_object_id()
1775 )],
1776 Options::default(),
1777 )
1778 .await
1779 .expect("new_transaction failed");
1780 root_directory
1781 .create_child_file(&mut transaction, &format!("{prefix} {i}"))
1782 .await
1783 .expect("create_child_file failed");
1784 transaction.commit().await.expect("commit failed");
1785 }
1786 }
1787
1788 create_files(&store, "A").await;
1790
1791 let mut transaction = fs
1794 .root_store()
1795 .new_transaction(
1796 lock_keys![LockKey::object(
1797 store.store_object_id(),
1798 store.root_directory_object_id()
1799 )],
1800 Options::default(),
1801 )
1802 .await
1803 .expect("new_transaction failed");
1804 let object = root_directory
1805 .create_child_file(&mut transaction, "test")
1806 .await
1807 .expect("create_child_file failed");
1808 transaction.commit().await.expect("commit failed");
1809
1810 let mut transaction =
1811 object.new_transaction().await.expect("new_transaction failed");
1812 let mut buffer = object.allocate_buffer(4096).await;
1813 buffer.fill(0xed);
1814 object
1815 .txn_write(&mut transaction, 0, buffer.as_ref())
1816 .await
1817 .expect("txn_write failed");
1818 transaction.commit().await.expect("commit failed");
1819
1820 create_files(&store, "B").await;
1822
1823 fs.sync(SyncOptions::default()).await.expect("sync failed");
1826
1827 fasync::Timer::new(Duration::from_millis(10)).await;
1833
1834 (
1835 store.store_object_id(),
1836 fs.device().snapshot().expect("snapshot failed"),
1837 object.object_id(),
1838 )
1839 };
1840
1841 device
1844 .discard_random_since_last_flush()
1845 .expect("discard_random_since_last_flush failed");
1846
1847 let fs = FxFilesystem::open(device).await.expect("open failed");
1848 fsck(fs.clone()).await.expect("fsck failed");
1849
1850 let mut check_test_file = false;
1851
1852 let object_id = if fs.object_manager().store(store_id).is_some() {
1855 fsck_volume(&fs, store_id, Some(Arc::new(new_insecure_crypt())))
1856 .await
1857 .expect("fsck_volume failed");
1858
1859 let store = root_volume(fs.clone())
1863 .await
1864 .expect("root_volume failed")
1865 .volume(
1866 "test",
1867 StoreOptions {
1868 crypt: Some(Arc::new(new_insecure_crypt())),
1869 ..StoreOptions::default()
1870 },
1871 )
1872 .await
1873 .expect("volume failed");
1874
1875 let root_directory = Directory::open(&store, store.root_directory_object_id())
1876 .await
1877 .expect("open failed");
1878
1879 let mut transaction = fs
1880 .root_store()
1881 .new_transaction(
1882 lock_keys![LockKey::object(
1883 store.store_object_id(),
1884 store.root_directory_object_id()
1885 )],
1886 Options::default(),
1887 )
1888 .await
1889 .expect("new_transaction failed");
1890 let object = root_directory
1891 .create_child_file(&mut transaction, &format!("C"))
1892 .await
1893 .expect("create_child_file failed");
1894 transaction.commit().await.expect("commit failed");
1895
1896 if let Ok(test_file) = ObjectStore::open_object(
1898 &store,
1899 test_file_object_id,
1900 HandleOptions::default(),
1901 None,
1902 )
1903 .await
1904 {
1905 let mut buffer = test_file.allocate_buffer(4096).await;
1907 let bytes =
1908 test_file.read_aligned(0, buffer.as_mut()).await.expect("read failed");
1909 if bytes == 4096 {
1910 let expected = [0xed; 4096];
1911 assert_eq!(buffer.to_vec(), expected);
1912 } else {
1913 assert_eq!(bytes, 0);
1915 }
1916
1917 let mut transaction =
1919 test_file.new_transaction().await.expect("new_transaction failed");
1920 buffer.fill(0x37);
1921 test_file
1922 .txn_write(&mut transaction, 0, buffer.as_ref())
1923 .await
1924 .expect("txn_write failed");
1925 transaction.commit().await.expect("commit failed");
1926 check_test_file = true;
1927 }
1928
1929 object.object_id()
1930 } else {
1931 INVALID_OBJECT_ID
1932 };
1933
1934 fs.close().await.expect("close failed");
1936 let device = fs.take_device().await;
1937 device.reopen(false);
1938
1939 let fs = FxFilesystem::open(device).await.expect("open failed");
1940 fsck(fs.clone()).await.expect("fsck failed");
1941
1942 if object_id != INVALID_OBJECT_ID {
1945 fsck_volume(&fs, store_id, Some(Arc::new(new_insecure_crypt())))
1946 .await
1947 .expect("fsck_volume failed");
1948
1949 let store = root_volume(fs.clone())
1950 .await
1951 .expect("root_volume failed")
1952 .volume(
1953 "test",
1954 StoreOptions {
1955 crypt: Some(Arc::new(new_insecure_crypt())),
1956 ..StoreOptions::default()
1957 },
1958 )
1959 .await
1960 .expect("volume failed");
1961 ObjectStore::open_object(&store, object_id, HandleOptions::default(), None)
1963 .await
1964 .expect("open_object failed");
1965
1966 if check_test_file {
1968 info!("Checking test file for modification");
1969 let test_file = ObjectStore::open_object(
1970 &store,
1971 test_file_object_id,
1972 HandleOptions::default(),
1973 None,
1974 )
1975 .await
1976 .expect("open_object failed");
1977 let mut buffer = test_file.allocate_buffer(4096).await;
1978 assert_eq!(
1979 test_file.read_aligned(0, buffer.as_mut()).await.expect("read failed"),
1980 4096
1981 );
1982 let expected = [0x37; 4096];
1983 let data = buffer.to_vec();
1984 assert_eq!(data, expected);
1985 }
1986 }
1987
1988 fs.close().await.expect("close failed");
1989 }
1990 }
1991
1992 #[fuchsia::test]
1993 async fn test_barrier_not_emitted_when_transaction_has_no_data() {
1994 let barrier_count = Arc::new(AtomicU32::new(0));
1995
1996 struct Observer(Arc<AtomicU32>);
1997
1998 impl fake_device::Observer for Observer {
1999 fn barrier(&self) {
2000 self.0.fetch_add(1, Ordering::Relaxed);
2001 }
2002 }
2003
2004 let mut fake_device = FakeDevice::new(8192, 4096);
2005 fake_device.set_observer(Box::new(Observer(barrier_count.clone())));
2006 let device = DeviceHolder::new(fake_device);
2007 let fs = FxFilesystemBuilder::new()
2008 .barriers_enabled(true)
2009 .format(true)
2010 .open(device)
2011 .await
2012 .expect("new filesystem failed");
2013
2014 {
2015 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2016 root_vol
2017 .new_volume(
2018 "test",
2019 NewChildStoreOptions {
2020 options: StoreOptions {
2021 crypt: Some(Arc::new(new_insecure_crypt())),
2022 ..StoreOptions::default()
2023 },
2024 ..NewChildStoreOptions::default()
2025 },
2026 )
2027 .await
2028 .expect("there is no test volume");
2029 fs.close().await.expect("close failed");
2030 }
2031 let device = fs.take_device().await;
2034 device.reopen(false);
2035 let fs = FxFilesystemBuilder::new()
2036 .barriers_enabled(true)
2037 .open(device)
2038 .await
2039 .expect("new filesystem failed");
2040 let expected_barrier_count = barrier_count.load(Ordering::Relaxed);
2041
2042 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2043 let store = root_vol
2044 .volume(
2045 "test",
2046 StoreOptions {
2047 crypt: Some(Arc::new(new_insecure_crypt())),
2048 ..StoreOptions::default()
2049 },
2050 )
2051 .await
2052 .expect("there is no test volume");
2053
2054 let fs = store.filesystem();
2056 let root_directory =
2057 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
2058 for i in 0..100 {
2059 let mut transaction = fs
2060 .root_store()
2061 .new_transaction(
2062 lock_keys![LockKey::object(
2063 store.store_object_id(),
2064 store.root_directory_object_id()
2065 )],
2066 Options::default(),
2067 )
2068 .await
2069 .expect("new_transaction failed");
2070 root_directory
2071 .create_child_file(&mut transaction, &format!("A {i}"))
2072 .await
2073 .expect("create_child_file failed");
2074 transaction.commit().await.expect("commit failed");
2075 }
2076
2077 fs.close().await.expect("close failed");
2079 assert_eq!(expected_barrier_count, barrier_count.load(Ordering::Relaxed));
2081 }
2082
2083 #[fuchsia::test]
2084 async fn test_barrier_emitted_when_transaction_includes_data() {
2085 let barrier_count = Arc::new(AtomicU32::new(0));
2086
2087 struct Observer(Arc<AtomicU32>);
2088
2089 impl fake_device::Observer for Observer {
2090 fn barrier(&self) {
2091 self.0.fetch_add(1, Ordering::Relaxed);
2092 }
2093 }
2094
2095 let mut fake_device = FakeDevice::new(8192, 4096);
2096 fake_device.set_observer(Box::new(Observer(barrier_count.clone())));
2097 let device = DeviceHolder::new(fake_device);
2098 let fs = FxFilesystemBuilder::new()
2099 .barriers_enabled(true)
2100 .format(true)
2101 .open(device)
2102 .await
2103 .expect("new filesystem failed");
2104
2105 {
2106 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2107 root_vol
2108 .new_volume(
2109 "test",
2110 NewChildStoreOptions {
2111 options: StoreOptions {
2112 crypt: Some(Arc::new(new_insecure_crypt())),
2113 ..StoreOptions::default()
2114 },
2115 ..NewChildStoreOptions::default()
2116 },
2117 )
2118 .await
2119 .expect("there is no test volume");
2120 fs.close().await.expect("close failed");
2121 }
2122 let device = fs.take_device().await;
2125 device.reopen(false);
2126 let fs = FxFilesystemBuilder::new()
2127 .barriers_enabled(true)
2128 .open(device)
2129 .await
2130 .expect("new filesystem failed");
2131 let expected_barrier_count = barrier_count.load(Ordering::Relaxed);
2132
2133 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2134 let store = root_vol
2135 .volume(
2136 "test",
2137 StoreOptions {
2138 crypt: Some(Arc::new(new_insecure_crypt())),
2139 ..StoreOptions::default()
2140 },
2141 )
2142 .await
2143 .expect("there is no test volume");
2144
2145 let fs: Arc<FxFilesystem> = store.filesystem();
2147 let root_directory =
2148 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
2149
2150 let mut transaction = fs
2151 .root_store()
2152 .new_transaction(
2153 lock_keys![LockKey::object(
2154 store.store_object_id(),
2155 store.root_directory_object_id()
2156 )],
2157 Options::default(),
2158 )
2159 .await
2160 .expect("new_transaction failed");
2161 let object = root_directory
2162 .create_child_file(&mut transaction, "test")
2163 .await
2164 .expect("create_child_file failed");
2165 transaction.commit().await.expect("commit failed");
2166
2167 let mut transaction = object.new_transaction().await.expect("new_transaction failed");
2168 let mut buffer = object.allocate_buffer(4096).await;
2169 buffer.fill(0xed);
2170 object.txn_write(&mut transaction, 0, buffer.as_ref()).await.expect("txn_write failed");
2171 transaction.commit().await.expect("commit failed");
2172
2173 fs.close().await.expect("close failed");
2175 assert!(expected_barrier_count < barrier_count.load(Ordering::Relaxed));
2177 }
2178
2179 #[test_case(true; "fail when original filesystem has barriers enabled")]
2180 #[test_case(false; "fail when original filesystem has barriers disabled")]
2181 #[fuchsia::test]
2182 async fn test_switching_barrier_mode_on_existing_filesystem(original_barrier_mode: bool) {
2183 let crypt = Some(Arc::new(new_insecure_crypt()) as Arc<dyn fxfs_crypto::Crypt>);
2184 let fake_device = FakeDevice::new(8192, 4096);
2185 let device = DeviceHolder::new(fake_device);
2186 let fs: super::OpenFxFilesystem = FxFilesystemBuilder::new()
2187 .barriers_enabled(original_barrier_mode)
2188 .format(true)
2189 .open(device)
2190 .await
2191 .expect("new filesystem failed");
2192
2193 {
2195 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2196 let store = root_vol
2197 .new_volume(
2198 "test",
2199 NewChildStoreOptions {
2200 options: StoreOptions { crypt: crypt.clone(), ..Default::default() },
2201 ..Default::default()
2202 },
2203 )
2204 .await
2205 .expect("creating test volume");
2206 let root_dir = Directory::open(&store, store.root_directory_object_id())
2207 .await
2208 .expect("open failed");
2209 let mut transaction = fs
2210 .root_store()
2211 .new_transaction(
2212 lock_keys![LockKey::object(
2213 store.store_object_id(),
2214 store.root_directory_object_id()
2215 )],
2216 Default::default(),
2217 )
2218 .await
2219 .expect("new_transaction failed");
2220 let object = root_dir
2221 .create_child_file(&mut transaction, "file")
2222 .await
2223 .expect("create_child_file failed");
2224 transaction.commit().await.expect("commit failed");
2225 let mut buffer = object.allocate_buffer(4096).await;
2226 buffer.fill(0xA7);
2227 let new_size = object.write_or_append(None, buffer.as_ref()).await.unwrap();
2228 assert_eq!(new_size, 4096);
2229 }
2230
2231 fs.close().await.expect("close failed");
2233 let device = fs.take_device().await;
2234 device.reopen(false);
2235 let fs = FxFilesystemBuilder::new()
2236 .barriers_enabled(!original_barrier_mode)
2237 .open(device)
2238 .await
2239 .expect("new filesystem failed");
2240 {
2241 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2242 let store = root_vol
2243 .volume("test", StoreOptions { crypt: crypt.clone(), ..Default::default() })
2244 .await
2245 .expect("opening test volume");
2246 let root_dir = Directory::open(&store, store.root_directory_object_id())
2247 .await
2248 .expect("open failed");
2249 let (object_id, _, _) =
2250 root_dir.lookup("file").await.expect("lookup failed").expect("missing file");
2251 let test_file = ObjectStore::open_object(&store, object_id, Default::default(), None)
2252 .await
2253 .expect("open failed");
2254 let mut buffer = test_file.allocate_buffer(4096).await;
2256 buffer.fill(0xA8);
2257 let new_size = test_file.write_or_append(None, buffer.as_ref()).await.unwrap();
2258 assert_eq!(new_size, 8192);
2259 }
2260
2261 fs.close().await.expect("close failed");
2264 let device = fs.take_device().await;
2265 device.reopen(false);
2266 let fs = FxFilesystemBuilder::new()
2267 .barriers_enabled(original_barrier_mode)
2268 .open(device)
2269 .await
2270 .expect("new filesystem failed");
2271 {
2272 let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2273 let store = root_vol
2274 .volume("test", StoreOptions { crypt: crypt.clone(), ..Default::default() })
2275 .await
2276 .expect("opening test volume");
2277 let root_dir = Directory::open(&store, store.root_directory_object_id())
2278 .await
2279 .expect("open failed");
2280 let (object_id, _, _) =
2281 root_dir.lookup("file").await.expect("lookup failed").expect("missing file");
2282 let test_file = ObjectStore::open_object(&store, object_id, Default::default(), None)
2283 .await
2284 .expect("open failed");
2285 let mut buffer = test_file.allocate_buffer(8192).await;
2286 assert_eq!(
2287 test_file.read_aligned(0, buffer.as_mut()).await.expect("read failed"),
2288 8192,
2289 "short read"
2290 );
2291 let data = buffer.to_vec();
2292 assert_eq!(data[0..4096], [0xA7; 4096]);
2293 assert_eq!(data[4096..8192], [0xA8; 4096]);
2294 }
2295 fs.close().await.expect("close failed");
2296 }
2297
2298 #[fuchsia::test]
2299 async fn test_image_builder_mode_no_early_writes() {
2300 const BLOCK_SIZE: u32 = 4096;
2301 let device = DeviceHolder::new(FakeDevice::new(2048, BLOCK_SIZE));
2302 device.reopen(true);
2303 let fs = FxFilesystemBuilder::new()
2304 .format(true)
2305 .image_builder_mode(Some(SuperBlockInstance::A))
2306 .open(device)
2307 .await
2308 .expect("open failed");
2309 fs.enable_allocations();
2310 fs.device().reopen(false);
2312 fs.close().await.expect("closed");
2313 }
2314
2315 #[fuchsia::test]
2316 async fn test_image_builder_mode() {
2317 const BLOCK_SIZE: u32 = 4096;
2318 const EXISTING_FILE_RANGE: Range<u64> = 4096 * 1024..4096 * 1025;
2319 let device = DeviceHolder::new(FakeDevice::new(2048, BLOCK_SIZE));
2320
2321 {
2323 let mut write_buf =
2324 device.allocate_buffer(EXISTING_FILE_RANGE.length().unwrap() as usize).await;
2325 write_buf.fill(0xf0);
2326 device.write(EXISTING_FILE_RANGE.start, write_buf.as_ref()).await.expect("write");
2327 }
2328
2329 device.reopen(true);
2330
2331 let device = {
2332 let fs = FxFilesystemBuilder::new()
2333 .format(true)
2334 .image_builder_mode(Some(SuperBlockInstance::B))
2335 .open(device)
2336 .await
2337 .expect("open failed");
2338 fs.enable_allocations();
2339 {
2340 let root_store = fs.root_store();
2341 let root_directory =
2342 Directory::open(&root_store, root_store.root_directory_object_id())
2343 .await
2344 .expect("open failed");
2345 let handle;
2347 {
2348 let mut transaction = fs
2349 .root_store()
2350 .new_transaction(
2351 lock_keys![LockKey::object(
2352 root_directory.store().store_object_id(),
2353 root_directory.object_id()
2354 )],
2355 Options::default(),
2356 )
2357 .await
2358 .expect("new transaction");
2359 handle = root_directory
2360 .create_child_file(&mut transaction, "test")
2361 .await
2362 .expect("create file");
2363 handle.extend(&mut transaction, EXISTING_FILE_RANGE).await.expect("extend");
2364 transaction.commit().await.expect("commit");
2365 }
2366 }
2367 fs.device().reopen(false);
2368 fs.close().await.expect("close");
2369 fs.take_device().await
2370 };
2371 device.reopen(false);
2372 let fs = FxFilesystem::open(device).await.expect("open failed");
2373 fsck(fs.clone()).await.expect("fsck failed");
2374
2375 let root_store = fs.root_store();
2377 let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
2378 .await
2379 .expect("open failed");
2380 let (object_id, descriptor, _) =
2381 root_directory.lookup("test").await.expect("lookup failed").unwrap();
2382 assert_eq!(descriptor, ObjectDescriptor::File);
2383 let test_file =
2384 ObjectStore::open_object(&root_store, object_id, HandleOptions::default(), None)
2385 .await
2386 .expect("open failed");
2387 let mut read_buf =
2388 test_file.allocate_buffer(EXISTING_FILE_RANGE.length().unwrap() as usize).await;
2389 test_file.read_aligned(0, read_buf.as_mut()).await.expect("read failed");
2390 let data = read_buf.to_vec();
2391 assert_eq!(data, [0xf0; 4096]);
2392 fs.close().await.expect("closed");
2393 }
2394
2395 #[fuchsia::test]
2396 async fn test_read_only_mount_on_full_filesystem() {
2397 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
2398 let fs =
2399 FxFilesystemBuilder::new().format(true).open(device).await.expect("new_empty failed");
2400 let root_store = fs.root_store();
2401 let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
2402 .await
2403 .expect("open failed");
2404
2405 let mut transaction = fs
2406 .root_store()
2407 .new_transaction(
2408 lock_keys![LockKey::object(
2409 root_store.store_object_id(),
2410 root_directory.object_id()
2411 )],
2412 Options::default(),
2413 )
2414 .await
2415 .expect("new_transaction failed");
2416 let handle = root_directory
2417 .create_child_file(&mut transaction, "test")
2418 .await
2419 .expect("create_child_file failed");
2420 transaction.commit().await.expect("commit failed");
2421
2422 let mut buf = handle.allocate_buffer(4096).await;
2423 buf.fill(0xaa);
2424 loop {
2425 if handle.write_or_append(None, buf.as_ref()).await.is_err() {
2426 break;
2427 }
2428 }
2429
2430 let max_offset = fs.allocator().maximum_offset();
2431 fs.close().await.expect("Close failed");
2432
2433 let device = fs.take_device().await;
2434 device.reopen(false);
2435 let mut buffer = device
2436 .allocate_buffer(
2437 crate::round::round_up(max_offset, TEST_DEVICE_BLOCK_SIZE).unwrap() as usize
2438 )
2439 .await;
2440 device.read(0, buffer.as_mut()).await.expect("read failed");
2441
2442 let image_data = buffer.to_vec();
2443 let device = DeviceHolder::new(
2444 FakeDevice::from_image(image_data.as_slice(), TEST_DEVICE_BLOCK_SIZE)
2445 .expect("from_image failed"),
2446 );
2447 let fs =
2448 FxFilesystemBuilder::new().read_only(true).open(device).await.expect("open failed");
2449 fs.close().await.expect("Close failed");
2450 }
2451
2452 #[test_case(SuperBlockInstance::A; "Superblock instance A")]
2453 #[test_case(SuperBlockInstance::B; "Superblock instance B")]
2454 #[fuchsia::test]
2455 async fn test_image_builder_mode_flush_on_close_sb_a(target_sb: SuperBlockInstance) {
2456 const BLOCK_SIZE: u32 = 4096;
2457 let device = DeviceHolder::new(FakeDevice::new(2048, BLOCK_SIZE));
2458
2459 device.reopen(true);
2461 let fs = FxFilesystemBuilder::new()
2462 .format(true)
2463 .image_builder_mode(Some(target_sb))
2464 .open(device)
2465 .await
2466 .expect("open failed");
2467
2468 fs.enable_allocations();
2469
2470 fs.device().reopen(false);
2472
2473 {
2475 let root_store = fs.root_store();
2476 let root_directory =
2477 Directory::open(&root_store, root_store.root_directory_object_id())
2478 .await
2479 .expect("open failed");
2480
2481 let mut transaction = fs
2482 .root_store()
2483 .new_transaction(
2484 lock_keys![LockKey::object(
2485 root_directory.store().store_object_id(),
2486 root_directory.object_id()
2487 )],
2488 Options::default(),
2489 )
2490 .await
2491 .expect("new transaction");
2492 let handle = root_directory
2493 .create_child_file(&mut transaction, "post_finalize_file")
2494 .await
2495 .expect("create file");
2496 transaction.commit().await.expect("commit");
2497
2498 let mut buf = handle.allocate_buffer(BLOCK_SIZE as usize).await;
2499 buf.fill(0xaa);
2500 handle.write_or_append(None, buf.as_ref()).await.expect("write failed");
2501 }
2502
2503 fs.close().await.expect("close failed");
2505
2506 let other_sb = target_sb.next();
2507
2508 let device = fs.take_device().await;
2510 device.reopen(true); let mut buf = device.allocate_buffer(BLOCK_SIZE as usize).await;
2512
2513 device.read(target_sb.first_extent().start, buf.as_mut()).await.expect("read target_sb");
2514 let data = buf.to_vec();
2515 assert_eq!(&data[..8], b"FxfsSupr", "target_sb should have magic bytes");
2516
2517 buf.fill(0); device.read(other_sb.first_extent().start, buf.as_mut()).await.expect("read other_sb");
2519 let data2 = buf.to_vec();
2521 assert_eq!(data2, &[0; 4096], "other_sb should be zeroed");
2522 }
2523
2524 #[cfg(target_os = "fuchsia")]
2525 #[fuchsia::test(allow_stalls = false)]
2526 async fn test_trim_with_power_manager() {
2527 use anyhow::Error;
2528 use async_trait::async_trait;
2529 use fuchsia_async::TestExecutor;
2530 use futures::StreamExt;
2531
2532 TestExecutor::advance_to(fasync::MonotonicInstant::ZERO).await;
2533
2534 #[derive(Default)]
2535 struct MockPowerManager {
2536 on_battery: Mutex<bool>,
2537 event: event_listener::Event,
2538 wake_lease: Mutex<Option<zx::EventPair>>,
2539 }
2540
2541 impl MockPowerManager {
2542 fn set_on_battery(&self, v: bool) {
2543 *self.on_battery.lock() = v;
2544 self.event.notify(usize::MAX);
2545 }
2546
2547 fn is_lease_held(&self) -> bool {
2548 self.wake_lease.lock().as_ref().is_some_and(|handle| {
2549 handle
2550 .wait_one(
2551 zx::Signals::EVENTPAIR_PEER_CLOSED,
2552 zx::MonotonicInstant::INFINITE_PAST,
2553 )
2554 .is_err()
2555 })
2556 }
2557 }
2558
2559 impl super::PowerManager for MockPowerManager {
2560 fn watch_battery(
2561 self: Arc<Self>,
2562 ) -> futures::stream::BoxStream<'static, (bool, Option<super::WakeLease>)> {
2563 futures::stream::unfold(true, move |first| {
2564 let this = self.clone();
2565 async move {
2566 if !first {
2567 this.event.listen().await;
2568 }
2569 let val = *this.on_battery.lock();
2570 let handle = if val {
2571 None
2572 } else {
2573 let (h1, h2) = zx::EventPair::create();
2574 *this.wake_lease.lock() = Some(h2);
2575 Some(super::WakeLease::new(h1))
2576 };
2577 Some(((val, handle), false))
2578 }
2579 })
2580 .boxed()
2581 }
2582 }
2583
2584 let trim_count = Arc::new(AtomicU32::new(0));
2585
2586 struct TrimTrackingDevice {
2587 inner: DeviceHolder,
2588 trim_count: Arc<AtomicU32>,
2589 power_manager: Arc<MockPowerManager>,
2590 }
2591
2592 #[async_trait]
2593 impl storage_device::Device for TrimTrackingDevice {
2594 fn allocate_buffer(&self, size: usize) -> storage_device::buffer::BufferFuture<'_> {
2595 self.inner.allocate_buffer(size)
2596 }
2597 fn block_size(&self) -> u32 {
2598 self.inner.block_size()
2599 }
2600 fn block_count(&self) -> u64 {
2601 self.inner.block_count()
2602 }
2603 async fn read_with_opts(
2604 &self,
2605 offset: u64,
2606 buffer: storage_device::buffer::MutableBufferRef<'_>,
2607 opts: storage_device::ReadOptions,
2608 ) -> Result<(), Error> {
2609 self.inner.read_with_opts(offset, buffer, opts).await
2610 }
2611 async fn write_with_opts(
2612 &self,
2613 offset: u64,
2614 buffer: storage_device::buffer::BufferRef<'_>,
2615 opts: storage_device::WriteOptions,
2616 ) -> Result<(), Error> {
2617 self.inner.write_with_opts(offset, buffer, opts).await
2618 }
2619 async fn trim(&self, range: std::ops::Range<u64>) -> Result<(), Error> {
2620 assert!(self.power_manager.is_lease_held());
2621 self.trim_count.fetch_add(1, Ordering::SeqCst);
2622 self.inner.trim(range).await
2623 }
2624 async fn flush(&self) -> Result<(), Error> {
2625 self.inner.flush().await
2626 }
2627 async fn close(&self) -> Result<(), Error> {
2628 self.inner.close().await
2629 }
2630 fn supports_trim(&self) -> bool {
2631 true
2632 }
2633 fn is_read_only(&self) -> bool {
2634 self.inner.is_read_only()
2635 }
2636 fn snapshot(&self) -> Result<DeviceHolder, Error> {
2637 Ok(DeviceHolder::new(TrimTrackingDevice {
2638 inner: self.inner.snapshot()?,
2639 trim_count: self.trim_count.clone(),
2640 power_manager: self.power_manager.clone(),
2641 }))
2642 }
2643 fn reopen(&self, read_only: bool) {
2644 self.inner.reopen(read_only)
2645 }
2646 }
2647
2648 let pm = Arc::new(MockPowerManager::default());
2649
2650 pm.set_on_battery(true);
2652
2653 let fake_device = FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE);
2654 let device = DeviceHolder::new(TrimTrackingDevice {
2655 inner: DeviceHolder::new(fake_device),
2656 trim_count: trim_count.clone(),
2657 power_manager: pm.clone(),
2658 });
2659
2660 let fs = FxFilesystemBuilder::new()
2661 .format(true)
2662 .power_manager(pm.clone())
2663 .trim_config(Some((Duration::ZERO, Duration::from_millis(100))))
2664 .trim_charger_wait(Duration::from_millis(10))
2665 .open(device)
2666 .await
2667 .expect("open failed");
2668
2669 TestExecutor::advance_to(fasync::MonotonicInstant::after(
2671 Duration::from_millis(500).into(),
2672 ))
2673 .await;
2674 let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2675
2676 assert_eq!(trim_count.load(Ordering::SeqCst), 0);
2677
2678 {
2680 let root_store = fs.root_store();
2681 let root_directory =
2682 Directory::open(&root_store, root_store.root_directory_object_id())
2683 .await
2684 .expect("open failed");
2685 let mut transaction = fs
2686 .root_store()
2687 .new_transaction(
2688 lock_keys![LockKey::object(
2689 root_store.store_object_id(),
2690 root_directory.object_id()
2691 )],
2692 Options::default(),
2693 )
2694 .await
2695 .expect("new_transaction failed");
2696 let handle = root_directory
2697 .create_child_file(&mut transaction, "test")
2698 .await
2699 .expect("create_child_file failed");
2700 transaction.commit().await.expect("commit failed");
2701 handle.allocate(0..4096).await.expect("allocate failed");
2702 let mut transaction = fs
2704 .root_store()
2705 .new_transaction(
2706 lock_keys![
2707 LockKey::object(root_store.store_object_id(), root_directory.object_id()),
2708 LockKey::object(root_store.store_object_id(), handle.object_id()),
2709 ],
2710 Options::default(),
2711 )
2712 .await
2713 .expect("new_transaction failed");
2714 replace_child(&mut transaction, None, (&root_directory, "test"))
2715 .await
2716 .expect("delete failed");
2717 transaction.commit().await.expect("commit failed");
2718 fs.root_store()
2719 .tombstone_object(handle.object_id(), Options::default(), None)
2720 .await
2721 .expect("tombstone failed");
2722 }
2723
2724 pm.set_on_battery(false);
2726
2727 TestExecutor::advance_to(fasync::MonotonicInstant::after(Duration::from_millis(10).into()))
2729 .await;
2730
2731 let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2732
2733 assert!(trim_count.load(Ordering::SeqCst) > 0);
2734
2735 trim_count.store(0, Ordering::SeqCst);
2737 pm.set_on_battery(true);
2738
2739 TestExecutor::advance_to(fasync::MonotonicInstant::after(
2741 Duration::from_millis(500).into(),
2742 ))
2743 .await;
2744
2745 let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2746
2747 assert_eq!(trim_count.load(Ordering::SeqCst), 0);
2748
2749 fs.close().await.expect("close failed");
2750 }
2751
2752 #[fuchsia::test]
2753 async fn test_concurrent_do_trim_returns_error() {
2754 let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
2755 let fs = FxFilesystemBuilder::new()
2756 .trim_config(None)
2757 .format(true)
2758 .open(device)
2759 .await
2760 .expect("open failed");
2761
2762 let max_extent_size = fs.device().size() as usize;
2763 const EXTENTS_PER_BATCH: usize = usize::MAX;
2764
2765 let allocator = fs.allocator();
2767 let _trimmable_extents = allocator
2768 .take_for_trimming(0, max_extent_size, EXTENTS_PER_BATCH)
2769 .await
2770 .expect("take_for_trimming failed");
2771
2772 let res = fs.do_trim(None).await;
2775 assert!(matches!(res, Err(e) if FxfsError::AlreadyBound.matches(&e)));
2776
2777 fs.close().await.expect("close failed");
2778 }
2779}