1use crate::errors::FxfsError;
8use crate::filesystem::{FlushReason, ForceMajor};
9use crate::log::*;
10use crate::lsm_tree::types::{ItemRef, LayerIterator};
11use crate::lsm_tree::{LSMTree, layer_from_handle};
12use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, ReadObjectHandle};
13use crate::object_store::extent_record::ExtentValue;
14use crate::object_store::object_manager::{ObjectManager, ReservationUpdate};
15use crate::object_store::object_record::{ObjectKey, ObjectValue};
16use crate::object_store::transaction::{AssociatedObject, LockKey, Mutation, lock_keys};
17use crate::object_store::{
18 AssocObj, DirectWriter, EncryptedMutations, HandleOptions, LastObjectId, LastObjectIdInfo,
19 LockState, MAX_ENCRYPTED_MUTATIONS_SIZE, ObjectStore, Options, ReservedId, StoreInfo,
20 layer_size_from_encrypted_mutations_size, tree,
21};
22use crate::serialized_types::{LATEST_VERSION, Version, VersionedLatest};
23use anyhow::{Context, Error, anyhow};
24use fxfs_crypto::UnwrappedKey;
25use std::sync::OnceLock;
26use std::sync::atomic::Ordering;
27
28#[fxfs_trace::trace]
29impl ObjectStore {
30 pub async fn flush(&self) -> Result<Version, Error> {
33 self.flush_with_reason(FlushReason::Journal(ForceMajor::False)).await
34 }
35
36 pub async fn flush_with_reason(&self, reason: FlushReason) -> Result<Version, Error> {
38 let filesystem = self.filesystem();
39
40 let keys = lock_keys![LockKey::flush(self.store_object_id())];
41 let _guard = filesystem.lock_manager().write_lock(keys).await;
42
43 self.flush_guarded_with_reason(reason).await
44 }
45
46 #[trace("store_object_id" => self.store_object_id)]
48 pub async fn flush_guarded_with_reason(&self, reason: FlushReason) -> Result<Version, Error> {
49 if self.parent_store.is_none() {
50 return Ok(self.tree.get_earliest_version());
52 }
53
54 if matches!(*self.lock_state.lock(), LockState::Deleted) {
56 return Ok(LATEST_VERSION);
59 }
60
61 let filesystem = self.filesystem();
62 let object_manager = filesystem.object_manager();
63 let earliest_version = self.tree.get_earliest_version();
64 let needs_flush = object_manager.needs_flush(self.store_object_id);
65 if (reason == FlushReason::Journal(ForceMajor::False) && !needs_flush)
67 || (reason == FlushReason::UpgradeVersion && earliest_version == LATEST_VERSION)
68 || (reason == FlushReason::EncryptedMutations
69 && self.store_info().unwrap().encrypted_mutations_object_id == INVALID_OBJECT_ID)
70 {
71 return Ok(earliest_version);
74 }
75
76 let trace = self.trace.load(Ordering::Relaxed);
77 if trace {
78 info!(store_id = self.store_object_id(); "OS: begin flush");
79 }
80
81 if matches!(&*self.lock_state.lock(), LockState::Locked) {
82 self.flush_locked().await.with_context(|| {
83 format!("Failed to flush object store {}", self.store_object_id)
84 })?;
85 } else {
86 self.flush_unlocked(reason).await.with_context(|| {
87 format!("Failed to flush object store {}", self.store_object_id)
88 })?;
89 }
90
91 if trace {
92 info!(store_id = self.store_object_id(); "OS: end flush");
93 }
94 if let Some(callback) = &*self.flush_callback.lock() {
95 callback(self);
96 }
97
98 let mut counters = self.counters.lock();
101 counters.num_flushes += 1;
102 counters.last_flush_time = Some(std::time::SystemTime::now());
103 Ok(self.tree.get_earliest_version())
105 }
106
107 async fn flush_unlocked(&self, reason: FlushReason) -> Result<Vec<u64>, Error> {
109 struct StoreInfoSnapshot<'a> {
110 store: &'a ObjectStore,
111 store_info: OnceLock<StoreInfo>,
112 }
113 impl AssociatedObject for StoreInfoSnapshot<'_> {
114 fn will_apply_mutation(
115 &self,
116 _mutation: &Mutation,
117 _object_id: u64,
118 _manager: &ObjectManager,
119 ) {
120 let mut store_info = self.store.store_info().unwrap();
121
122 let lock_state = self.store.lock_state.lock();
123 if let LockState::Unlocked { mutations_cipher, .. } = &*lock_state {
124 store_info.mutations_cipher_offset = mutations_cipher.sequence_number();
125 }
126
127 self.store_info.set(store_info).unwrap();
128 }
129 }
130
131 let store_info_snapshot = StoreInfoSnapshot { store: self, store_info: OnceLock::new() };
132
133 let filesystem = self.filesystem();
134 let object_manager = filesystem.object_manager();
135 let reservation = object_manager.metadata_reservation();
136 let txn_options = Options {
137 skip_journal_checks: true,
138 borrow_metadata_space: true,
139 allocator_reservation: Some(reservation),
140 ..Default::default()
141 };
142
143 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
146 transaction.add_with_object(
147 self.store_object_id(),
148 Mutation::BeginFlush,
149 AssocObj::Borrowed(&store_info_snapshot),
150 );
151 transaction.commit().await.context("Failed to commit BeginFlush transaction")?;
152
153 let mut new_store_info = store_info_snapshot.store_info.into_inner().unwrap();
154
155 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
160
161 let parent_store = self.parent_store.as_ref().unwrap();
163 let handle_options = HandleOptions { skip_journal_checks: true, ..Default::default() };
164 let id_and_key = {
165 let mut lock_state = self.lock_state.lock();
166 match &mut *lock_state {
167 LockState::Unlocked { cached_keys, .. } => {
168 if let Some(item) = cached_keys.pop() {
169 Ok(Some(item))
170 } else {
171 log::warn!("No cached keys for flush for store {}", self.store_object_id());
172 Err(anyhow!(FxfsError::Internal).context("No cached keys for flush"))
173 }
174 }
175 LockState::UnlockedReadOnly(..) => {
176 Err(anyhow!(FxfsError::Internal).context("Flush on read-only store"))
177 }
178 LockState::Unencrypted => Ok(None),
179 _ => Err(anyhow!(FxfsError::Internal))
180 .with_context(|| format!("Invalid lock state ({:?}) for flush", *lock_state)),
181 }
182 }?;
183
184 let (new_object_tree_layer, unwrapped_key) = if let Some((raw_id, key, unwrapped_key)) =
185 id_and_key
186 {
187 let object_id = ReservedId::new(parent_store, raw_id);
188 let handle = ObjectStore::create_object_with_key(
189 parent_store,
190 &mut transaction,
191 object_id,
192 handle_options,
193 key,
194 UnwrappedKey::new(unwrapped_key.clone()),
195 )
196 .await?;
197 (handle, Some(unwrapped_key))
198 } else {
199 let handle =
200 ObjectStore::create_object(parent_store, &mut transaction, handle_options, None)
201 .await?;
202 (handle, None)
203 };
204 let writer = DirectWriter::new(&new_object_tree_layer, txn_options).await;
205 let new_object_tree_layer_object_id = new_object_tree_layer.object_id();
206 parent_store.add_to_graveyard(&mut transaction, new_object_tree_layer_object_id);
207
208 transaction.commit().await.context("Failed to commit create layer transaction")?;
209
210 let full_compaction = match reason {
212 FlushReason::UpgradeVersion | FlushReason::Journal(ForceMajor::True) => true,
213 _ => false,
214 };
215 let (layers_to_keep, old_layers) = tree::flush(
216 &self.tree,
217 writer,
218 matches!(reason, FlushReason::Journal(_))
219 .then(|| filesystem.journal().get_compaction_yielder()),
220 full_compaction,
221 )
222 .await
223 .context("Failed to flush tree")?;
224
225 let mut new_layers = vec![layer_from_handle(new_object_tree_layer, unwrapped_key).await?];
227 new_layers.extend(layers_to_keep.iter().map(|l| (*l).clone()));
228
229 new_store_info.layers = Vec::new();
230 for layer in &new_layers {
231 if let Some(handle) = layer.handle() {
232 new_store_info.layers.push(handle.object_id());
233 }
234 }
235
236 let reservation_update: ReservationUpdate; let mut end_transaction = parent_store
238 .new_transaction(
239 lock_keys![LockKey::object(
240 self.parent_store.as_ref().unwrap().store_object_id(),
241 self.store_info_handle_object_id().unwrap(),
242 )],
243 txn_options,
244 )
245 .await?;
246
247 parent_store.remove_from_graveyard(&mut end_transaction, new_object_tree_layer_object_id);
248
249 for layer in &old_layers {
251 if let Some(handle) = layer.handle() {
252 parent_store.add_to_graveyard(&mut end_transaction, handle.object_id());
253 }
254 }
255
256 let old_encrypted_mutations_object_id =
257 std::mem::replace(&mut new_store_info.encrypted_mutations_object_id, INVALID_OBJECT_ID);
258 if old_encrypted_mutations_object_id != INVALID_OBJECT_ID {
259 parent_store.add_to_graveyard(&mut end_transaction, old_encrypted_mutations_object_id);
260 }
261
262 match &mut new_store_info.last_object_id {
270 LastObjectIdInfo::Unencrypted { id } => {
271 let LastObjectId::Unencrypted { id: in_memory_value } =
272 &*self.last_object_id.lock()
273 else {
274 unreachable!()
275 };
276 *id = *in_memory_value;
277 }
278 LastObjectIdInfo::Encrypted { id, key } => {
279 let LastObjectId::Encrypted { id: in_memory_value, .. } =
280 &*self.last_object_id.lock()
281 else {
282 unreachable!()
283 };
284 *id = *in_memory_value;
285 let guard = self.store_info.lock();
286 let current_store_info = guard.as_ref().unwrap();
287 let LastObjectIdInfo::Encrypted { key: in_memory_value, .. } =
288 ¤t_store_info.last_object_id
289 else {
290 unreachable!()
291 };
292 *key = in_memory_value.clone();
293 }
294 LastObjectIdInfo::Low32Bit => {}
295 }
296
297 self.write_store_info(&mut end_transaction, &new_store_info).await?;
298
299 let layer_file_sizes = new_layers
300 .iter()
301 .map(|l| l.handle().map(ReadObjectHandle::get_size).unwrap_or(0))
302 .collect::<Vec<u64>>();
303
304 let total_layer_size = layer_file_sizes.iter().sum();
305 reservation_update =
306 ReservationUpdate::new(tree::reservation_amount_from_layer_size(total_layer_size));
307
308 end_transaction.add_with_object(
309 self.store_object_id(),
310 Mutation::EndFlush,
311 AssocObj::Borrowed(&reservation_update),
312 );
313
314 if self.trace.load(Ordering::Relaxed) {
315 info!(
316 store_id = self.store_object_id(),
317 old_layer_count = old_layers.len(),
318 new_layer_count = new_layers.len(),
319 total_layer_size,
320 new_store_info:?;
321 "OS: compacting"
322 );
323 }
324
325 end_transaction
326 .commit_with_callback(|_| {
327 let mut store_info = self.store_info.lock();
328 let info = store_info.as_mut().unwrap();
329 info.layers = new_store_info.layers;
330 info.encrypted_mutations_object_id = new_store_info.encrypted_mutations_object_id;
331 info.mutations_cipher_offset = new_store_info.mutations_cipher_offset;
332 self.tree.set_layers(new_layers);
333 })
334 .await
335 .context("Failed to commit EndFlush transaction")?;
336
337 for layer in old_layers {
339 let object_id = layer.handle().map(|h| h.object_id());
340 layer.close_layer().await;
341 if let Some(object_id) = object_id {
342 parent_store
343 .tombstone_object(object_id, txn_options, None)
344 .await
345 .context("Failed to tombstone old layer")?;
346 }
347 }
348
349 if old_encrypted_mutations_object_id != INVALID_OBJECT_ID {
350 parent_store
351 .tombstone_object(old_encrypted_mutations_object_id, txn_options, None)
352 .await
353 .context("Failed to tombstone old encrypted mutations")?;
354 }
355
356 Ok(layer_file_sizes)
357 }
358
359 async fn flush_locked(&self) -> Result<(), Error> {
361 let filesystem = self.filesystem();
362 let object_manager = filesystem.object_manager();
363 let reservation = object_manager.metadata_reservation();
364 let txn_options = Options {
365 skip_journal_checks: true,
366 borrow_metadata_space: true,
367 allocator_reservation: Some(reservation),
368 ..Default::default()
369 };
370
371 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
372 transaction.add(self.store_object_id(), Mutation::BeginFlush);
373 transaction.commit().await.context("Failed to commit BeginFlush transaction")?;
374
375 let mut new_store_info = self.load_store_info().await?;
376
377 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
382
383 let reservation_update: ReservationUpdate; let handle; let mut end_transaction;
386
387 let parent_store = self.parent_store.as_ref().unwrap();
390 handle = if new_store_info.encrypted_mutations_object_id == INVALID_OBJECT_ID {
391 let handle = ObjectStore::create_object(
392 parent_store,
393 &mut transaction,
394 HandleOptions { skip_journal_checks: true, ..Default::default() },
395 None,
396 )
397 .await?;
398 let oid = handle.object_id();
399 end_transaction = parent_store
400 .new_transaction(
401 lock_keys![
402 LockKey::object(parent_store.store_object_id(), oid),
403 LockKey::object(
404 parent_store.store_object_id(),
405 self.store_info_handle_object_id().unwrap(),
406 ),
407 ],
408 txn_options,
409 )
410 .await?;
411 new_store_info.encrypted_mutations_object_id = oid;
412 parent_store.add_to_graveyard(&mut transaction, oid);
413 parent_store.remove_from_graveyard(&mut end_transaction, oid);
414 handle
415 } else {
416 end_transaction = parent_store
417 .new_transaction(
418 lock_keys![
419 LockKey::object(
420 parent_store.store_object_id(),
421 new_store_info.encrypted_mutations_object_id,
422 ),
423 LockKey::object(
424 parent_store.store_object_id(),
425 self.store_info_handle_object_id().unwrap(),
426 ),
427 ],
428 txn_options,
429 )
430 .await?;
431 ObjectStore::open_object(
432 parent_store,
433 new_store_info.encrypted_mutations_object_id,
434 HandleOptions { skip_journal_checks: true, ..Default::default() },
435 None,
436 )
437 .await?
438 };
439 transaction
440 .commit()
441 .await
442 .context("Failed to commit create encrypted mutations transaction")?;
443
444 let journaled = filesystem
447 .journal()
448 .read_transactions_for_object(self.store_object_id)
449 .await
450 .context("Failed to read encrypted mutations from journal")?;
451 let mut buffer = handle.allocate_buffer(MAX_ENCRYPTED_MUTATIONS_SIZE).await;
452 let mut writer = buffer.writer();
453 EncryptedMutations::from_replayed_mutations(self.store_object_id, journaled)
454 .serialize_with_version(&mut writer)?;
455 let len = writer.position();
456 handle
457 .txn_write(&mut end_transaction, handle.get_size(), buffer.subslice(..len))
458 .await
459 .context("Failed to write encrypted mutations")?;
460
461 self.write_store_info(&mut end_transaction, &new_store_info)
462 .await
463 .context("Failed to write store info")?;
464
465 let mut total_layer_size = 0;
466 for &oid in &new_store_info.layers {
467 total_layer_size += parent_store.get_file_size(oid).await?;
468 }
469 total_layer_size +=
470 layer_size_from_encrypted_mutations_size(handle.get_size() + len as u64);
471
472 reservation_update =
473 ReservationUpdate::new(tree::reservation_amount_from_layer_size(total_layer_size));
474
475 end_transaction.add_with_object(
476 self.store_object_id(),
477 Mutation::EndFlush,
478 AssocObj::Borrowed(&reservation_update),
479 );
480
481 end_transaction.commit().await.context("Failed to commit EndFlush transaction")?;
482
483 Ok(())
484 }
485}
486
487#[cfg(test)]
488mod tests {
489 use super::{FlushReason, ForceMajor};
490 use crate::filesystem::{FxFilesystem, FxFilesystemBuilder};
491 use crate::object_handle::{
492 INVALID_OBJECT_ID, ObjectHandle, ReadObjectHandle, WriteObjectHandle,
493 };
494 use crate::object_store::directory::Directory;
495 use crate::object_store::journal::JournalOptions;
496 use crate::object_store::transaction::{Options, lock_keys};
497 use crate::object_store::volume::root_volume;
498 use crate::object_store::{
499 HandleOptions, LockKey, NewChildStoreOptions, ObjectStore, StoreOptions,
500 layer_size_from_encrypted_mutations_size, tree,
501 };
502 use fxfs_insecure_crypto::new_insecure_crypt;
503 use std::sync::Arc;
504 use storage_device::DeviceHolder;
505 use storage_device::fake_device::FakeDevice;
506
507 #[fuchsia::test]
508 async fn test_flush_when_locked() {
509 let device = DeviceHolder::new(FakeDevice::new(8192, 1024));
510 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
511 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
512 let crypt = Arc::new(new_insecure_crypt());
513 let store = root_volume
514 .new_volume(
515 "test",
516 NewChildStoreOptions {
517 options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
518 ..NewChildStoreOptions::default()
519 },
520 )
521 .await
522 .expect("new_volume failed");
523 let root_dir =
524 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
525 let mut transaction = fs
526 .root_store()
527 .new_transaction(
528 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
529 Options::default(),
530 )
531 .await
532 .expect("new_transaction failed");
533 let foo = root_dir
534 .create_child_file(&mut transaction, "foo")
535 .await
536 .expect("create_child_file failed");
537 transaction.commit().await.expect("commit failed");
538
539 store.flush().await.expect("flush failed");
543
544 let mut transaction = fs
545 .root_store()
546 .new_transaction(
547 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
548 Options::default(),
549 )
550 .await
551 .expect("new_transaction failed");
552 let bar = root_dir
553 .create_child_file(&mut transaction, "bar")
554 .await
555 .expect("create_child_file failed");
556 transaction.commit().await.expect("commit failed");
557
558 store.lock().await.expect("lock failed");
559
560 store.flush().await.expect("flush failed");
562
563 let info = store.load_store_info().await.unwrap();
565 let parent_store = store.parent_store().unwrap();
566 let mut total_layer_size = 0;
567 for &oid in &info.layers {
568 total_layer_size +=
569 parent_store.get_file_size(oid).await.expect("get_file_size failed");
570 }
571 assert_ne!(info.encrypted_mutations_object_id, INVALID_OBJECT_ID);
572 total_layer_size += layer_size_from_encrypted_mutations_size(
573 parent_store
574 .get_file_size(info.encrypted_mutations_object_id)
575 .await
576 .expect("get_file_size failed"),
577 );
578 assert_eq!(
579 fs.object_manager().reservation(store.store_object_id()),
580 Some(tree::reservation_amount_from_layer_size(total_layer_size))
581 );
582
583 store.unlock(crypt).await.expect("unlock failed");
585
586 ObjectStore::open_object(&store, foo.object_id(), HandleOptions::default(), None)
587 .await
588 .expect("open_object failed");
589
590 ObjectStore::open_object(&store, bar.object_id(), HandleOptions::default(), None)
591 .await
592 .expect("open_object failed");
593
594 fs.close().await.expect("close failed");
595 }
596
597 #[fuchsia::test]
598 async fn test_major_compaction_frees_reservation_and_merges_layers() {
599 let device = DeviceHolder::new(FakeDevice::new(32768, 512));
600 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
601 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
602 let store = root_volume
603 .new_volume("test", NewChildStoreOptions::default())
604 .await
605 .expect("new_volume failed");
606 let root_dir =
607 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
608
609 let num_files = 2000;
612 let mut file_ids = Vec::with_capacity(num_files);
613 for i in 0..num_files {
614 let filename = format!("file_{:04}_{:<200}", i, i);
615 let mut transaction = store
616 .new_transaction(
617 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
618 Options::default(),
619 )
620 .await
621 .expect("new_transaction failed");
622 let file = root_dir
623 .create_child_file(&mut transaction, &filename)
624 .await
625 .expect("create_child_file failed");
626 transaction.commit().await.expect("commit failed");
627 file_ids.push((file.object_id(), filename));
628 }
629
630 store.flush().await.expect("flush failed");
632
633 let info_initial = store.load_store_info().await.expect("load_store_info failed");
634 assert_eq!(info_initial.layers.len(), 1);
635 let initial_reservation = fs
636 .object_manager()
637 .reservation(store.store_object_id())
638 .expect("reservation not found");
639
640 for (oid, filename) in &file_ids[0..1000] {
642 let mut transaction = store
643 .new_transaction(
644 lock_keys![
645 LockKey::object(store.store_object_id(), root_dir.object_id()),
646 LockKey::object(store.store_object_id(), *oid),
647 ],
648 Options::default(),
649 )
650 .await
651 .expect("new_transaction failed");
652 let replaced = crate::object_store::directory::replace_child(
653 &mut transaction,
654 None,
655 (&root_dir, filename.as_str()),
656 )
657 .await
658 .expect("replace_child failed");
659 assert_matches::assert_matches!(
660 replaced,
661 crate::object_store::directory::ReplacedChild::Object(id) if id == *oid
662 );
663 transaction.commit().await.expect("commit failed");
664 }
665
666 store
669 .flush_with_reason(FlushReason::Journal(ForceMajor::False))
670 .await
671 .expect("minor flush failed");
672
673 let info_after_minor = store.load_store_info().await.expect("load_store_info failed");
674 assert_eq!(
675 info_after_minor.layers.len(),
676 2,
677 "Minor compaction should keep the base layer, creating 2 layers"
678 );
679 let res_after_minor = fs
680 .object_manager()
681 .reservation(store.store_object_id())
682 .expect("reservation not found");
683 assert!(
684 res_after_minor >= initial_reservation,
685 "Minor compaction should not decrease reservation because base layer is kept"
686 );
687
688 store
691 .flush_with_reason(FlushReason::Journal(ForceMajor::True))
692 .await
693 .expect("major flush failed");
694
695 let info_after_major = store.load_store_info().await.expect("load_store_info failed");
696 assert_eq!(
697 info_after_major.layers.len(),
698 1,
699 "Major compaction should merge all layers into 1"
700 );
701 let res_after_major = fs
702 .object_manager()
703 .reservation(store.store_object_id())
704 .expect("reservation not found");
705 assert!(
706 res_after_major < initial_reservation,
707 "Major compaction should decrease reservation (was {}, initial was {})",
708 res_after_major,
709 initial_reservation
710 );
711
712 for (oid, filename) in &file_ids[0..1000] {
715 assert!(
716 root_dir.lookup(filename).await.expect("lookup failed").is_none(),
717 "Deleted file {} should not exist",
718 oid
719 );
720 }
721 for (oid, filename) in &file_ids[1000..num_files] {
722 let res = root_dir.lookup(filename).await.expect("lookup failed");
723 assert!(res.is_some(), "Surviving file {} should exist", oid);
724 assert_eq!(res.unwrap().0, *oid);
725 }
726
727 fs.close().await.expect("close failed");
728 }
729
730 #[fuchsia::test]
731 async fn test_major_compaction_without_journal_mutations() {
732 let device = DeviceHolder::new(FakeDevice::new(32768, 512));
733 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
734 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
735 let store = root_volume
736 .new_volume("test", NewChildStoreOptions::default())
737 .await
738 .expect("new_volume failed");
739 let root_dir =
740 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
741
742 for i in 0..2000 {
744 let filename = format!("file_{:04}_{:<200}", i, i);
745 let mut transaction = store
746 .new_transaction(
747 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
748 Options::default(),
749 )
750 .await
751 .expect("new_transaction failed");
752 root_dir.create_child_file(&mut transaction, &filename).await.expect("create failed");
753 transaction.commit().await.expect("commit failed");
754 }
755 store.flush().await.expect("flush failed");
756
757 let mut transaction = store
759 .new_transaction(
760 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
761 Options::default(),
762 )
763 .await
764 .expect("new_transaction failed");
765 root_dir.create_child_file(&mut transaction, "extra_file").await.expect("create failed");
766 transaction.commit().await.expect("commit failed");
767 store
768 .flush_with_reason(FlushReason::Journal(ForceMajor::False))
769 .await
770 .expect("flush failed");
771
772 let info = store.load_store_info().await.expect("load_store_info failed");
773 assert_eq!(info.layers.len(), 2, "Should have 2 layers before major flush");
774
775 assert!(!fs.object_manager().needs_flush(store.store_object_id()));
777 store
778 .flush_with_reason(FlushReason::Journal(ForceMajor::True))
779 .await
780 .expect("major flush failed");
781
782 let info_after_major = store.load_store_info().await.expect("load_store_info failed");
783 assert_eq!(
784 info_after_major.layers.len(),
785 1,
786 "Major flush should merge 2 layers into 1 even without journal mutations"
787 );
788
789 assert!(!fs.object_manager().needs_flush(store.store_object_id()));
792 store
793 .flush_with_reason(FlushReason::Journal(ForceMajor::True))
794 .await
795 .expect("second major flush failed");
796 let info_after_second = store.load_store_info().await.expect("load_store_info failed");
797 assert_eq!(
798 info_after_second.layers.len(),
799 1,
800 "Major flush when at 1 layer should succeed and leave 1 layer"
801 );
802
803 fs.close().await.expect("close failed");
804 }
805
806 #[fuchsia::test]
807 async fn test_major_compaction_purges_deleted_extents() {
808 let device = DeviceHolder::new(FakeDevice::new(32768, 512));
809 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
810 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
811 let store = root_volume
812 .new_volume("test", NewChildStoreOptions::default())
813 .await
814 .expect("new_volume failed");
815 let root_dir =
816 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
817
818 for i in 0..2000 {
820 let filename = format!("file_{:04}_{:<200}", i, i);
821 let mut transaction = store
822 .new_transaction(
823 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
824 Options::default(),
825 )
826 .await
827 .expect("new_transaction failed");
828 root_dir.create_child_file(&mut transaction, &filename).await.expect("create failed");
829 transaction.commit().await.expect("commit failed");
830 }
831
832 let data_oid = {
834 let mut transaction = store
835 .new_transaction(
836 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
837 Options::default(),
838 )
839 .await
840 .expect("new_transaction failed");
841 let file = root_dir
842 .create_child_file(&mut transaction, "data_file")
843 .await
844 .expect("create_child_file failed");
845 let oid = file.object_id();
846 transaction.commit().await.expect("commit failed");
847 oid
848 };
849
850 let data_file = ObjectStore::open_object(&store, data_oid, HandleOptions::default(), None)
851 .await
852 .expect("open_object failed");
853 {
854 let mut transaction = store
855 .new_transaction(
856 lock_keys![LockKey::object(store.store_object_id(), data_oid)],
857 Options::default(),
858 )
859 .await
860 .expect("new_transaction failed");
861 let mut buffer = data_file.allocate_buffer(131072).await;
862 buffer.fill(0xAB);
863 data_file
864 .txn_write(&mut transaction, 0, buffer.as_ref())
865 .await
866 .expect("txn_write failed");
867 transaction.commit().await.expect("commit failed");
868 }
869
870 store.flush().await.expect("flush failed");
872
873 let info0 = store.load_store_info().await.expect("load_store_info failed");
874 assert_eq!(info0.layers.len(), 1);
875
876 data_file.truncate(0).await.expect("truncate failed");
878
879 store
881 .flush_with_reason(FlushReason::Journal(ForceMajor::False))
882 .await
883 .expect("minor flush failed");
884 let info1 = store.load_store_info().await.expect("load_store_info failed");
885 assert_eq!(info1.layers.len(), 2, "Minor flush should keep Layer 0 and create Layer 1");
886
887 store
889 .flush_with_reason(FlushReason::Journal(ForceMajor::True))
890 .await
891 .expect("major flush failed");
892 let info2 = store.load_store_info().await.expect("load_store_info failed");
893 assert_eq!(info2.layers.len(), 1, "Major flush should merge into 1 layer");
894
895 assert_eq!(data_file.get_size(), 0);
897
898 fs.close().await.expect("close failed");
899 }
900
901 #[fuchsia::test]
907 async fn test_borrow_metadata_space_fails_without_major_compaction() {
908 let reclaim_size = 65536;
909 let device = DeviceHolder::new(FakeDevice::new(102400, 512));
911 let fs = FxFilesystemBuilder::new()
912 .journal_options(JournalOptions { reclaim_size, ..Default::default() })
913 .format(true)
914 .open(device)
915 .await
916 .expect("open failed");
917
918 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
923 let num_stores = 3;
924 let files_per_store = 3500;
925 let mut stores_and_files = Vec::new();
926
927 for s in 0..num_stores {
928 let store = root_volume
929 .new_volume(&format!("test_{s}"), NewChildStoreOptions::default())
930 .await
931 .expect("new_volume failed");
932 let root_dir = Directory::open(&store, store.root_directory_object_id())
933 .await
934 .expect("open failed");
935 let mut file_ids = Vec::with_capacity(files_per_store);
936 for i in 0..files_per_store {
937 let filename = format!("file_{:04}_{:<200}", i, i);
938 let mut transaction = store
939 .new_transaction(
940 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
941 Options::default(),
942 )
943 .await
944 .expect("new_transaction failed");
945 let file = root_dir
946 .create_child_file(&mut transaction, &filename)
947 .await
948 .expect("create_child_file failed");
949 transaction.commit().await.expect("commit failed");
950 file_ids.push((file.object_id(), filename));
951 }
952 store
953 .flush_with_reason(FlushReason::Journal(ForceMajor::True))
954 .await
955 .expect("flush failed");
956 stores_and_files.push((store, root_dir, file_ids));
957 }
958
959 for (_, root_dir, file_ids) in stores_and_files {
960 for (oid, filename) in &file_ids[0..2500] {
961 let mut context = root_dir
962 .acquire_context_for_replace(None, filename.as_str(), true)
963 .await
964 .expect("acquire_context failed");
965 let replaced = crate::object_store::directory::replace_child(
966 &mut context.transaction,
967 None,
968 (&root_dir, filename.as_str()),
969 )
970 .await
971 .expect("replace_child failed");
972 assert_matches::assert_matches!(
973 replaced,
974 crate::object_store::directory::ReplacedChild::Object(id) if id == *oid
975 );
976 context.transaction.commit().await.expect("commit failed");
977 }
978 }
979
980 fs.close().await.expect("close failed");
981 }
982}
983
984impl tree::MajorCompactable<ObjectKey, ObjectValue> for LSMTree<ObjectKey, ObjectValue> {
985 async fn major_iter(
986 iter: impl LayerIterator<ObjectKey, ObjectValue>,
987 ) -> Result<impl LayerIterator<ObjectKey, ObjectValue>, Error> {
988 iter.filter(|item: ItemRef<'_, _, _>| match item {
989 ItemRef { value: ObjectValue::None, .. } => false,
991 ItemRef { value: ObjectValue::Extent(ExtentValue::None), .. } => false,
993 _ => true,
994 })
995 .await
996 }
997}