1use crate::checksum::Checksum;
6use crate::filesystem::FxFilesystem;
7use crate::log::*;
8use crate::lsm_tree::types::Item;
9use crate::object_handle::INVALID_OBJECT_ID;
10use crate::object_store::allocator::{AllocatorItem, Hold, Reservation};
11use crate::object_store::object_manager::{ObjectManager, reserved_space_from_journal_usage};
12use crate::object_store::object_record::{
13 BytesAndNodes, FxfsKey, FxfsKeyV49, ObjectItem, ObjectItemV56, ObjectItemV59, ObjectKey,
14 ObjectKeyData, ObjectValue, ProjectProperty,
15};
16use crate::object_store::{AttributeId, AttributeKey, ProjectId};
17use crate::serialized_types::{Migrate, Versioned, migrate_nodefault, migrate_to_version};
18use anyhow::Error;
19use either::{Either, Left, Right};
20use fprint::TypeFingerprint;
21use fuchsia_sync::Mutex;
22use futures::future::poll_fn;
23use futures::pin_mut;
24use rustc_hash::FxHashMap as HashMap;
25use scopeguard::ScopeGuard;
26use serde::{Deserialize, Serialize};
27use std::cell::UnsafeCell;
28use std::cmp::Ordering;
29use std::collections::hash_map::Entry;
30use std::collections::{BTreeSet, btree_set};
31use std::iter::Peekable;
32use std::marker::PhantomPinned;
33use std::ops::{Deref, DerefMut, Range};
34use std::sync::Arc;
35use std::task::{Poll, Waker};
36use std::{fmt, mem};
37
38#[derive(Clone, Copy, Default)]
60pub enum ReservationOptions<'a> {
61 #[default]
69 New,
70
71 BorrowedMetadata,
90
91 BorrowedMetadataAndData,
101
102 Hold(&'a Reservation),
105}
106
107#[derive(Clone, Copy, Default)]
111pub struct Options<'a> {
112 pub skip_journal_checks: bool,
115
116 pub reservation: ReservationOptions<'a>,
118
119 pub parent_transaction: Option<&'a Transaction<'a>>,
127}
128
129pub const TRANSACTION_MAX_JOURNAL_USAGE: u64 = 24_576;
136pub const TRANSACTION_METADATA_MAX_AMOUNT: u64 =
137 reserved_space_from_journal_usage(TRANSACTION_MAX_JOURNAL_USAGE);
138
139#[must_use]
140pub struct TransactionLocks<'a>(pub WriteGuard<'a>);
141
142pub type Mutation = MutationV59;
147
148#[derive(
149 Clone, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize, TypeFingerprint, Versioned,
150)]
151#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
152pub enum MutationV59 {
153 ObjectStore(ObjectStoreMutationV59),
154 EncryptedObjectStore(#[serde(with = "crate::zerocopy_serialization")] Box<[u8]>),
155 Allocator(AllocatorMutationV32),
156 BeginFlush,
158 EndFlush,
161 DeleteVolume,
163 UpdateBorrowed(u64),
164 UpdateMutationsKey(UpdateMutationsKey),
165 CreateInternalDir(u64),
166}
167
168#[derive(Migrate, Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
169#[migrate_to_version(MutationV59)]
170pub enum MutationV57 {
171 ObjectStore(ObjectStoreMutationV56),
172 EncryptedObjectStore(#[serde(with = "crate::zerocopy_serialization")] Box<[u8]>),
173 Allocator(AllocatorMutationV32),
174 BeginFlush,
175 EndFlush,
176 DeleteVolume,
177 UpdateBorrowed(u64),
178 UpdateMutationsKey(UpdateMutationsKey),
179 CreateInternalDir(u64),
180}
181
182#[derive(Migrate, Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
183#[migrate_to_version(MutationV57)]
184pub enum MutationV56 {
185 ObjectStore(ObjectStoreMutationV56),
186 EncryptedObjectStore(#[serde(with = "crate::zerocopy_serialization")] Box<[u8]>),
187 Allocator(AllocatorMutationV32),
188 BeginFlush,
189 EndFlush,
190 DeleteVolume,
191 UpdateBorrowed(u64),
192 UpdateMutationsKey(UpdateMutationsKey),
193 CreateInternalDir(u64),
194}
195
196impl Mutation {
197 pub fn insert_object(key: ObjectKey, value: ObjectValue) -> Self {
198 Mutation::ObjectStore(ObjectStoreMutation {
199 item: Item::new(key, value),
200 op: Operation::Insert,
201 })
202 }
203
204 pub fn replace_or_insert_object(key: ObjectKey, value: ObjectValue) -> Self {
205 Mutation::ObjectStore(ObjectStoreMutation {
206 item: Item::new(key, value),
207 op: Operation::ReplaceOrInsert,
208 })
209 }
210
211 pub fn merge_object(key: ObjectKey, value: ObjectValue) -> Self {
212 Mutation::ObjectStore(ObjectStoreMutation {
213 item: Item::new(key, value),
214 op: Operation::Merge,
215 })
216 }
217
218 pub fn update_mutations_key(key: FxfsKey) -> Self {
219 Mutation::UpdateMutationsKey(key.into())
220 }
221}
222
223pub type ObjectStoreMutation = ObjectStoreMutationV59;
227
228#[derive(Clone, Debug, Serialize, Deserialize, TypeFingerprint)]
229#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
230pub struct ObjectStoreMutationV59 {
231 pub item: ObjectItemV59,
232 pub op: Operation,
233}
234
235#[derive(Migrate, Clone, Debug, PartialEq, Serialize, Deserialize, TypeFingerprint, Versioned)]
236#[migrate_to_version(ObjectStoreMutationV59)]
237#[migrate_nodefault]
238pub struct ObjectStoreMutationV56 {
239 pub item: ObjectItemV56,
240 pub op: Operation,
241}
242
243pub type Operation = OperationV32;
245
246#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, TypeFingerprint)]
247#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
248pub enum OperationV32 {
249 Insert,
250 ReplaceOrInsert,
251 Merge,
252}
253
254impl Ord for ObjectStoreMutation {
255 fn cmp(&self, other: &Self) -> Ordering {
256 self.item.key.cmp(&other.item.key)
257 }
258}
259
260impl PartialOrd for ObjectStoreMutation {
261 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
262 Some(self.cmp(other))
263 }
264}
265
266impl PartialEq for ObjectStoreMutation {
267 fn eq(&self, other: &Self) -> bool {
268 self.item.key.eq(&other.item.key)
269 }
270}
271
272impl Eq for ObjectStoreMutation {}
273
274impl Ord for AllocatorItem {
275 fn cmp(&self, other: &Self) -> Ordering {
276 self.key.cmp(&other.key)
277 }
278}
279
280impl PartialOrd for AllocatorItem {
281 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
282 Some(self.cmp(other))
283 }
284}
285
286pub type DeviceRange = DeviceRangeV32;
289
290#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TypeFingerprint)]
291#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
292pub struct DeviceRangeV32(pub Range<u64>);
293
294impl Deref for DeviceRange {
295 type Target = Range<u64>;
296
297 fn deref(&self) -> &Self::Target {
298 &self.0
299 }
300}
301
302impl DerefMut for DeviceRange {
303 fn deref_mut(&mut self) -> &mut Self::Target {
304 &mut self.0
305 }
306}
307
308impl From<Range<u64>> for DeviceRange {
309 fn from(range: Range<u64>) -> Self {
310 Self(range)
311 }
312}
313
314impl Into<Range<u64>> for DeviceRange {
315 fn into(self) -> Range<u64> {
316 self.0
317 }
318}
319
320impl Ord for DeviceRange {
321 fn cmp(&self, other: &Self) -> Ordering {
322 self.start.cmp(&other.start).then(self.end.cmp(&other.end))
323 }
324}
325
326impl PartialOrd for DeviceRange {
327 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
328 Some(self.cmp(other))
329 }
330}
331
332pub type AllocatorMutation = AllocatorMutationV32;
333
334#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize, TypeFingerprint)]
335#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
336pub enum AllocatorMutationV32 {
337 Allocate {
338 device_range: DeviceRangeV32,
339 owner_object_id: u64,
340 },
341 Deallocate {
342 device_range: DeviceRangeV32,
343 owner_object_id: u64,
344 },
345 SetLimit {
346 owner_object_id: u64,
347 bytes: u64,
348 },
349 MarkForDeletion(u64),
355}
356
357pub type UpdateMutationsKey = UpdateMutationsKeyV49;
358
359#[derive(Clone, Debug, Serialize, Deserialize, TypeFingerprint)]
360pub struct UpdateMutationsKeyV49(pub FxfsKeyV49);
361
362impl From<UpdateMutationsKey> for FxfsKey {
363 fn from(outer: UpdateMutationsKey) -> Self {
364 outer.0
365 }
366}
367
368impl From<FxfsKey> for UpdateMutationsKey {
369 fn from(inner: FxfsKey) -> Self {
370 Self(inner)
371 }
372}
373
374#[cfg(fuzz)]
375impl<'a> arbitrary::Arbitrary<'a> for UpdateMutationsKey {
376 fn arbitrary(u: &mut arbitrary::Unstructured<'a>) -> arbitrary::Result<Self> {
377 Ok(UpdateMutationsKey::from(FxfsKey::arbitrary(u).unwrap()))
378 }
379}
380
381impl Ord for UpdateMutationsKey {
382 fn cmp(&self, other: &Self) -> Ordering {
383 (self as *const UpdateMutationsKey).cmp(&(other as *const _))
384 }
385}
386
387impl PartialOrd for UpdateMutationsKey {
388 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
389 Some(self.cmp(other))
390 }
391}
392
393impl Eq for UpdateMutationsKey {}
394
395impl PartialEq for UpdateMutationsKey {
396 fn eq(&self, other: &Self) -> bool {
397 std::ptr::eq(self, other)
398 }
399}
400
401#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Copy)]
409pub enum LockKey {
410 Flush {
412 object_id: u64,
413 },
414
415 ObjectAttribute {
417 store_object_id: u64,
418 object_id: u64,
419 attribute_id: AttributeId,
420 },
421
422 Object {
424 store_object_id: u64,
425 object_id: u64,
426 },
427
428 ProjectId {
429 store_object_id: u64,
430 project_id: ProjectId,
431 },
432
433 Truncate {
435 store_object_id: u64,
436 object_id: u64,
437 },
438
439 InternalDirectory {
441 store_object_id: u64,
442 },
443
444 PreCacheKeys {
448 store_object_id: u64,
449 },
450}
451
452impl LockKey {
453 pub const fn object_attribute(
454 store_object_id: u64,
455 object_id: u64,
456 attribute_id: AttributeId,
457 ) -> Self {
458 LockKey::ObjectAttribute { store_object_id, object_id, attribute_id }
459 }
460
461 pub const fn object(store_object_id: u64, object_id: u64) -> Self {
462 LockKey::Object { store_object_id, object_id }
463 }
464
465 pub const fn flush(object_id: u64) -> Self {
466 LockKey::Flush { object_id }
467 }
468
469 pub const fn truncate(store_object_id: u64, object_id: u64) -> Self {
470 LockKey::Truncate { store_object_id, object_id }
471 }
472
473 pub const fn pre_cache_keys(store_object_id: u64) -> Self {
474 LockKey::PreCacheKeys { store_object_id }
475 }
476}
477
478#[derive(Clone, Debug)]
480pub enum LockKeys {
481 None,
482 Inline(LockKey),
483 Vec(Vec<LockKey>),
484}
485
486impl LockKeys {
487 pub fn with_capacity(capacity: usize) -> Self {
488 if capacity > 1 { LockKeys::Vec(Vec::with_capacity(capacity)) } else { LockKeys::None }
489 }
490
491 pub fn push(&mut self, key: LockKey) {
492 match self {
493 Self::None => *self = LockKeys::Inline(key),
494 Self::Inline(inline) => {
495 *self = LockKeys::Vec(vec![*inline, key]);
496 }
497 Self::Vec(vec) => vec.push(key),
498 }
499 }
500
501 pub fn truncate(&mut self, len: usize) {
502 match self {
503 Self::None => {}
504 Self::Inline(_) => {
505 if len == 0 {
506 *self = Self::None;
507 }
508 }
509 Self::Vec(vec) => vec.truncate(len),
510 }
511 }
512
513 fn len(&self) -> usize {
514 match self {
515 Self::None => 0,
516 Self::Inline(_) => 1,
517 Self::Vec(vec) => vec.len(),
518 }
519 }
520
521 fn contains(&self, key: &LockKey) -> bool {
522 match self {
523 Self::None => false,
524 Self::Inline(single) => single == key,
525 Self::Vec(vec) => vec.contains(key),
526 }
527 }
528
529 fn sort_unstable(&mut self) {
530 match self {
531 Self::Vec(vec) => vec.sort_unstable(),
532 _ => {}
533 }
534 }
535
536 fn dedup(&mut self) {
537 match self {
538 Self::Vec(vec) => vec.dedup(),
539 _ => {}
540 }
541 }
542
543 fn iter(&self) -> LockKeysIter<'_> {
544 match self {
545 LockKeys::None => LockKeysIter::None,
546 LockKeys::Inline(key) => LockKeysIter::Inline(key),
547 LockKeys::Vec(keys) => LockKeysIter::Vec(keys.iter()),
548 }
549 }
550}
551
552enum LockKeysIter<'a> {
553 None,
554 Inline(&'a LockKey),
555 Vec(std::slice::Iter<'a, LockKey>),
556}
557
558impl<'a> Iterator for LockKeysIter<'a> {
559 type Item = &'a LockKey;
560 fn next(&mut self) -> Option<Self::Item> {
561 match self {
562 Self::None => None,
563 Self::Inline(inline) => {
564 let next = *inline;
565 *self = Self::None;
566 Some(next)
567 }
568 Self::Vec(vec) => vec.next(),
569 }
570 }
571}
572
573impl Default for LockKeys {
574 fn default() -> Self {
575 LockKeys::None
576 }
577}
578
579#[macro_export]
580macro_rules! lock_keys {
581 () => {
582 $crate::object_store::transaction::LockKeys::None
583 };
584 ($lock_key:expr $(,)?) => {
585 $crate::object_store::transaction::LockKeys::Inline($lock_key)
586 };
587 ($($lock_keys:expr),+ $(,)?) => {
588 $crate::object_store::transaction::LockKeys::Vec(vec![$($lock_keys),+])
589 };
590}
591pub use lock_keys;
592
593pub trait AssociatedObject: Send + Sync {
597 fn will_apply_mutation(&self, _mutation: &Mutation, _object_id: u64, _manager: &ObjectManager) {
598 }
599}
600
601pub enum AssocObj<'a> {
602 None,
603 Borrowed(&'a dyn AssociatedObject),
604 Owned(Box<dyn AssociatedObject>),
605}
606
607impl AssocObj<'_> {
608 pub fn map<R, F: FnOnce(&dyn AssociatedObject) -> R>(&self, f: F) -> Option<R> {
609 match self {
610 AssocObj::None => None,
611 AssocObj::Borrowed(b) => Some(f(*b)),
612 AssocObj::Owned(o) => Some(f(o.as_ref())),
613 }
614 }
615}
616
617pub struct TxnMutation<'a> {
618 pub object_id: u64,
622
623 pub mutation: Mutation,
625
626 pub associated_object: AssocObj<'a>,
629}
630
631impl Ord for TxnMutation<'_> {
639 fn cmp(&self, other: &Self) -> Ordering {
640 self.object_id.cmp(&other.object_id).then_with(|| self.mutation.cmp(&other.mutation))
641 }
642}
643
644impl PartialOrd for TxnMutation<'_> {
645 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
646 Some(self.cmp(other))
647 }
648}
649
650impl PartialEq for TxnMutation<'_> {
651 fn eq(&self, other: &Self) -> bool {
652 self.object_id.eq(&other.object_id) && self.mutation.eq(&other.mutation)
653 }
654}
655
656impl Eq for TxnMutation<'_> {}
657
658impl std::fmt::Debug for TxnMutation<'_> {
659 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
660 f.debug_struct("TxnMutation")
661 .field("object_id", &self.object_id)
662 .field("mutation", &self.mutation)
663 .finish()
664 }
665}
666
667pub struct ObjectMutationIterator<'a, 'b> {
671 iter: &'a mut Peekable<btree_set::Iter<'b, TxnMutation<'b>>>,
672 object_id: u64,
673}
674
675impl<'a, 'b> ObjectMutationIterator<'a, 'b> {
676 pub fn new(iter: &'a mut Peekable<btree_set::Iter<'b, TxnMutation<'b>>>) -> Option<Self> {
677 let object_id = iter.peek()?.object_id;
678 Some(Self { iter, object_id })
679 }
680
681 pub fn object_id(&self) -> u64 {
682 self.object_id
683 }
684}
685
686impl<'b> Iterator for ObjectMutationIterator<'_, 'b> {
687 type Item = &'b Mutation;
688
689 fn next(&mut self) -> Option<Self::Item> {
690 if self.iter.peek().is_some_and(|m| m.object_id == self.object_id) {
691 Some(&self.iter.next().unwrap().mutation)
692 } else {
693 None
694 }
695 }
696}
697
698impl Drop for ObjectMutationIterator<'_, '_> {
699 fn drop(&mut self) {
702 for _ in self.by_ref() {}
703 }
704}
705
706pub enum MetadataReservation<'a> {
707 BorrowedMetadata,
710
711 BorrowedMetadataAndData,
715
716 Reservation(Reservation),
718
719 Hold(Hold<'a>),
721}
722
723pub struct Transaction<'a> {
725 fs: Arc<FxFilesystem>,
726
727 mutations: BTreeSet<TxnMutation<'a>>,
729
730 txn_locks: LockKeys,
732
733 pub metadata_reservation: MetadataReservation<'a>,
735
736 new_objects: BTreeSet<(u64, u64)>,
739
740 checksums: Vec<(Range<u64>, Vec<Checksum>, bool)>,
742
743 includes_write: bool,
745
746 pub skip_journal_checks: bool,
748
749 parent_transaction: Option<&'a Transaction<'a>>,
752}
753
754impl<'a> Transaction<'a> {
755 pub async fn new(
758 fs: Arc<FxFilesystem>,
759 options: Options<'a>,
760 txn_locks: LockKeys,
761 ) -> Result<Transaction<'a>, Error> {
762 if options.parent_transaction.is_none() {
763 fs.add_transaction(options.skip_journal_checks).await;
764 }
765 let metadata_reservation = if fs.options().image_builder_mode.is_some() {
779 MetadataReservation::BorrowedMetadata
780 } else {
781 match options.reservation {
782 ReservationOptions::New => {
783 MetadataReservation::Reservation(fs.allocator().reserve(None, 0).unwrap())
784 }
785 ReservationOptions::BorrowedMetadata => MetadataReservation::BorrowedMetadata,
786 ReservationOptions::BorrowedMetadataAndData => {
787 MetadataReservation::BorrowedMetadataAndData
788 }
789 ReservationOptions::Hold(reservation) => {
790 MetadataReservation::Hold(reservation.reserve(0).unwrap())
791 }
792 }
793 };
794 let mut transaction = Transaction {
795 fs: fs.clone(),
796 mutations: BTreeSet::new(),
797 txn_locks: LockKeys::default(),
798 metadata_reservation,
799 new_objects: BTreeSet::new(),
800 checksums: Vec::new(),
801 includes_write: false,
802 skip_journal_checks: options.skip_journal_checks,
803 parent_transaction: options.parent_transaction,
804 };
805 fs.add_transaction_reservation(&mut transaction).await?;
806
807 transaction.txn_locks = {
808 let lock_manager = fs.lock_manager();
809 let mut write_guard = lock_manager.txn_lock(txn_locks).await;
810 std::mem::take(&mut write_guard.0.lock_keys)
811 };
812 Ok(transaction)
813 }
814
815 pub fn allocator_reservation(&self) -> Option<&Reservation> {
818 match &self.metadata_reservation {
819 MetadataReservation::BorrowedMetadataAndData => {
820 Some(self.fs.object_manager().metadata_reservation())
821 }
822 MetadataReservation::Hold(hold) => Some(hold.owner()),
823 MetadataReservation::BorrowedMetadata | MetadataReservation::Reservation(_) => None,
824 }
825 }
826
827 pub fn mutations(&self) -> &BTreeSet<TxnMutation<'a>> {
828 &self.mutations
829 }
830
831 pub fn take_mutations(&mut self) -> BTreeSet<TxnMutation<'a>> {
832 self.new_objects.clear();
833 mem::take(&mut self.mutations)
834 }
835
836 pub fn add(&mut self, object_id: u64, mutation: Mutation) -> Option<Mutation> {
839 self.add_with_object(object_id, mutation, AssocObj::None)
840 }
841
842 pub fn merge_bytes_and_nodes(
846 &mut self,
847 store_object_id: u64,
848 key: ObjectKey,
849 delta: BytesAndNodes,
850 ) {
851 let delta = match self.get_object_mutation(store_object_id, key.clone()) {
852 Some(ObjectStoreMutation {
853 item: Item { value: ObjectValue::BytesAndNodes { bytes, nodes }, .. },
854 ..
855 }) => delta + BytesAndNodes { bytes: *bytes, nodes: *nodes },
856 _ => delta,
857 };
858 if delta.is_zero() {
859 self.remove(store_object_id, Mutation::merge_object(key, ObjectValue::None));
860 } else {
861 self.add_with_object_internal(
862 store_object_id,
863 Mutation::merge_object(key, delta.into()),
864 AssocObj::None,
865 );
866 }
867 }
868
869 pub fn remove(&mut self, object_id: u64, mutation: Mutation) {
871 let txn_mutation = TxnMutation { object_id, mutation, associated_object: AssocObj::None };
872 if self.mutations.remove(&txn_mutation) {
873 if let Mutation::ObjectStore(ObjectStoreMutation {
874 item:
875 ObjectItem {
876 key: ObjectKey { object_id: new_object_id, data: ObjectKeyData::Object },
877 ..
878 },
879 op: Operation::Insert,
880 }) = txn_mutation.mutation
881 {
882 self.new_objects.remove(&(object_id, new_object_id));
883 }
884 }
885 }
886
887 pub fn add_with_object(
890 &mut self,
891 object_id: u64,
892 mutation: Mutation,
893 associated_object: AssocObj<'a>,
894 ) -> Option<Mutation> {
895 debug_assert!(
896 !matches!(
897 mutation,
898 Mutation::ObjectStore(ObjectStoreMutation {
899 item: ObjectItem {
900 key: ObjectKey {
901 data: ObjectKeyData::Project { property: ProjectProperty::Usage, .. },
902 ..
903 },
904 ..
905 },
906 ..
907 })
908 ),
909 "Use merge_bytes_and_nodes"
910 );
911 self.add_with_object_internal(object_id, mutation, associated_object)
912 }
913
914 fn add_with_object_internal(
917 &mut self,
918 object_id: u64,
919 mutation: Mutation,
920 associated_object: AssocObj<'a>,
921 ) -> Option<Mutation> {
922 assert!(object_id != INVALID_OBJECT_ID);
923 let mut is_allocate = false;
924 match &mutation {
925 Mutation::ObjectStore(ObjectStoreMutation {
926 item:
927 Item {
928 key:
929 ObjectKey {
930 data: ObjectKeyData::Attribute(_, AttributeKey::Extent(_)), ..
931 },
932 ..
933 },
934 ..
935 }) => {
936 self.includes_write = true;
937 }
938 Mutation::Allocator(AllocatorMutation::Allocate { .. }) => {
939 is_allocate = true;
940 }
941 _ => {}
942 }
943 let txn_mutation = TxnMutation { object_id, mutation, associated_object };
944 self.verify_locks(&txn_mutation);
945 let old = self.mutations.replace(txn_mutation).map(|m| m.mutation);
946 if is_allocate {
947 assert!(
948 !matches!(self.metadata_reservation, MetadataReservation::BorrowedMetadata)
949 || self.fs.options().image_builder_mode.is_some(),
950 "Allocations are not allowed in BorrowedMetadata transactions"
951 );
952 }
953 old
954 }
955
956 pub fn add_checksum(&mut self, range: Range<u64>, checksums: Vec<Checksum>, first_write: bool) {
957 self.checksums.push((range, checksums, first_write));
958 }
959
960 pub fn includes_write(&self) -> bool {
961 self.includes_write
962 }
963
964 pub fn checksums(&self) -> &[(Range<u64>, Vec<Checksum>, bool)] {
965 &self.checksums
966 }
967
968 pub fn take_checksums(&mut self) -> Vec<(Range<u64>, Vec<Checksum>, bool)> {
969 std::mem::replace(&mut self.checksums, Vec::new())
970 }
971
972 fn verify_locks(&mut self, mutation: &TxnMutation<'_>) {
973 match mutation {
977 TxnMutation {
978 mutation:
979 Mutation::ObjectStore {
980 0: ObjectStoreMutation { item: ObjectItem { key, .. }, op },
981 },
982 object_id: store_object_id,
983 ..
984 } => {
985 match &key.data {
986 ObjectKeyData::Attribute(..) => {
987 }
989 ObjectKeyData::Child { .. }
990 | ObjectKeyData::EncryptedChild(_)
991 | ObjectKeyData::EncryptedCasefoldChild(_)
992 | ObjectKeyData::CasefoldChild { .. }
993 | ObjectKeyData::LegacyCasefoldChild(_) => {
994 let id = key.object_id;
995 if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
996 && !self.new_objects.contains(&(*store_object_id, id))
997 {
998 debug_assert!(
999 false,
1000 "Not holding required lock for object {id} \
1001 in store {store_object_id}"
1002 );
1003 error!(
1004 "Not holding required lock for object {id} in store \
1005 {store_object_id}"
1006 )
1007 }
1008 }
1009 ObjectKeyData::GraveyardEntry { .. } => {
1010 }
1012 ObjectKeyData::GraveyardAttributeEntry { .. } => {
1013 }
1015 ObjectKeyData::Keys => {
1016 let id = key.object_id;
1017 if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
1018 && !self.new_objects.contains(&(*store_object_id, id))
1019 {
1020 debug_assert!(
1021 false,
1022 "Not holding required lock for object {id} \
1023 in store {store_object_id}"
1024 );
1025 error!(
1026 "Not holding required lock for object {id} in store \
1027 {store_object_id}"
1028 )
1029 }
1030 }
1031 ObjectKeyData::Object => match op {
1032 Operation::Insert => {
1034 self.new_objects.insert((*store_object_id, key.object_id));
1035 }
1036 Operation::Merge | Operation::ReplaceOrInsert => {
1037 let id = key.object_id;
1038 if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
1039 && !self.new_objects.contains(&(*store_object_id, id))
1040 {
1041 debug_assert!(
1042 false,
1043 "Not holding required lock for object {id} \
1044 in store {store_object_id}"
1045 );
1046 error!(
1047 "Not holding required lock for object {id} in store \
1048 {store_object_id}"
1049 )
1050 }
1051 }
1052 },
1053 ObjectKeyData::Project { project_id, property: ProjectProperty::Limit } => {
1054 if !self.txn_locks.contains(&LockKey::ProjectId {
1055 store_object_id: *store_object_id,
1056 project_id: *project_id,
1057 }) {
1058 debug_assert!(
1059 false,
1060 "Not holding required lock for project limit id {project_id} \
1061 in store {store_object_id}"
1062 );
1063 error!(
1064 "Not holding required lock for project limit id {project_id} in \
1065 store {store_object_id}"
1066 )
1067 }
1068 }
1069 ObjectKeyData::Project { property: ProjectProperty::Usage, .. } => match op {
1070 Operation::Insert | Operation::ReplaceOrInsert => {
1071 panic!(
1072 "Project usage is all handled by merging deltas, no inserts or \
1073 replacements should be used"
1074 );
1075 }
1076 Operation::Merge => {}
1078 },
1079 ObjectKeyData::ExtendedAttribute { .. } => {
1080 let id = key.object_id;
1081 if !self.txn_locks.contains(&LockKey::object(*store_object_id, id))
1082 && !self.new_objects.contains(&(*store_object_id, id))
1083 {
1084 debug_assert!(
1085 false,
1086 "Not holding required lock for object {id} \
1087 in store {store_object_id} while mutating extended attribute"
1088 );
1089 error!(
1090 "Not holding required lock for object {id} in store \
1091 {store_object_id} while mutating extended attribute"
1092 )
1093 }
1094 }
1095 }
1096 }
1097 TxnMutation { mutation: Mutation::DeleteVolume, object_id, .. } => {
1098 if !self.txn_locks.contains(&LockKey::flush(*object_id)) {
1099 debug_assert!(false, "Not holding required lock for DeleteVolume");
1100 error!("Not holding required lock for DeleteVolume");
1101 }
1102 }
1103 _ => {}
1104 }
1105 }
1106
1107 pub fn is_empty(&self) -> bool {
1109 self.mutations.is_empty()
1110 }
1111
1112 pub fn get_object_mutation(
1115 &self,
1116 store_object_id: u64,
1117 key: ObjectKey,
1118 ) -> Option<&ObjectStoreMutation> {
1119 if let Some(TxnMutation { mutation: Mutation::ObjectStore(mutation), .. }) =
1120 self.mutations.get(&TxnMutation {
1121 object_id: store_object_id,
1122 mutation: Mutation::insert_object(key, ObjectValue::None),
1123 associated_object: AssocObj::None,
1124 })
1125 {
1126 Some(mutation)
1127 } else {
1128 None
1129 }
1130 }
1131
1132 pub async fn commit(mut self) -> Result<u64, Error> {
1134 debug!(txn:? = &self; "Commit");
1135 self.fs.clone().commit_transaction(&mut self, |x| x).await
1136 }
1137
1138 pub async fn commit_with_callback<R: Send>(
1141 mut self,
1142 f: impl FnOnce(u64) -> R + Send,
1143 ) -> Result<R, Error> {
1144 debug!(txn:? = &self; "Commit");
1145 self.fs.clone().commit_transaction(&mut self, f).await
1146 }
1147
1148 pub async fn commit_and_continue(&mut self) -> Result<(), Error> {
1151 debug!(txn:? = self; "Commit");
1152 self.fs.clone().commit_transaction(self, |_| {}).await?;
1153 assert!(self.mutations.is_empty());
1154 assert!(self.new_objects.is_empty());
1155 assert!(self.checksums.is_empty());
1156 self.includes_write = false;
1157 self.fs.lock_manager().downgrade_locks(&self.txn_locks);
1158 self.fs.clone().add_transaction_reservation(self).await
1159 }
1160
1161 pub async fn commit_prepare(&self) {
1164 self.fs.lock_manager().commit_prepare(self).await;
1165 }
1166}
1167
1168impl Drop for Transaction<'_> {
1169 fn drop(&mut self) {
1170 debug!(txn:? = &self; "Drop");
1173 if self.parent_transaction.is_none() {
1174 self.fs.sub_transaction();
1175 }
1176 self.fs.clone().drop_transaction(self);
1177 }
1178}
1179
1180impl std::fmt::Debug for Transaction<'_> {
1181 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1182 f.debug_struct("Transaction")
1183 .field("mutations", &self.mutations)
1184 .field("txn_locks", &self.txn_locks)
1185 .field("reservation", &self.allocator_reservation())
1186 .finish()
1187 }
1188}
1189
1190pub enum BorrowedOrOwned<'a, T> {
1191 Borrowed(&'a T),
1192 Owned(T),
1193}
1194
1195impl<T> Deref for BorrowedOrOwned<'_, T> {
1196 type Target = T;
1197
1198 fn deref(&self) -> &Self::Target {
1199 match self {
1200 BorrowedOrOwned::Borrowed(b) => b,
1201 BorrowedOrOwned::Owned(o) => &o,
1202 }
1203 }
1204}
1205
1206impl<'a, T> From<&'a T> for BorrowedOrOwned<'a, T> {
1207 fn from(value: &'a T) -> Self {
1208 BorrowedOrOwned::Borrowed(value)
1209 }
1210}
1211
1212impl<T> From<T> for BorrowedOrOwned<'_, T> {
1213 fn from(value: T) -> Self {
1214 BorrowedOrOwned::Owned(value)
1215 }
1216}
1217
1218pub struct LockManager {
1247 locks: Mutex<Locks>,
1248}
1249
1250struct Locks {
1251 keys: HashMap<LockKey, LockEntry>,
1252}
1253
1254impl Locks {
1255 fn drop_lock(&mut self, key: LockKey, state: LockState) {
1256 if let Entry::Occupied(mut occupied) = self.keys.entry(key) {
1257 let entry = occupied.get_mut();
1258 let wake = match state {
1259 LockState::ReadLock => {
1260 entry.read_count -= 1;
1261 entry.read_count == 0
1262 }
1263 LockState::Locked | LockState::WriteLock => {
1265 entry.state = LockState::ReadLock;
1266 true
1267 }
1268 };
1269 if wake {
1270 unsafe {
1272 entry.wake();
1273 }
1274 if entry.can_remove() {
1275 occupied.remove_entry();
1276 }
1277 }
1278 } else {
1279 unreachable!();
1280 }
1281 }
1282
1283 fn drop_read_locks(&mut self, lock_keys: LockKeys) {
1284 for lock in lock_keys.iter() {
1285 self.drop_lock(*lock, LockState::ReadLock);
1286 }
1287 }
1288
1289 fn drop_write_locks(&mut self, lock_keys: LockKeys) {
1290 for lock in lock_keys.iter() {
1291 self.drop_lock(*lock, LockState::WriteLock);
1294 }
1295 }
1296
1297 fn downgrade_locks(&mut self, lock_keys: &LockKeys) {
1299 for lock in lock_keys.iter() {
1300 unsafe {
1302 self.keys.get_mut(lock).unwrap().downgrade_lock();
1303 }
1304 }
1305 }
1306}
1307
1308#[derive(Debug)]
1309struct LockEntry {
1310 read_count: u64,
1313
1314 state: LockState,
1316
1317 head: *const LockWaker,
1322 tail: *const LockWaker,
1323}
1324
1325unsafe impl Send for LockEntry {}
1326
1327struct LockWaker {
1330 next: UnsafeCell<*const LockWaker>,
1332 prev: UnsafeCell<*const LockWaker>,
1333
1334 key: LockKey,
1337
1338 waker: UnsafeCell<WakerState>,
1340
1341 target_state: LockState,
1343
1344 is_upgrade: bool,
1346
1347 _pin: PhantomPinned,
1349}
1350
1351enum WakerState {
1352 Pending,
1354
1355 Registered(Waker),
1357
1358 Woken,
1360}
1361
1362impl WakerState {
1363 fn is_woken(&self) -> bool {
1364 matches!(self, WakerState::Woken)
1365 }
1366}
1367
1368unsafe impl Send for LockWaker {}
1369unsafe impl Sync for LockWaker {}
1370
1371impl LockWaker {
1372 async fn wait(&self, manager: &LockManager) {
1374 let waker_guard = scopeguard::guard((), |_| {
1376 let mut locks = manager.locks.lock();
1377 unsafe {
1379 if (*self.waker.get()).is_woken() {
1380 if self.is_upgrade {
1382 locks.keys.get_mut(&self.key).unwrap().downgrade_lock();
1383 } else {
1384 locks.drop_lock(self.key, self.target_state);
1385 }
1386 } else {
1387 locks.keys.get_mut(&self.key).unwrap().remove_waker(self);
1390 }
1391 }
1392 });
1393
1394 poll_fn(|cx| {
1395 let _locks = manager.locks.lock();
1396 unsafe {
1398 if (*self.waker.get()).is_woken() {
1399 Poll::Ready(())
1400 } else {
1401 *self.waker.get() = WakerState::Registered(cx.waker().clone());
1402 Poll::Pending
1403 }
1404 }
1405 })
1406 .await;
1407
1408 ScopeGuard::into_inner(waker_guard);
1409 }
1410}
1411
1412#[derive(Copy, Clone, Debug, PartialEq)]
1413enum LockState {
1414 ReadLock,
1416
1417 Locked,
1420
1421 WriteLock,
1423}
1424
1425impl LockManager {
1426 pub fn new() -> Self {
1427 LockManager { locks: Mutex::new(Locks { keys: HashMap::default() }) }
1428 }
1429
1430 pub async fn txn_lock<'a>(&'a self, lock_keys: LockKeys) -> TransactionLocks<'a> {
1434 TransactionLocks(
1435 debug_assert_not_too_long!(self.lock(lock_keys, LockState::Locked)).right().unwrap(),
1436 )
1437 }
1438
1439 async fn lock<'a>(
1442 &'a self,
1443 mut lock_keys: LockKeys,
1444 target_state: LockState,
1445 ) -> Either<ReadGuard<'a>, WriteGuard<'a>> {
1446 let mut guard = match &target_state {
1447 LockState::ReadLock => Left(ReadGuard {
1448 manager: self.into(),
1449 lock_keys: LockKeys::with_capacity(lock_keys.len()),
1450 }),
1451 LockState::Locked | LockState::WriteLock => Right(WriteGuard {
1452 manager: self.into(),
1453 lock_keys: LockKeys::with_capacity(lock_keys.len()),
1454 }),
1455 };
1456 let guard_keys = match &mut guard {
1457 Left(g) => &mut g.lock_keys,
1458 Right(g) => &mut g.lock_keys,
1459 };
1460 lock_keys.sort_unstable();
1461 lock_keys.dedup();
1462 for lock in lock_keys.iter() {
1463 let lock_waker = None;
1464 pin_mut!(lock_waker);
1465 {
1466 let mut locks = self.locks.lock();
1467 match locks.keys.entry(*lock) {
1468 Entry::Vacant(vacant) => {
1469 vacant.insert(LockEntry {
1470 read_count: if let LockState::ReadLock = target_state {
1471 guard_keys.push(*lock);
1472 1
1473 } else {
1474 guard_keys.push(*lock);
1475 0
1476 },
1477 state: target_state,
1478 head: std::ptr::null(),
1479 tail: std::ptr::null(),
1480 });
1481 }
1482 Entry::Occupied(mut occupied) => {
1483 let entry = occupied.get_mut();
1484 if unsafe { entry.is_allowed(target_state, entry.head.is_null()) } {
1486 if let LockState::ReadLock = target_state {
1487 entry.read_count += 1;
1488 guard_keys.push(*lock);
1489 } else {
1490 entry.state = target_state;
1491 guard_keys.push(*lock);
1492 }
1493 } else {
1494 unsafe {
1497 *lock_waker.as_mut().get_unchecked_mut() = Some(LockWaker {
1498 next: UnsafeCell::new(std::ptr::null()),
1499 prev: UnsafeCell::new(entry.tail),
1500 key: *lock,
1501 waker: UnsafeCell::new(WakerState::Pending),
1502 target_state: target_state,
1503 is_upgrade: false,
1504 _pin: PhantomPinned,
1505 });
1506 }
1507 let waker = (*lock_waker).as_ref().unwrap();
1508 if entry.tail.is_null() {
1509 entry.head = waker;
1510 } else {
1511 unsafe {
1513 *(*entry.tail).next.get() = waker;
1514 }
1515 }
1516 entry.tail = waker;
1517 }
1518 }
1519 }
1520 }
1521 if let Some(waker) = &*lock_waker {
1522 waker.wait(self).await;
1523 guard_keys.push(*lock);
1524 }
1525 }
1526 guard
1527 }
1528
1529 pub fn drop_transaction(&self, transaction: &mut Transaction<'_>) {
1531 let mut locks = self.locks.lock();
1532 locks.drop_write_locks(std::mem::take(&mut transaction.txn_locks));
1533 }
1534
1535 pub async fn commit_prepare(&self, transaction: &Transaction<'_>) {
1537 self.commit_prepare_keys(&transaction.txn_locks).await;
1538 }
1539
1540 async fn commit_prepare_keys(&self, lock_keys: &LockKeys) {
1541 for lock in lock_keys.iter() {
1542 let lock_waker = None;
1543 pin_mut!(lock_waker);
1544 {
1545 let mut locks = self.locks.lock();
1546 let entry = locks.keys.get_mut(lock).unwrap();
1547 if entry.state == LockState::WriteLock {
1552 continue;
1553 }
1554 assert_eq!(entry.state, LockState::Locked);
1555
1556 if entry.read_count == 0 {
1557 entry.state = LockState::WriteLock;
1558 } else {
1559 unsafe {
1562 *lock_waker.as_mut().get_unchecked_mut() = Some(LockWaker {
1563 next: UnsafeCell::new(entry.head),
1564 prev: UnsafeCell::new(std::ptr::null()),
1565 key: *lock,
1566 waker: UnsafeCell::new(WakerState::Pending),
1567 target_state: LockState::WriteLock,
1568 is_upgrade: true,
1569 _pin: PhantomPinned,
1570 });
1571 }
1572 let waker = (*lock_waker).as_ref().unwrap();
1573 if entry.head.is_null() {
1574 entry.tail = (*lock_waker).as_ref().unwrap();
1575 } else {
1576 unsafe {
1578 *(*entry.head).prev.get() = waker;
1579 }
1580 }
1581 entry.head = waker;
1582 }
1583 }
1584
1585 if let Some(waker) = &*lock_waker {
1586 waker.wait(self).await;
1587 }
1588 }
1589 }
1590
1591 pub async fn read_lock<'a>(&'a self, lock_keys: LockKeys) -> ReadGuard<'a> {
1597 debug_assert_not_too_long!(self.lock(lock_keys, LockState::ReadLock)).left().unwrap()
1598 }
1599
1600 pub async fn write_lock<'a>(&'a self, lock_keys: LockKeys) -> WriteGuard<'a> {
1603 debug_assert_not_too_long!(self.lock(lock_keys, LockState::WriteLock)).right().unwrap()
1604 }
1605
1606 pub fn downgrade_locks(&self, lock_keys: &LockKeys) {
1609 self.locks.lock().downgrade_locks(lock_keys);
1610 }
1611}
1612
1613impl LockEntry {
1615 unsafe fn wake(&mut self) {
1616 if self.head.is_null() || self.state == LockState::WriteLock {
1618 return;
1619 }
1620
1621 let waker = unsafe { &*self.head };
1622
1623 if waker.is_upgrade {
1624 if self.read_count > 0 {
1625 return;
1626 }
1627 } else if !unsafe { self.is_allowed(waker.target_state, true) } {
1628 return;
1629 }
1630
1631 unsafe { self.pop_and_wake() };
1632
1633 if waker.target_state == LockState::WriteLock {
1636 return;
1637 }
1638
1639 while !self.head.is_null() && unsafe { (*self.head).target_state } == LockState::ReadLock {
1640 unsafe { self.pop_and_wake() };
1641 }
1642 }
1643
1644 unsafe fn pop_and_wake(&mut self) {
1645 let waker = unsafe { &*self.head };
1646
1647 self.head = unsafe { *waker.next.get() };
1649 if self.head.is_null() {
1650 self.tail = std::ptr::null()
1651 } else {
1652 unsafe { *(*self.head).prev.get() = std::ptr::null() };
1653 }
1654
1655 if waker.target_state == LockState::ReadLock {
1657 self.read_count += 1;
1658 } else {
1659 self.state = waker.target_state;
1660 }
1661
1662 if let WakerState::Registered(waker) =
1664 std::mem::replace(unsafe { &mut *waker.waker.get() }, WakerState::Woken)
1665 {
1666 waker.wake();
1667 }
1668 }
1669
1670 fn can_remove(&self) -> bool {
1671 self.state == LockState::ReadLock && self.read_count == 0
1672 }
1673
1674 unsafe fn remove_waker(&mut self, waker: &LockWaker) {
1675 unsafe {
1676 let is_first = (*waker.prev.get()).is_null();
1677 if is_first {
1678 self.head = *waker.next.get();
1679 } else {
1680 *(**waker.prev.get()).next.get() = *waker.next.get();
1681 }
1682 if (*waker.next.get()).is_null() {
1683 self.tail = *waker.prev.get();
1684 } else {
1685 *(**waker.next.get()).prev.get() = *waker.prev.get();
1686 }
1687 if is_first {
1688 self.wake();
1691 }
1692 }
1693 }
1694
1695 unsafe fn is_allowed(&self, target_state: LockState, is_head: bool) -> bool {
1699 match self.state {
1700 LockState::ReadLock => {
1701 (self.read_count == 0
1703 || target_state == LockState::Locked
1704 || target_state == LockState::ReadLock)
1705 && is_head
1706 }
1707 LockState::Locked => {
1708 target_state == LockState::ReadLock
1712 && (is_head || unsafe { !(*self.head).is_upgrade })
1713 }
1714 LockState::WriteLock => false,
1715 }
1716 }
1717
1718 unsafe fn downgrade_lock(&mut self) {
1719 assert_eq!(std::mem::replace(&mut self.state, LockState::Locked), LockState::WriteLock);
1720 unsafe { self.wake() };
1721 }
1722}
1723
1724#[must_use]
1725pub struct ReadGuard<'a> {
1726 manager: LockManagerRef<'a>,
1727 lock_keys: LockKeys,
1728}
1729
1730impl ReadGuard<'_> {
1731 pub fn fs(&self) -> Option<&Arc<FxFilesystem>> {
1732 if let LockManagerRef::Owned(fs) = &self.manager { Some(fs) } else { None }
1733 }
1734
1735 pub fn into_owned(mut self, fs: Arc<FxFilesystem>) -> ReadGuard<'static> {
1736 ReadGuard {
1737 manager: LockManagerRef::Owned(fs),
1738 lock_keys: std::mem::replace(&mut self.lock_keys, LockKeys::None),
1739 }
1740 }
1741}
1742
1743impl Drop for ReadGuard<'_> {
1744 fn drop(&mut self) {
1745 let mut locks = self.manager.locks.lock();
1746 locks.drop_read_locks(std::mem::take(&mut self.lock_keys));
1747 }
1748}
1749
1750impl fmt::Debug for ReadGuard<'_> {
1751 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1752 f.debug_struct("ReadGuard")
1753 .field("manager", &(&self.manager as *const _))
1754 .field("lock_keys", &self.lock_keys)
1755 .finish()
1756 }
1757}
1758
1759#[must_use]
1760pub struct WriteGuard<'a> {
1761 manager: LockManagerRef<'a>,
1762 lock_keys: LockKeys,
1763}
1764
1765impl WriteGuard<'_> {
1766 pub fn into_owned(mut self, fs: Arc<FxFilesystem>) -> WriteGuard<'static> {
1767 WriteGuard {
1768 manager: LockManagerRef::Owned(fs),
1769 lock_keys: std::mem::replace(&mut self.lock_keys, LockKeys::None),
1770 }
1771 }
1772}
1773
1774impl Drop for WriteGuard<'_> {
1775 fn drop(&mut self) {
1776 let mut locks = self.manager.locks.lock();
1777 locks.drop_write_locks(std::mem::take(&mut self.lock_keys));
1778 }
1779}
1780
1781impl fmt::Debug for WriteGuard<'_> {
1782 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1783 f.debug_struct("WriteGuard")
1784 .field("manager", &(&self.manager as *const _))
1785 .field("lock_keys", &self.lock_keys)
1786 .finish()
1787 }
1788}
1789
1790enum LockManagerRef<'a> {
1791 Borrowed(&'a LockManager),
1792 Owned(Arc<FxFilesystem>),
1793}
1794
1795impl Deref for LockManagerRef<'_> {
1796 type Target = LockManager;
1797
1798 fn deref(&self) -> &Self::Target {
1799 match self {
1800 LockManagerRef::Borrowed(m) => m,
1801 LockManagerRef::Owned(f) => f.lock_manager(),
1802 }
1803 }
1804}
1805
1806impl<'a> From<&'a LockManager> for LockManagerRef<'a> {
1807 fn from(value: &'a LockManager) -> Self {
1808 LockManagerRef::Borrowed(value)
1809 }
1810}
1811
1812#[cfg(test)]
1813mod tests {
1814 use super::{
1815 AssocObj, AttributeId, LockKey, LockKeys, LockManager, LockState, Mutation,
1816 ObjectMutationIterator, Options, ReservationOptions, TxnMutation,
1817 };
1818 use crate::filesystem::FxFilesystem;
1819 use crate::object_store::{BytesAndNodes, ObjectKey};
1820 use fuchsia_async as fasync;
1821 use fuchsia_sync::Mutex;
1822 use futures::channel::oneshot::channel;
1823 use futures::future::FutureExt;
1824 use futures::stream::FuturesUnordered;
1825 use futures::{StreamExt, join, pin_mut};
1826 use std::collections::BTreeSet;
1827 use std::task::Poll;
1828 use std::time::Duration;
1829 use storage_device::DeviceHolder;
1830 use storage_device::fake_device::FakeDevice;
1831
1832 #[fuchsia::test]
1833 async fn test_simple() {
1834 let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1835 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1836 let mut t = fs
1837 .root_store()
1838 .new_transaction(lock_keys![], Options::default())
1839 .await
1840 .expect("new_transaction failed");
1841 t.add(1, Mutation::BeginFlush);
1842 assert!(!t.is_empty());
1843 }
1844
1845 #[fuchsia::test]
1846 async fn test_locks() {
1847 let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1848 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1849 let (send1, recv1) = channel();
1850 let (send2, recv2) = channel();
1851 let (send3, recv3) = channel();
1852 let done = Mutex::new(false);
1853 let mut futures = FuturesUnordered::new();
1854 futures.push(
1855 async {
1856 let _t = fs
1857 .root_store()
1858 .new_transaction(
1859 lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1860 Options::default(),
1861 )
1862 .await
1863 .expect("new_transaction failed");
1864 send1.send(()).unwrap(); send3.send(()).unwrap(); recv2.await.unwrap();
1867 fasync::Timer::new(Duration::from_millis(100)).await;
1869 assert!(!*done.lock());
1870 }
1871 .boxed(),
1872 );
1873 futures.push(
1874 async {
1875 recv1.await.unwrap();
1876 let _t = fs
1878 .root_store()
1879 .new_transaction(
1880 lock_keys![LockKey::object_attribute(2, 2, AttributeId::TEST_ID)],
1881 Options::default(),
1882 )
1883 .await
1884 .expect("new_transaction failed");
1885 send2.send(()).unwrap();
1887 }
1888 .boxed(),
1889 );
1890 futures.push(
1891 async {
1892 recv3.await.unwrap();
1894 let _t = fs
1895 .root_store()
1896 .new_transaction(
1897 lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1898 Options::default(),
1899 )
1900 .await;
1901 *done.lock() = true;
1902 }
1903 .boxed(),
1904 );
1905 while let Some(()) = futures.next().await {}
1906 }
1907
1908 #[fuchsia::test]
1909 async fn test_read_lock_after_write_lock() {
1910 let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1911 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1912 let (send1, recv1) = channel();
1913 let (send2, recv2) = channel();
1914 let done = Mutex::new(false);
1915 join!(
1916 async {
1917 let t = fs
1918 .root_store()
1919 .new_transaction(
1920 lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1921 Options::default(),
1922 )
1923 .await
1924 .expect("new_transaction failed");
1925 send1.send(()).unwrap(); recv2.await.unwrap();
1927 t.commit().await.expect("commit failed");
1928 *done.lock() = true;
1929 },
1930 async {
1931 recv1.await.unwrap();
1932 let _guard = fs
1934 .lock_manager()
1935 .read_lock(lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)])
1936 .await;
1937 send2.send(()).unwrap();
1939 fasync::Timer::new(Duration::from_millis(100)).await;
1942 assert!(!*done.lock());
1943 },
1944 );
1945 }
1946
1947 #[fuchsia::test]
1948 async fn test_write_lock_after_read_lock() {
1949 let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1950 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1951 let (send1, recv1) = channel();
1952 let (send2, recv2) = channel();
1953 let done = Mutex::new(false);
1954 join!(
1955 async {
1956 let _guard = fs
1958 .lock_manager()
1959 .read_lock(lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)])
1960 .await;
1961 send1.send(()).unwrap();
1963 recv2.await.unwrap();
1964 fasync::Timer::new(Duration::from_millis(100)).await;
1967 assert!(!*done.lock());
1968 },
1969 async {
1970 recv1.await.unwrap();
1971 let t = fs
1972 .root_store()
1973 .new_transaction(
1974 lock_keys![LockKey::object_attribute(1, 2, AttributeId::TEST_ID)],
1975 Options::default(),
1976 )
1977 .await
1978 .expect("new_transaction failed");
1979 send2.send(()).unwrap(); t.commit().await.expect("commit failed");
1981 *done.lock() = true;
1982 },
1983 );
1984 }
1985
1986 #[fuchsia::test]
1987 async fn test_drop_uncommitted_transaction() {
1988 let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
1989 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1990 let key = lock_keys![LockKey::object(1, 1)];
1991
1992 {
1994 let _write_lock = fs
1995 .root_store()
1996 .new_transaction(key.clone(), Options::default())
1997 .await
1998 .expect("new_transaction failed");
1999 let _read_lock = fs.lock_manager().read_lock(key.clone()).await;
2000 }
2001 {
2003 let _write_lock = fs
2004 .root_store()
2005 .new_transaction(key.clone(), Options::default())
2006 .await
2007 .expect("new_transaction failed");
2008 }
2009 fs.root_store()
2011 .new_transaction(key.clone(), Options::default())
2012 .await
2013 .expect("new_transaction failed");
2014 }
2015
2016 #[fuchsia::test]
2017 async fn test_drop_waiting_write_lock() {
2018 let manager = LockManager::new();
2019 let keys = lock_keys![LockKey::object(1, 1)];
2020 {
2021 let _guard = manager.lock(keys.clone(), LockState::ReadLock).await;
2022 if let Poll::Ready(_) =
2023 futures::poll!(manager.lock(keys.clone(), LockState::WriteLock).boxed())
2024 {
2025 assert!(false);
2026 }
2027 }
2028 let _ = manager.lock(keys, LockState::WriteLock).await;
2029 }
2030
2031 #[fuchsia::test]
2032 async fn test_write_lock_blocks_everything() {
2033 let manager = LockManager::new();
2034 let keys = lock_keys![LockKey::object(1, 1)];
2035 {
2036 let _guard = manager.lock(keys.clone(), LockState::WriteLock).await;
2037 if let Poll::Ready(_) =
2038 futures::poll!(manager.lock(keys.clone(), LockState::WriteLock).boxed())
2039 {
2040 assert!(false);
2041 }
2042 if let Poll::Ready(_) =
2043 futures::poll!(manager.lock(keys.clone(), LockState::ReadLock).boxed())
2044 {
2045 assert!(false);
2046 }
2047 }
2048 {
2049 let _guard = manager.lock(keys.clone(), LockState::WriteLock).await;
2050 }
2051 {
2052 let _guard = manager.lock(keys, LockState::ReadLock).await;
2053 }
2054 }
2055
2056 #[fuchsia::test]
2057 async fn test_downgrade_locks() {
2058 let manager = LockManager::new();
2059 let keys = lock_keys![LockKey::object(1, 1)];
2060 let _guard = manager.txn_lock(keys.clone()).await;
2061 manager.commit_prepare_keys(&keys).await;
2062
2063 let mut read_lock: FuturesUnordered<_> =
2065 std::iter::once(manager.read_lock(keys.clone())).collect();
2066
2067 assert!(futures::poll!(read_lock.next()).is_pending());
2069
2070 manager.downgrade_locks(&keys);
2071
2072 assert!(futures::poll!(read_lock.next()).is_ready());
2074 }
2075
2076 #[fuchsia::test]
2077 async fn test_dropped_write_lock_wakes() {
2078 let manager = LockManager::new();
2079 let keys = lock_keys![LockKey::object(1, 1)];
2080 let _guard = manager.lock(keys.clone(), LockState::ReadLock).await;
2081 let mut read_lock = FuturesUnordered::new();
2082 read_lock.push(manager.lock(keys.clone(), LockState::ReadLock));
2083
2084 {
2085 let write_lock = manager.lock(keys, LockState::WriteLock);
2086 pin_mut!(write_lock);
2087
2088 assert!(futures::poll!(write_lock).is_pending());
2090
2091 assert!(futures::poll!(read_lock.next()).is_pending());
2093 }
2094
2095 assert!(futures::poll!(read_lock.next()).is_ready());
2097 }
2098
2099 #[fuchsia::test]
2100 async fn test_drop_upgrade() {
2101 let manager = LockManager::new();
2102 let keys = lock_keys![LockKey::object(1, 1)];
2103 let _guard = manager.lock(keys.clone(), LockState::Locked).await;
2104
2105 {
2106 let commit_prepare = manager.commit_prepare_keys(&keys);
2107 pin_mut!(commit_prepare);
2108 let _read_guard = manager.lock(keys.clone(), LockState::ReadLock).await;
2109 assert!(futures::poll!(commit_prepare).is_pending());
2110
2111 }
2114
2115 manager.commit_prepare_keys(&keys).await;
2117 }
2118
2119 #[fuchsia::test]
2120 async fn test_woken_upgrade_blocks_reads() {
2121 let manager = LockManager::new();
2122 let keys = lock_keys![LockKey::object(1, 1)];
2123 let guard = manager.lock(keys.clone(), LockState::Locked).await;
2125
2126 let read1 = manager.lock(keys.clone(), LockState::ReadLock).await;
2128
2129 let commit_prepare = manager.commit_prepare_keys(&keys);
2131 pin_mut!(commit_prepare);
2132 assert!(futures::poll!(commit_prepare.as_mut()).is_pending());
2133
2134 let read2 = manager.lock(keys.clone(), LockState::ReadLock);
2136 pin_mut!(read2);
2137 assert!(futures::poll!(read2.as_mut()).is_pending());
2138
2139 std::mem::drop(read1);
2141 assert!(futures::poll!(commit_prepare).is_ready());
2142
2143 assert!(futures::poll!(read2.as_mut()).is_pending());
2145
2146 std::mem::drop(guard);
2148 assert!(futures::poll!(read2).is_ready());
2149 }
2150
2151 static LOCK_KEY_1: LockKey = LockKey::flush(1);
2152 static LOCK_KEY_2: LockKey = LockKey::flush(2);
2153 static LOCK_KEY_3: LockKey = LockKey::flush(3);
2154
2155 fn assert_lock_keys_equal(value: &LockKeys, expected: &LockKeys) {
2157 match (value, expected) {
2158 (LockKeys::None, LockKeys::None) => {}
2159 (LockKeys::Inline(key1), LockKeys::Inline(key2)) => {
2160 if key1 != key2 {
2161 panic!("{key1:?} != {key2:?}");
2162 }
2163 }
2164 (LockKeys::Vec(vec1), LockKeys::Vec(vec2)) => {
2165 if vec1 != vec2 {
2166 panic!("{vec1:?} != {vec2:?}");
2167 }
2168 if vec1.capacity() != vec2.capacity() {
2169 panic!(
2170 "LockKeys have different capacity: {} != {}",
2171 vec1.capacity(),
2172 vec2.capacity()
2173 );
2174 }
2175 }
2176 (_, _) => panic!("{value:?} != {expected:?}"),
2177 }
2178 }
2179
2180 fn assert_lock_keys_equivalent(value: &LockKeys, expected: &LockKeys) {
2182 let value: Vec<_> = value.iter().collect();
2183 let expected: Vec<_> = expected.iter().collect();
2184 assert_eq!(value, expected);
2185 }
2186
2187 #[test]
2188 fn test_lock_keys_macro() {
2189 assert_lock_keys_equal(&lock_keys![], &LockKeys::None);
2190 assert_lock_keys_equal(&lock_keys![LOCK_KEY_1], &LockKeys::Inline(LOCK_KEY_1));
2191 assert_lock_keys_equal(
2192 &lock_keys![LOCK_KEY_1, LOCK_KEY_2],
2193 &LockKeys::Vec(vec![LOCK_KEY_1, LOCK_KEY_2]),
2194 );
2195 }
2196
2197 #[test]
2198 fn test_lock_keys_with_capacity() {
2199 assert_lock_keys_equal(&LockKeys::with_capacity(0), &LockKeys::None);
2200 assert_lock_keys_equal(&LockKeys::with_capacity(1), &LockKeys::None);
2201 assert_lock_keys_equal(&LockKeys::with_capacity(2), &LockKeys::Vec(Vec::with_capacity(2)));
2202 }
2203
2204 #[test]
2205 fn test_lock_keys_len() {
2206 assert_eq!(lock_keys![].len(), 0);
2207 assert_eq!(lock_keys![LOCK_KEY_1].len(), 1);
2208 assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].len(), 2);
2209 }
2210
2211 #[test]
2212 fn test_lock_keys_contains() {
2213 assert_eq!(lock_keys![].contains(&LOCK_KEY_1), false);
2214 assert_eq!(lock_keys![LOCK_KEY_1].contains(&LOCK_KEY_1), true);
2215 assert_eq!(lock_keys![LOCK_KEY_1].contains(&LOCK_KEY_2), false);
2216 assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].contains(&LOCK_KEY_1), true);
2217 assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].contains(&LOCK_KEY_2), true);
2218 assert_eq!(lock_keys![LOCK_KEY_1, LOCK_KEY_2].contains(&LOCK_KEY_3), false);
2219 }
2220
2221 #[test]
2222 fn test_lock_keys_push() {
2223 let mut keys = lock_keys![];
2224 keys.push(LOCK_KEY_1);
2225 assert_lock_keys_equal(&keys, &LockKeys::Inline(LOCK_KEY_1));
2226 keys.push(LOCK_KEY_2);
2227 assert_lock_keys_equal(&keys, &LockKeys::Vec(vec![LOCK_KEY_1, LOCK_KEY_2]));
2228 keys.push(LOCK_KEY_3);
2229 assert_lock_keys_equivalent(
2230 &keys,
2231 &LockKeys::Vec(vec![LOCK_KEY_1, LOCK_KEY_2, LOCK_KEY_3]),
2232 );
2233 }
2234
2235 #[test]
2236 fn test_lock_keys_sort_unstable() {
2237 let mut keys = lock_keys![];
2238 keys.sort_unstable();
2239 assert_lock_keys_equal(&keys, &lock_keys![]);
2240
2241 let mut keys = lock_keys![LOCK_KEY_1];
2242 keys.sort_unstable();
2243 assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1]);
2244
2245 let mut keys = lock_keys![LOCK_KEY_2, LOCK_KEY_1];
2246 keys.sort_unstable();
2247 assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1, LOCK_KEY_2]);
2248 }
2249
2250 #[test]
2251 fn test_lock_keys_dedup() {
2252 let mut keys = lock_keys![];
2253 keys.dedup();
2254 assert_lock_keys_equal(&keys, &lock_keys![]);
2255
2256 let mut keys = lock_keys![LOCK_KEY_1];
2257 keys.dedup();
2258 assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1]);
2259
2260 let mut keys = lock_keys![LOCK_KEY_1, LOCK_KEY_1];
2261 keys.dedup();
2262 assert_lock_keys_equivalent(&keys, &lock_keys![LOCK_KEY_1]);
2263 }
2264
2265 #[test]
2266 fn test_lock_keys_truncate() {
2267 let mut keys = lock_keys![];
2268 keys.truncate(5);
2269 assert_lock_keys_equal(&keys, &lock_keys![]);
2270 keys.truncate(0);
2271 assert_lock_keys_equal(&keys, &lock_keys![]);
2272
2273 let mut keys = lock_keys![LOCK_KEY_1];
2274 keys.truncate(5);
2275 assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1]);
2276 keys.truncate(0);
2277 assert_lock_keys_equal(&keys, &lock_keys![]);
2278
2279 let mut keys = lock_keys![LOCK_KEY_1, LOCK_KEY_2];
2280 keys.truncate(5);
2281 assert_lock_keys_equal(&keys, &lock_keys![LOCK_KEY_1, LOCK_KEY_2]);
2282 keys.truncate(1);
2283 assert_lock_keys_equivalent(&keys, &lock_keys![LOCK_KEY_1]);
2285 }
2286
2287 #[test]
2288 fn test_lock_keys_iter() {
2289 assert_eq!(lock_keys![].iter().collect::<Vec<_>>(), Vec::<&LockKey>::new());
2290
2291 assert_eq!(lock_keys![LOCK_KEY_1].iter().collect::<Vec<_>>(), vec![&LOCK_KEY_1]);
2292
2293 assert_eq!(
2294 lock_keys![LOCK_KEY_1, LOCK_KEY_2].iter().collect::<Vec<_>>(),
2295 vec![&LOCK_KEY_1, &LOCK_KEY_2]
2296 );
2297 }
2298
2299 #[test]
2300 fn test_object_mutation_iterator() {
2301 let mut mutations = BTreeSet::new();
2302 mutations.insert(TxnMutation {
2303 object_id: 1,
2304 mutation: Mutation::BeginFlush,
2305 associated_object: AssocObj::None,
2306 });
2307 mutations.insert(TxnMutation {
2308 object_id: 1,
2309 mutation: Mutation::EndFlush,
2310 associated_object: AssocObj::None,
2311 });
2312 mutations.insert(TxnMutation {
2313 object_id: 2,
2314 mutation: Mutation::DeleteVolume,
2315 associated_object: AssocObj::None,
2316 });
2317
2318 let mut iter = mutations.iter().peekable();
2319
2320 {
2321 let mut obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 1");
2322 assert_eq!(obj_iter.object_id(), 1);
2323 assert_eq!(obj_iter.next(), Some(&Mutation::BeginFlush));
2324 assert_eq!(obj_iter.next(), Some(&Mutation::EndFlush));
2325 assert_eq!(obj_iter.next(), None);
2326 }
2327
2328 {
2329 let mut obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 2");
2330 assert_eq!(obj_iter.object_id(), 2);
2331 assert_eq!(obj_iter.next(), Some(&Mutation::DeleteVolume));
2332 assert_eq!(obj_iter.next(), None);
2333 }
2334
2335 assert!(ObjectMutationIterator::new(&mut iter).is_none());
2337 }
2338
2339 #[test]
2340 fn test_object_mutation_iterator_drop_drains_remaining() {
2341 let mut mutations = BTreeSet::new();
2342 mutations.insert(TxnMutation {
2344 object_id: 1,
2345 mutation: Mutation::BeginFlush,
2346 associated_object: AssocObj::None,
2347 });
2348 mutations.insert(TxnMutation {
2349 object_id: 1,
2350 mutation: Mutation::EndFlush,
2351 associated_object: AssocObj::None,
2352 });
2353 mutations.insert(TxnMutation {
2354 object_id: 1,
2355 mutation: Mutation::DeleteVolume,
2356 associated_object: AssocObj::None,
2357 });
2358 mutations.insert(TxnMutation {
2360 object_id: 2,
2361 mutation: Mutation::BeginFlush,
2362 associated_object: AssocObj::None,
2363 });
2364 mutations.insert(TxnMutation {
2365 object_id: 2,
2366 mutation: Mutation::EndFlush,
2367 associated_object: AssocObj::None,
2368 });
2369
2370 let mut iter = mutations.iter().peekable();
2371
2372 {
2373 let mut obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 1");
2374 assert_eq!(obj_iter.object_id(), 1);
2375 assert_eq!(obj_iter.next(), Some(&Mutation::BeginFlush));
2376 }
2378
2379 {
2380 let obj_iter = ObjectMutationIterator::new(&mut iter).expect("expected object 2");
2381 assert_eq!(obj_iter.object_id(), 2);
2382 }
2384
2385 assert!(ObjectMutationIterator::new(&mut iter).is_none());
2387 }
2388
2389 #[fuchsia::test]
2390 async fn test_merge_bytes_and_nodes() {
2391 let device = DeviceHolder::new(FakeDevice::new(4096, 1024));
2392 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2393 let root_store = fs.root_store();
2394 let store_id = root_store.store_object_id();
2395 let mut t = root_store
2396 .new_transaction(lock_keys![], Options::default())
2397 .await
2398 .expect("new_transaction failed");
2399
2400 let key = ObjectKey::project_usage(
2401 root_store.root_directory_object_id(),
2402 crate::object_store::ProjectId::new(1).unwrap(),
2403 );
2404
2405 t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: 100, nodes: 2 });
2407 let mutation = t.get_object_mutation(store_id, key.clone()).expect("mutation expected");
2408 assert_eq!(mutation.op, super::Operation::Merge);
2409 assert_eq!(
2410 mutation.item.value,
2411 crate::object_store::object_record::ObjectValue::BytesAndNodes { bytes: 100, nodes: 2 }
2412 );
2413
2414 t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: 50, nodes: 1 });
2416 let mutation = t.get_object_mutation(store_id, key.clone()).expect("mutation expected");
2417 assert_eq!(mutation.op, super::Operation::Merge);
2418 assert_eq!(
2419 mutation.item.value,
2420 crate::object_store::object_record::ObjectValue::BytesAndNodes { bytes: 150, nodes: 3 }
2421 );
2422
2423 t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: -50, nodes: -1 });
2425 let mutation = t.get_object_mutation(store_id, key.clone()).expect("mutation expected");
2426 assert_eq!(mutation.op, super::Operation::Merge);
2427 assert_eq!(
2428 mutation.item.value,
2429 crate::object_store::object_record::ObjectValue::BytesAndNodes { bytes: 100, nodes: 2 }
2430 );
2431
2432 t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: -100, nodes: -2 });
2434 assert!(t.get_object_mutation(store_id, key.clone()).is_none());
2435
2436 t.merge_bytes_and_nodes(store_id, key.clone(), BytesAndNodes { bytes: 0, nodes: 0 });
2438 assert!(t.get_object_mutation(store_id, key.clone()).is_none());
2439 }
2440
2441 #[fuchsia::test]
2442 async fn test_commit_and_continue_checks_journal_space() {
2443 use crate::filesystem::FxFilesystemBuilder;
2444 use crate::hooks::Hooks;
2445 use crate::object_store::journal::JournalOptions;
2446 use crate::object_store::volume::root_volume;
2447 use crate::object_store::{NewChildStoreOptions, ObjectKey, ObjectValue};
2448 use storage_device::DeviceHolder;
2449 use storage_device::fake_device::FakeDevice;
2450
2451 let reclaim_size = 65536;
2452 let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
2453 let (mut hooks, fs_hooks) = Hooks::new();
2454 let fs = FxFilesystemBuilder::new()
2455 .hooks(fs_hooks)
2456 .journal_options(JournalOptions { reclaim_size, ..Default::default() })
2457 .format(true)
2458 .open(device)
2459 .await
2460 .expect("open failed");
2461
2462 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
2463 let store = root_volume
2464 .new_volume("test", NewChildStoreOptions::default())
2465 .await
2466 .expect("new_volume failed");
2467
2468 fs.journal().force_compact().await.expect("force_compact failed");
2469 fs.journal().pause_compactions().await;
2470
2471 let journal_clone = fs.journal().clone();
2472 hooks.set_waiting_for_journal_space(move || {
2473 journal_clone.resume_compactions();
2474 });
2475
2476 let _reservation = fs.allocator().reserve_with(None, |limit| limit);
2479
2480 let mut transaction = store
2481 .new_transaction(
2482 lock_keys![LockKey::object(store.store_object_id(), 1000)],
2483 Options { reservation: ReservationOptions::BorrowedMetadata, ..Default::default() },
2484 )
2485 .await
2486 .expect("new_transaction failed");
2487
2488 for _ in 0..64 {
2492 for i in 0..16 {
2493 transaction.add(
2494 store.store_object_id(),
2495 Mutation::replace_or_insert_object(
2496 ObjectKey::extended_attribute(1000, format!("attr_{i}").into_bytes()),
2497 ObjectValue::inline_extended_attribute(vec![0u8; 1024]),
2498 ),
2499 );
2500 }
2501 transaction.commit_and_continue().await.expect("commit_and_continue failed");
2502 }
2503 transaction.commit().await.expect("commit failed");
2504
2505 fs.close().await.expect("close failed");
2506 }
2507
2508 #[fuchsia::test]
2509 async fn test_commit_and_continue_resets_state() {
2510 use crate::object_store::volume::root_volume;
2511 use crate::object_store::{ExtentValue, NewChildStoreOptions, ObjectKey, ObjectValue};
2512 use storage_device::DeviceHolder;
2513 use storage_device::fake_device::FakeDevice;
2514
2515 let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
2516 let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
2517 let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
2518 let store = root_volume
2519 .new_volume("test", NewChildStoreOptions::default())
2520 .await
2521 .expect("new_volume failed");
2522
2523 let mut transaction = store
2524 .new_transaction(
2525 lock_keys![LockKey::object(store.store_object_id(), 1000)],
2526 Options::default(),
2527 )
2528 .await
2529 .expect("new_transaction failed");
2530
2531 assert!(!transaction.includes_write());
2532 transaction.add(
2533 store.store_object_id(),
2534 Mutation::merge_object(
2535 ObjectKey::extent(1000, AttributeId::DATA, 0..4096),
2536 ObjectValue::Extent(ExtentValue::deleted_extent()),
2537 ),
2538 );
2539 transaction.add_checksum(0..4096, vec![0], true);
2540 assert!(transaction.includes_write());
2541 assert!(!transaction.checksums().is_empty());
2542
2543 transaction.commit_and_continue().await.expect("commit_and_continue failed");
2544
2545 assert!(!transaction.includes_write());
2546 assert!(transaction.checksums().is_empty());
2547
2548 transaction.commit().await.expect("commit failed");
2549 fs.close().await.expect("close failed");
2550 }
2551}