1use crate::errors::FxfsError;
8use crate::log::*;
9use crate::lsm_tree::types::{ItemRef, LayerIterator};
10use crate::lsm_tree::{LSMTree, layers_from_handles};
11use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, ReadObjectHandle};
12use crate::object_store::extent_record::ExtentValue;
13use crate::object_store::object_manager::{ObjectManager, ReservationUpdate};
14use crate::object_store::object_record::{ObjectKey, ObjectValue};
15use crate::object_store::transaction::{AssociatedObject, LockKey, Mutation, lock_keys};
16use crate::object_store::{
17 AssocObj, DirectWriter, EncryptedMutations, HandleOptions, LastObjectId, LastObjectIdInfo,
18 LockState, MAX_ENCRYPTED_MUTATIONS_SIZE, ObjectStore, Options, ReservedId, StoreInfo,
19 layer_size_from_encrypted_mutations_size, tree,
20};
21use crate::serialized_types::{LATEST_VERSION, Version, VersionedLatest};
22use anyhow::{Context, Error, anyhow};
23use std::sync::OnceLock;
24use std::sync::atomic::Ordering;
25
26#[derive(Copy, Clone, Debug, PartialEq, Eq)]
27pub enum Reason {
28 Journal,
30
31 EncryptedMutations,
33
34 UpgradeVersion,
36}
37
38#[fxfs_trace::trace]
39impl ObjectStore {
40 pub async fn flush_with_reason(&self, reason: Reason) -> Result<Version, Error> {
42 let filesystem = self.filesystem();
43
44 let keys = lock_keys![LockKey::flush(self.store_object_id())];
45 let _guard = filesystem.lock_manager().write_lock(keys).await;
46
47 self.flush_guarded_with_reason(reason).await
48 }
49
50 #[trace("store_object_id" => self.store_object_id)]
52 pub async fn flush_guarded_with_reason(&self, reason: Reason) -> Result<Version, Error> {
53 if self.parent_store.is_none() {
54 return Ok(self.tree.get_earliest_version());
56 }
57
58 if matches!(*self.lock_state.lock(), LockState::Deleted) {
60 return Ok(LATEST_VERSION);
63 }
64
65 let filesystem = self.filesystem();
66 let object_manager = filesystem.object_manager();
67 let earliest_version = self.tree.get_earliest_version();
68 if (reason == Reason::Journal && !object_manager.needs_flush(self.store_object_id))
70 || (reason == Reason::UpgradeVersion && earliest_version == LATEST_VERSION)
71 || (reason == Reason::EncryptedMutations
72 && self.store_info().unwrap().encrypted_mutations_object_id == INVALID_OBJECT_ID)
73 {
74 return Ok(earliest_version);
77 }
78
79 let trace = self.trace.load(Ordering::Relaxed);
80 if trace {
81 info!(store_id = self.store_object_id(); "OS: begin flush");
82 }
83
84 if matches!(&*self.lock_state.lock(), LockState::Locked) {
85 self.flush_locked().await.with_context(|| {
86 format!("Failed to flush object store {}", self.store_object_id)
87 })?;
88 } else {
89 self.flush_unlocked(reason).await.with_context(|| {
90 format!("Failed to flush object store {}", self.store_object_id)
91 })?;
92 }
93
94 if trace {
95 info!(store_id = self.store_object_id(); "OS: end flush");
96 }
97 if let Some(callback) = &*self.flush_callback.lock() {
98 callback(self);
99 }
100
101 let mut counters = self.counters.lock();
104 counters.num_flushes += 1;
105 counters.last_flush_time = Some(std::time::SystemTime::now());
106 Ok(self.tree.get_earliest_version())
108 }
109
110 async fn flush_unlocked(&self, reason: Reason) -> Result<Vec<u64>, Error> {
112 struct StoreInfoSnapshot<'a> {
113 store: &'a ObjectStore,
114 store_info: OnceLock<StoreInfo>,
115 }
116 impl AssociatedObject for StoreInfoSnapshot<'_> {
117 fn will_apply_mutation(
118 &self,
119 _mutation: &Mutation,
120 _object_id: u64,
121 _manager: &ObjectManager,
122 ) {
123 let mut store_info = self.store.store_info().unwrap();
124
125 let mutations_cipher = self.store.mutations_cipher.lock();
127 if let Some(cipher) = mutations_cipher.as_ref() {
128 store_info.mutations_cipher_offset = cipher.offset();
129 }
130
131 self.store_info.set(store_info).unwrap();
132 }
133 }
134
135 let store_info_snapshot = StoreInfoSnapshot { store: self, store_info: OnceLock::new() };
136
137 let filesystem = self.filesystem();
138 let object_manager = filesystem.object_manager();
139 let reservation = object_manager.metadata_reservation();
140 let txn_options = Options {
141 skip_journal_checks: true,
142 skip_key_roll: true,
143 borrow_metadata_space: true,
144 allocator_reservation: Some(reservation),
145 ..Default::default()
146 };
147
148 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
151 transaction.add_with_object(
152 self.store_object_id(),
153 Mutation::BeginFlush,
154 AssocObj::Borrowed(&store_info_snapshot),
155 );
156 transaction.commit().await?;
157
158 let mut new_store_info = store_info_snapshot.store_info.into_inner().unwrap();
159
160 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
165
166 let parent_store = self.parent_store.as_ref().unwrap();
168 let handle_options = HandleOptions { skip_journal_checks: true, ..Default::default() };
169 let id_and_key = {
170 let mut lock_state = self.lock_state.lock();
171 match &mut *lock_state {
172 LockState::Unlocked { cached_keys, .. } => {
173 if let Some(item) = cached_keys.pop() {
174 Ok(Some(item))
175 } else {
176 log::warn!("No cached keys for flush for store {}", self.store_object_id());
177 Err(anyhow!(FxfsError::Internal).context("No cached keys for flush"))
178 }
179 }
180 LockState::UnlockedReadOnly(..) => {
181 Err(anyhow!(FxfsError::Internal).context("Flush on read-only store"))
182 }
183 LockState::Unencrypted => Ok(None),
184 _ => Err(anyhow!(FxfsError::Internal))
185 .with_context(|| format!("Invalid lock state ({:?}) for flush", *lock_state)),
186 }
187 }?;
188
189 let new_object_tree_layer = if let Some((raw_id, key, unwrapped_key)) = id_and_key {
190 let object_id = ReservedId::new(parent_store, raw_id);
191 ObjectStore::create_object_with_key(
192 parent_store,
193 &mut transaction,
194 object_id,
195 handle_options,
196 key,
197 unwrapped_key,
198 )
199 .await?
200 } else {
201 ObjectStore::create_object(parent_store, &mut transaction, handle_options, None).await?
202 };
203 let writer = DirectWriter::new(&new_object_tree_layer, txn_options).await;
204 let new_object_tree_layer_object_id = new_object_tree_layer.object_id();
205 parent_store.add_to_graveyard(&mut transaction, new_object_tree_layer_object_id);
206
207 transaction.commit().await?;
208
209 let (layers_to_keep, old_layers) = tree::flush(
211 &self.tree,
212 writer,
213 (reason == Reason::Journal).then(|| filesystem.journal().get_compaction_yielder()),
214 reason == Reason::UpgradeVersion,
215 )
216 .await
217 .context("Failed to flush tree")?;
218
219 let mut new_layers = layers_from_handles([new_object_tree_layer]).await?;
221 new_layers.extend(layers_to_keep.iter().map(|l| (*l).clone()));
222
223 new_store_info.layers = Vec::new();
224 for layer in &new_layers {
225 if let Some(handle) = layer.handle() {
226 new_store_info.layers.push(handle.object_id());
227 }
228 }
229
230 let reservation_update: ReservationUpdate; let mut end_transaction = parent_store
232 .new_transaction(
233 lock_keys![LockKey::object(
234 self.parent_store.as_ref().unwrap().store_object_id(),
235 self.store_info_handle_object_id().unwrap(),
236 )],
237 txn_options,
238 )
239 .await?;
240
241 parent_store.remove_from_graveyard(&mut end_transaction, new_object_tree_layer_object_id);
242
243 for layer in &old_layers {
245 if let Some(handle) = layer.handle() {
246 parent_store.add_to_graveyard(&mut end_transaction, handle.object_id());
247 }
248 }
249
250 let old_encrypted_mutations_object_id =
251 std::mem::replace(&mut new_store_info.encrypted_mutations_object_id, INVALID_OBJECT_ID);
252 if old_encrypted_mutations_object_id != INVALID_OBJECT_ID {
253 parent_store.add_to_graveyard(&mut end_transaction, old_encrypted_mutations_object_id);
254 }
255
256 match &mut new_store_info.last_object_id {
264 LastObjectIdInfo::Unencrypted { id } => {
265 let LastObjectId::Unencrypted { id: in_memory_value } =
266 &*self.last_object_id.lock()
267 else {
268 unreachable!()
269 };
270 *id = *in_memory_value;
271 }
272 LastObjectIdInfo::Encrypted { id, key } => {
273 let LastObjectId::Encrypted { id: in_memory_value, .. } =
274 &*self.last_object_id.lock()
275 else {
276 unreachable!()
277 };
278 *id = *in_memory_value;
279 let guard = self.store_info.lock();
280 let current_store_info = guard.as_ref().unwrap();
281 let LastObjectIdInfo::Encrypted { key: in_memory_value, .. } =
282 ¤t_store_info.last_object_id
283 else {
284 unreachable!()
285 };
286 *key = in_memory_value.clone();
287 }
288 LastObjectIdInfo::Low32Bit => {}
289 }
290
291 self.write_store_info(&mut end_transaction, &new_store_info).await?;
292
293 let layer_file_sizes = new_layers
294 .iter()
295 .map(|l| l.handle().map(ReadObjectHandle::get_size).unwrap_or(0))
296 .collect::<Vec<u64>>();
297
298 let total_layer_size = layer_file_sizes.iter().sum();
299 reservation_update =
300 ReservationUpdate::new(tree::reservation_amount_from_layer_size(total_layer_size));
301
302 end_transaction.add_with_object(
303 self.store_object_id(),
304 Mutation::EndFlush,
305 AssocObj::Borrowed(&reservation_update),
306 );
307
308 if self.trace.load(Ordering::Relaxed) {
309 info!(
310 store_id = self.store_object_id(),
311 old_layer_count = old_layers.len(),
312 new_layer_count = new_layers.len(),
313 total_layer_size,
314 new_store_info:?;
315 "OS: compacting"
316 );
317 }
318
319 end_transaction
320 .commit_with_callback(|_| {
321 let mut store_info = self.store_info.lock();
322 let info = store_info.as_mut().unwrap();
323 info.layers = new_store_info.layers;
324 info.encrypted_mutations_object_id = new_store_info.encrypted_mutations_object_id;
325 info.mutations_cipher_offset = new_store_info.mutations_cipher_offset;
326 self.tree.set_layers(new_layers);
327 })
328 .await?;
329
330 for layer in old_layers {
332 let object_id = layer.handle().map(|h| h.object_id());
333 layer.close_layer().await;
334 if let Some(object_id) = object_id {
335 parent_store.tombstone_object(object_id, txn_options).await?;
336 }
337 }
338
339 if old_encrypted_mutations_object_id != INVALID_OBJECT_ID {
340 parent_store.tombstone_object(old_encrypted_mutations_object_id, txn_options).await?;
341 }
342
343 Ok(layer_file_sizes)
344 }
345
346 async fn flush_locked(&self) -> Result<(), Error> {
348 let filesystem = self.filesystem();
349 let object_manager = filesystem.object_manager();
350 let reservation = object_manager.metadata_reservation();
351 let txn_options = Options {
352 skip_journal_checks: true,
353 skip_key_roll: true,
354 borrow_metadata_space: true,
355 allocator_reservation: Some(reservation),
356 ..Default::default()
357 };
358
359 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
360 transaction.add(self.store_object_id(), Mutation::BeginFlush);
361 transaction.commit().await?;
362
363 let mut new_store_info = self.load_store_info().await?;
364
365 let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
370
371 let reservation_update: ReservationUpdate; let handle; let mut end_transaction;
374
375 let parent_store = self.parent_store.as_ref().unwrap();
378 handle = if new_store_info.encrypted_mutations_object_id == INVALID_OBJECT_ID {
379 let handle = ObjectStore::create_object(
380 parent_store,
381 &mut transaction,
382 HandleOptions { skip_journal_checks: true, ..Default::default() },
383 None,
384 )
385 .await?;
386 let oid = handle.object_id();
387 end_transaction = parent_store
388 .new_transaction(
389 lock_keys![
390 LockKey::object(parent_store.store_object_id(), oid),
391 LockKey::object(
392 parent_store.store_object_id(),
393 self.store_info_handle_object_id().unwrap(),
394 ),
395 ],
396 txn_options,
397 )
398 .await?;
399 new_store_info.encrypted_mutations_object_id = oid;
400 parent_store.add_to_graveyard(&mut transaction, oid);
401 parent_store.remove_from_graveyard(&mut end_transaction, oid);
402 handle
403 } else {
404 end_transaction = parent_store
405 .new_transaction(
406 lock_keys![
407 LockKey::object(
408 parent_store.store_object_id(),
409 new_store_info.encrypted_mutations_object_id,
410 ),
411 LockKey::object(
412 parent_store.store_object_id(),
413 self.store_info_handle_object_id().unwrap(),
414 ),
415 ],
416 txn_options,
417 )
418 .await?;
419 ObjectStore::open_object(
420 parent_store,
421 new_store_info.encrypted_mutations_object_id,
422 HandleOptions { skip_journal_checks: true, ..Default::default() },
423 None,
424 )
425 .await?
426 };
427 transaction.commit().await?;
428
429 let journaled = filesystem
432 .journal()
433 .read_transactions_for_object(self.store_object_id)
434 .await
435 .context("Failed to read encrypted mutations from journal")?;
436 let mut buffer = handle.allocate_buffer(MAX_ENCRYPTED_MUTATIONS_SIZE).await;
437 let mut writer = buffer.writer();
438 EncryptedMutations::from_replayed_mutations(self.store_object_id, journaled)
439 .serialize_with_version(&mut writer)?;
440 let len = writer.position();
441 handle.txn_write(&mut end_transaction, handle.get_size(), buffer.subslice(..len)).await?;
442
443 self.write_store_info(&mut end_transaction, &new_store_info).await?;
444
445 let mut total_layer_size = 0;
446 for &oid in &new_store_info.layers {
447 total_layer_size += parent_store.get_file_size(oid).await?;
448 }
449 total_layer_size +=
450 layer_size_from_encrypted_mutations_size(handle.get_size() + len as u64);
451
452 reservation_update =
453 ReservationUpdate::new(tree::reservation_amount_from_layer_size(total_layer_size));
454
455 end_transaction.add_with_object(
456 self.store_object_id(),
457 Mutation::EndFlush,
458 AssocObj::Borrowed(&reservation_update),
459 );
460
461 end_transaction.commit().await?;
462
463 Ok(())
464 }
465}
466
467#[cfg(test)]
468mod tests {
469 use crate::filesystem::{FxFilesystem, FxFilesystemBuilder, JournalingObject, SyncOptions};
470 use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle};
471 use crate::object_store::directory::Directory;
472 use crate::object_store::transaction::{Options, lock_keys};
473 use crate::object_store::volume::root_volume;
474 use crate::object_store::{
475 HandleOptions, LockKey, NewChildStoreOptions, ObjectStore, StoreOptions,
476 layer_size_from_encrypted_mutations_size, tree,
477 };
478 use fxfs_insecure_crypto::new_insecure_crypt;
479 use std::sync::Arc;
480 use storage_device::DeviceHolder;
481 use storage_device::fake_device::FakeDevice;
482
483 async fn run_key_roll_test(flush_before_unlock: bool) {
484 let device = DeviceHolder::new(FakeDevice::new(8192, 1024));
485 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
486 let store_id = {
487 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
488 root_volume
489 .new_volume(
490 "test",
491 NewChildStoreOptions {
492 options: StoreOptions {
493 crypt: Some(Arc::new(new_insecure_crypt())),
494 ..StoreOptions::default()
495 },
496 ..NewChildStoreOptions::default()
497 },
498 )
499 .await
500 .expect("new_volume failed")
501 .store_object_id()
502 };
503
504 fs.close().await.expect("close failed");
505 let device = fs.take_device().await;
506 device.reopen(false);
507
508 let fs = FxFilesystemBuilder::new()
509 .roll_metadata_key_byte_count(512 * 1024)
510 .open(device)
511 .await
512 .expect("open failed");
513
514 let (first_filename, last_filename) = {
515 let store = fs.object_manager().store(store_id).expect("store not found");
516 store.unlock(Arc::new(new_insecure_crypt())).await.expect("unlock failed");
517
518 let root_dir = Directory::open(&store, store.root_directory_object_id())
520 .await
521 .expect("open failed");
522
523 let mut last_mutations_cipher_offset = 0;
524 let mut i = 0;
525 let first_filename = format!("{:<200}", i);
526 loop {
527 let mut transaction = store
528 .new_transaction(
529 lock_keys![LockKey::object(store_id, root_dir.object_id())],
530 Options::default(),
531 )
532 .await
533 .expect("new_transaction failed");
534 root_dir
535 .create_child_file(&mut transaction, &format!("{:<200}", i))
536 .await
537 .expect("create_child_file failed");
538 i += 1;
539 transaction.commit().await.expect("commit failed");
540 let cipher_offset = store.mutations_cipher.lock().as_ref().unwrap().offset();
541 if cipher_offset < last_mutations_cipher_offset {
542 break;
543 }
544 last_mutations_cipher_offset = cipher_offset;
545 }
546
547 fs.sync(SyncOptions::default()).await.expect("sync failed");
550
551 let mut transaction = store
553 .new_transaction(
554 lock_keys![LockKey::object(store_id, root_dir.object_id())],
555 Options::default(),
556 )
557 .await
558 .expect("new_transaction failed");
559 let last_filename = format!("{:<200}", i);
560 root_dir
561 .create_child_file(&mut transaction, &last_filename)
562 .await
563 .expect("create_child_file failed");
564 transaction.commit().await.expect("commit failed");
565 (first_filename, last_filename)
566 };
567
568 fs.close().await.expect("close failed");
569
570 let device = fs.take_device().await;
572 device.reopen(false);
573 let fs = FxFilesystemBuilder::new()
574 .roll_metadata_key_byte_count(512 * 1024)
575 .open(device)
576 .await
577 .expect("open failed");
578
579 if flush_before_unlock {
580 fs.object_manager().flush().await.expect("flush failed");
583 }
584
585 {
586 let store = fs.object_manager().store(store_id).expect("store not found");
587 store.unlock(Arc::new(new_insecure_crypt())).await.expect("unlock failed");
588
589 assert_eq!(store.mutations_cipher.lock().as_ref().unwrap().offset(), 0);
591
592 let root_dir = Directory::open(&store, store.root_directory_object_id())
593 .await
594 .expect("open failed");
595 root_dir
596 .lookup(&first_filename)
597 .await
598 .expect("Lookup failed")
599 .expect("First created file wasn't present");
600 root_dir
601 .lookup(&last_filename)
602 .await
603 .expect("Lookup failed")
604 .expect("Last created file wasn't present");
605 }
606
607 fs.close().await.expect("close failed");
608 }
609
610 #[fuchsia::test(threads = 10)]
611 async fn test_metadata_key_roll() {
612 run_key_roll_test(false).await;
613 }
614
615 #[fuchsia::test(threads = 10)]
616 async fn test_metadata_key_roll_with_flush_before_unlock() {
617 run_key_roll_test(true).await;
618 }
619
620 #[fuchsia::test]
621 async fn test_flush_when_locked() {
622 let device = DeviceHolder::new(FakeDevice::new(8192, 1024));
623 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
624 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
625 let crypt = Arc::new(new_insecure_crypt());
626 let store = root_volume
627 .new_volume(
628 "test",
629 NewChildStoreOptions {
630 options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
631 ..NewChildStoreOptions::default()
632 },
633 )
634 .await
635 .expect("new_volume failed");
636 let root_dir =
637 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
638 let mut transaction = fs
639 .root_store()
640 .new_transaction(
641 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
642 Options::default(),
643 )
644 .await
645 .expect("new_transaction failed");
646 let foo = root_dir
647 .create_child_file(&mut transaction, "foo")
648 .await
649 .expect("create_child_file failed");
650 transaction.commit().await.expect("commit failed");
651
652 store.flush().await.expect("flush failed");
656
657 let mut transaction = fs
658 .root_store()
659 .new_transaction(
660 lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
661 Options::default(),
662 )
663 .await
664 .expect("new_transaction failed");
665 let bar = root_dir
666 .create_child_file(&mut transaction, "bar")
667 .await
668 .expect("create_child_file failed");
669 transaction.commit().await.expect("commit failed");
670
671 store.lock().await.expect("lock failed");
672
673 store.flush().await.expect("flush failed");
675
676 let info = store.load_store_info().await.unwrap();
678 let parent_store = store.parent_store().unwrap();
679 let mut total_layer_size = 0;
680 for &oid in &info.layers {
681 total_layer_size +=
682 parent_store.get_file_size(oid).await.expect("get_file_size failed");
683 }
684 assert_ne!(info.encrypted_mutations_object_id, INVALID_OBJECT_ID);
685 total_layer_size += layer_size_from_encrypted_mutations_size(
686 parent_store
687 .get_file_size(info.encrypted_mutations_object_id)
688 .await
689 .expect("get_file_size failed"),
690 );
691 assert_eq!(
692 fs.object_manager().reservation(store.store_object_id()),
693 Some(tree::reservation_amount_from_layer_size(total_layer_size))
694 );
695
696 store.unlock(crypt).await.expect("unlock failed");
698
699 ObjectStore::open_object(&store, foo.object_id(), HandleOptions::default(), None)
700 .await
701 .expect("open_object failed");
702
703 ObjectStore::open_object(&store, bar.object_id(), HandleOptions::default(), None)
704 .await
705 .expect("open_object failed");
706
707 fs.close().await.expect("close failed");
708 }
709}
710
711impl tree::MajorCompactable<ObjectKey, ObjectValue> for LSMTree<ObjectKey, ObjectValue> {
712 async fn major_iter(
713 iter: impl LayerIterator<ObjectKey, ObjectValue>,
714 ) -> Result<impl LayerIterator<ObjectKey, ObjectValue>, Error> {
715 iter.filter(|item: ItemRef<'_, _, _>| match item {
716 ItemRef { value: ObjectValue::None, .. } => false,
718 ItemRef { value: ObjectValue::Extent(ExtentValue::None), .. } => false,
720 _ => true,
721 })
722 .await
723 }
724}