1use crate::checksum::{Checksum, Checksums};
6use crate::errors::FxfsError;
7use crate::log::*;
8use crate::lsm_tree::Query;
9use crate::lsm_tree::merge::{Merger, MergerIterator};
10use crate::lsm_tree::types::{ItemRef, LayerIterator};
11use crate::object_handle::ObjectHandle;
12use crate::object_store::extent_record::{ExtentMode, ExtentValue};
13use crate::object_store::object_manager::ObjectManager;
14use crate::object_store::object_record::{
15 AttributeKey, BytesAndNodes, ExtendedAttributeValue, ObjectAttributes, ObjectKey,
16 ObjectKeyData, ObjectValue, Timestamp,
17};
18use crate::object_store::transaction::{
19 AssocObj, AssociatedObject, LockKey, Mutation, ObjectStoreMutation, Options, ReadGuard,
20 Transaction, lock_keys,
21};
22use crate::object_store::{
23 AttributeId, Extent, FileExtent, HandleOptions, HandleOwner, ObjectStore, TrimMode, TrimResult,
24 VOLUME_DATA_KEY_ID,
25};
26use crate::range::RangeExt;
27use anyhow::{Context, Error, anyhow, bail, ensure};
28use assert_matches::assert_matches;
29use bit_vec::BitVec;
30use futures::future::try_join_all;
31use futures::stream::{FuturesOrdered, FuturesUnordered, unfold};
32use futures::{Stream, TryStreamExt, try_join};
33use fxfs_crypto::{
34 Cipher, CipherHolder, CipherSet, EncryptionKey, FindKeyResult, FxfsCipher, KeyPurpose,
35 MutPtrByteSlice,
36};
37use fxfs_trace::{TraceFutureExt, trace, trace_future_args};
38use static_assertions::const_assert;
39use std::cmp::min;
40use std::future::Future;
41use std::ops::Range;
42use std::sync::Arc;
43use std::sync::atomic::{self, AtomicBool, Ordering};
44use storage_device::buffer::{Buffer, BufferFuture, BufferRef, MutableBufferRef};
45use storage_device::{InlineCryptoOptions, ReadOptions, WriteFlags, WriteOptions};
46use storage_units::BlockSize;
47
48use fidl_fuchsia_io as fio;
49use fuchsia_async as fasync;
50
51pub const MAX_XATTR_NAME_SIZE: usize = 255;
53pub const MAX_INLINE_XATTR_SIZE: usize = 256;
56pub const MAX_XATTR_VALUE_SIZE: usize = 64000;
60
61fn apply_bitmap_zeroing(
63 block_size: BlockSize,
64 bitmap: &bit_vec::BitVec,
65 mut buffer: MutableBufferRef<'_>,
66) {
67 let mut buf = buffer.as_mut_ptr_slice();
68 debug_assert_eq!(bitmap.len() as u64 * block_size, buf.len() as u64);
69 for (i, block) in bitmap.iter().enumerate() {
70 if !block {
71 let start = (i as u64 * block_size) as usize;
72 buf.subslice_mut(start..start + block_size.get() as usize).fill(0);
73 }
74 }
75}
76
77#[derive(Debug, Clone, PartialEq)]
81pub enum MaybeChecksums {
82 None,
83 Fletcher(Vec<Checksum>),
84}
85
86impl MaybeChecksums {
87 pub fn maybe_as_ref(&self) -> Option<&[Checksum]> {
88 match self {
89 Self::None => None,
90 Self::Fletcher(sums) => Some(&sums),
91 }
92 }
93
94 pub fn split_off(&mut self, at: usize) -> Self {
95 match self {
96 Self::None => Self::None,
97 Self::Fletcher(sums) => Self::Fletcher(sums.split_off(at)),
98 }
99 }
100
101 pub fn to_mode(self) -> ExtentMode {
102 match self {
103 Self::None => ExtentMode::Raw,
104 Self::Fletcher(sums) => ExtentMode::Cow(Checksums::fletcher(sums)),
105 }
106 }
107
108 pub fn into_option(self) -> Option<Vec<Checksum>> {
109 match self {
110 Self::None => None,
111 Self::Fletcher(sums) => Some(sums),
112 }
113 }
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub enum SetExtendedAttributeMode {
121 Set,
123 Create,
125 Replace,
127}
128
129impl From<fio::SetExtendedAttributeMode> for SetExtendedAttributeMode {
130 fn from(other: fio::SetExtendedAttributeMode) -> SetExtendedAttributeMode {
131 match other {
132 fio::SetExtendedAttributeMode::Set => SetExtendedAttributeMode::Set,
133 fio::SetExtendedAttributeMode::Create => SetExtendedAttributeMode::Create,
134 fio::SetExtendedAttributeMode::Replace => SetExtendedAttributeMode::Replace,
135 }
136 }
137}
138
139enum Encryption {
140 None,
142
143 CachedKeys,
146
147 PermanentKeys,
149}
150
151#[derive(PartialEq, Debug)]
152enum OverwriteBitmaps {
153 None,
154 Some {
155 extent_bitmap: BitVec,
157 write_bitmap: BitVec,
159 bitmap_offset: usize,
162 },
163}
164
165impl OverwriteBitmaps {
166 fn new(extent_bitmap: BitVec) -> Self {
167 OverwriteBitmaps::Some {
168 write_bitmap: BitVec::from_elem(extent_bitmap.len(), false),
169 extent_bitmap,
170 bitmap_offset: 0,
171 }
172 }
173
174 fn is_none(&self) -> bool {
175 *self == OverwriteBitmaps::None
176 }
177
178 fn set_offset(&mut self, new_offset: usize) {
179 match self {
180 OverwriteBitmaps::None => (),
181 OverwriteBitmaps::Some { bitmap_offset, .. } => *bitmap_offset = new_offset,
182 }
183 }
184
185 fn get_from_extent_bitmap(&self, i: usize) -> Option<bool> {
186 match self {
187 OverwriteBitmaps::None => None,
188 OverwriteBitmaps::Some { extent_bitmap, bitmap_offset, .. } => {
189 extent_bitmap.get(*bitmap_offset + i)
190 }
191 }
192 }
193
194 fn set_in_write_bitmap(&mut self, i: usize, x: bool) {
195 match self {
196 OverwriteBitmaps::None => (),
197 OverwriteBitmaps::Some { write_bitmap, bitmap_offset, .. } => {
198 write_bitmap.set(*bitmap_offset + i, x)
199 }
200 }
201 }
202
203 fn take_bitmaps(self) -> Option<(BitVec, BitVec)> {
204 match self {
205 OverwriteBitmaps::None => None,
206 OverwriteBitmaps::Some { extent_bitmap, write_bitmap, .. } => {
207 Some((extent_bitmap, write_bitmap))
208 }
209 }
210 }
211}
212
213#[derive(PartialEq, Debug)]
217struct ChecksumRangeChunk {
218 checksum_range: Range<usize>,
219 device_range: Range<u64>,
220 is_first_write: bool,
221}
222
223impl ChecksumRangeChunk {
224 fn group_first_write_ranges(
225 bitmaps: &mut OverwriteBitmaps,
226 block_size: BlockSize,
227 write_device_range: Range<u64>,
228 ) -> Vec<ChecksumRangeChunk> {
229 let write_block_len = (write_device_range.length().unwrap() / block_size) as usize;
230 if bitmaps.is_none() {
231 vec![ChecksumRangeChunk {
236 checksum_range: 0..write_block_len,
237 device_range: write_device_range,
238 is_first_write: false,
239 }]
240 } else {
241 let mut checksum_ranges = vec![ChecksumRangeChunk {
242 checksum_range: 0..0,
243 device_range: write_device_range.start..write_device_range.start,
244 is_first_write: !bitmaps.get_from_extent_bitmap(0).unwrap(),
245 }];
246 let mut working_range = checksum_ranges.last_mut().unwrap();
247 for i in 0..write_block_len {
248 bitmaps.set_in_write_bitmap(i, true);
249
250 if working_range.is_first_write != bitmaps.get_from_extent_bitmap(i).unwrap() {
253 working_range.checksum_range.end += 1;
256 working_range.device_range.end += block_size;
257 } else {
258 let new_chunk = ChecksumRangeChunk {
260 checksum_range: working_range.checksum_range.end
261 ..working_range.checksum_range.end + 1,
262 device_range: working_range.device_range.end
263 ..working_range.device_range.end + block_size,
264 is_first_write: !working_range.is_first_write,
265 };
266 checksum_ranges.push(new_chunk);
267 working_range = checksum_ranges.last_mut().unwrap();
268 }
269 }
270 checksum_ranges
271 }
272 }
273}
274
275pub struct StoreObjectHandle<S: HandleOwner> {
288 owner: Arc<S>,
289 object_id: u64,
290 options: HandleOptions,
291 trace: AtomicBool,
292 encryption: Encryption,
293}
294
295impl<S: HandleOwner> ObjectHandle for StoreObjectHandle<S> {
296 fn set_trace(&self, v: bool) {
297 info!(store_id = self.store().store_object_id, oid = self.object_id(), trace = v; "trace");
298 self.trace.store(v, atomic::Ordering::Relaxed);
299 }
300
301 fn object_id(&self) -> u64 {
302 return self.object_id;
303 }
304
305 fn allocate_buffer(&self, size: usize) -> BufferFuture<'_> {
306 self.store().device.allocate_buffer(size)
307 }
308
309 fn block_size(&self) -> BlockSize {
310 self.store().block_size()
311 }
312}
313
314struct Watchdog {
315 _task: fasync::Task<()>,
316}
317
318impl Watchdog {
319 fn new(increment_seconds: u64, cb: impl Fn(u64) + Send + 'static) -> Self {
320 Self {
321 _task: fasync::Task::spawn(
322 async move {
323 let increment = increment_seconds.try_into().unwrap();
324 let mut fired_counter = 0;
325 let mut next_wake = fasync::MonotonicInstant::now();
326 loop {
327 next_wake += std::time::Duration::from_secs(increment).into();
328 if fasync::MonotonicInstant::now() < next_wake {
332 fasync::Timer::new(next_wake).await;
333 }
334 fired_counter += 1;
335 cb(fired_counter);
336 }
337 }
338 .trace(trace_future_args!("StoreObjectHandle::Watchdog")),
339 ),
340 }
341 }
342}
343
344impl<S: HandleOwner> StoreObjectHandle<S> {
345 pub fn new(
347 owner: Arc<S>,
348 object_id: u64,
349 permanent_keys: bool,
350 options: HandleOptions,
351 trace: bool,
352 ) -> Self {
353 let encryption = if permanent_keys {
354 Encryption::PermanentKeys
355 } else if owner.as_ref().as_ref().is_encrypted() {
356 Encryption::CachedKeys
357 } else {
358 Encryption::None
359 };
360 Self { owner, object_id, encryption, options, trace: AtomicBool::new(trace) }
361 }
362
363 pub fn owner(&self) -> &Arc<S> {
364 &self.owner
365 }
366
367 pub fn store(&self) -> &ObjectStore {
368 self.owner.as_ref().as_ref()
369 }
370
371 pub fn trace(&self) -> bool {
372 self.trace.load(atomic::Ordering::Relaxed)
373 }
374
375 pub fn is_encrypted(&self) -> bool {
376 !matches!(self.encryption, Encryption::None)
377 }
378
379 pub fn default_transaction_options<'b>(&self) -> Options<'b> {
382 Options { skip_journal_checks: self.options.skip_journal_checks, ..Default::default() }
383 }
384
385 pub async fn new_transaction_with_options<'b>(
386 &self,
387 attribute_id: AttributeId,
388 options: Options<'b>,
389 ) -> Result<Transaction<'b>, Error> {
390 Ok(self
391 .store()
392 .new_transaction(
393 lock_keys![
394 LockKey::object_attribute(
395 self.store().store_object_id(),
396 self.object_id(),
397 attribute_id,
398 ),
399 LockKey::object(self.store().store_object_id(), self.object_id()),
400 ],
401 options,
402 )
403 .await?)
404 }
405
406 pub async fn new_transaction<'b>(
407 &self,
408 attribute_id: AttributeId,
409 ) -> Result<Transaction<'b>, Error> {
410 self.new_transaction_with_options(attribute_id, self.default_transaction_options()).await
411 }
412
413 async fn txn_get_object_mutation(
416 &self,
417 transaction: &Transaction<'_>,
418 ) -> Result<ObjectStoreMutation, Error> {
419 self.store().txn_get_object_mutation(transaction, self.object_id()).await
420 }
421
422 async fn deallocate_old_extents(
424 &self,
425 transaction: &mut Transaction<'_>,
426 attribute_id: AttributeId,
427 range: Range<u64>,
428 ) -> Result<u64, Error> {
429 let block_size = self.block_size();
430 assert!(block_size.is_aligned(&range));
431 if range.start == range.end {
432 return Ok(0);
433 }
434 let tree = &self.store().tree;
435 let layer_set = tree.layer_set();
436 let key = Extent(range);
437 let lower_bound = ObjectKey::attribute(
438 self.object_id(),
439 attribute_id,
440 AttributeKey::Extent(key.search_key()),
441 );
442 let mut merger = layer_set.merger();
443 let mut iter = merger.query(Query::FullRange(&lower_bound)).await?;
444 let allocator = self.store().allocator();
445 let mut deallocated = 0;
446 let trace = self.trace();
447 while let Some(ItemRef {
448 key:
449 ObjectKey {
450 object_id,
451 data: ObjectKeyData::Attribute(attr_id, AttributeKey::Extent(extent_key)),
452 },
453 value: ObjectValue::Extent(value),
454 ..
455 }) = iter.get()
456 {
457 if *object_id != self.object_id() || *attr_id != attribute_id {
458 break;
459 }
460 if let ExtentValue::Some { device_offset, .. } = value {
461 if let Some(overlap) = key.overlap(extent_key) {
462 let range = device_offset + overlap.start - extent_key.start
463 ..device_offset + overlap.end - extent_key.start;
464 ensure!(block_size.is_aligned(&range), FxfsError::Inconsistent);
465 if trace {
466 info!(
467 store_id = self.store().store_object_id(),
468 oid = self.object_id(),
469 device_range:? = range,
470 len = range.end - range.start,
471 extent_key:?;
472 "D",
473 );
474 }
475 allocator
476 .deallocate(transaction, self.store().store_object_id(), range)
477 .await?;
478 deallocated += overlap.end - overlap.start;
479 } else {
480 break;
481 }
482 }
483 iter.advance().await?;
484 }
485 Ok(deallocated)
486 }
487
488 async fn write_aligned(
494 &self,
495 buf: BufferRef<'_>,
496 device_offset: u64,
497 crypt_ctx: Option<(u64, u8)>,
498 flags: WriteFlags,
499 ) -> Result<MaybeChecksums, Error> {
500 if self.trace() {
501 info!(
502 store_id = self.store().store_object_id(),
503 oid = self.object_id(),
504 device_range:? = (device_offset..device_offset + buf.len() as u64),
505 len = buf.len();
506 "W",
507 );
508 }
509 let store = self.store();
510 store.device_write_ops.fetch_add(1, Ordering::Relaxed);
511 let _watchdog = Watchdog::new(10, |count| {
512 warn!("Write I/O request blocked for {} seconds", count * 10);
513 });
514
515 let (opts, compute_checksums) = match crypt_ctx {
516 Some((dun, slot)) => {
517 if !store.filesystem().options().barriers_enabled {
518 return Err(anyhow!(FxfsError::InvalidArgs)
519 .context("Barriers must be enabled for inline encrypted writes."));
520 }
521 (
522 WriteOptions { inline_crypto: InlineCryptoOptions::enabled(slot, dun), flags },
523 false,
524 )
525 }
526 None => (WriteOptions { flags, ..Default::default() }, !self.options.skip_checksums),
527 };
528
529 if compute_checksums {
530 let mut checksums = Vec::new();
531 try_join!(store.device.write_with_opts(device_offset, buf, opts), async {
532 let block_size = self.block_size().get() as usize;
533 for chunk in buf.as_ptr_slice().chunks(block_size) {
534 checksums.push(crate::checksum::fletcher64_ptr(chunk, 0));
535 }
536 Ok(())
537 })?;
538 Ok(MaybeChecksums::Fletcher(checksums))
539 } else {
540 store.device.write_with_opts(device_offset, buf, opts).await?;
541 Ok(MaybeChecksums::None)
542 }
543 }
544
545 pub async fn flush_device(&self) -> Result<(), Error> {
547 self.store().device().flush().await
548 }
549
550 pub async fn update_allocated_size(
551 &self,
552 transaction: &mut Transaction<'_>,
553 allocated: u64,
554 deallocated: u64,
555 ) -> Result<(), Error> {
556 if allocated == deallocated {
557 return Ok(());
558 }
559 let mut mutation = self.txn_get_object_mutation(transaction).await?;
560 if let ObjectValue::Object {
561 attributes: ObjectAttributes { project_id, allocated_size, .. },
562 ..
563 } = &mut mutation.item.value
564 {
565 *allocated_size = allocated_size
567 .checked_add(allocated)
568 .ok_or_else(|| anyhow!(FxfsError::Inconsistent).context("Allocated size overflow"))?
569 .checked_sub(deallocated)
570 .ok_or_else(|| {
571 anyhow!(FxfsError::Inconsistent).context("Allocated size underflow")
572 })?;
573
574 if let Some(project_id) = project_id {
575 let diff = i64::try_from(allocated).unwrap() - i64::try_from(deallocated).unwrap();
578 transaction.merge_bytes_and_nodes(
579 self.store().store_object_id(),
580 ObjectKey::project_usage(self.store().root_directory_object_id(), *project_id),
581 BytesAndNodes { bytes: diff, nodes: 0 },
582 );
583 }
584 } else {
585 bail!(anyhow!(FxfsError::Inconsistent).context("Unexpected object value"));
588 }
589 transaction.add(self.store().store_object_id, Mutation::ObjectStore(mutation));
590 Ok(())
591 }
592
593 pub async fn update_attributes<'a>(
594 &self,
595 transaction: &mut Transaction<'a>,
596 node_attributes: Option<&fio::MutableNodeAttributes>,
597 change_time: Option<Timestamp>,
598 ) -> Result<(), Error> {
599 if let Some(&fio::MutableNodeAttributes { selinux_context: Some(ref context), .. }) =
600 node_attributes
601 {
602 if let fio::SelinuxContext::Data(context) = context {
603 self.set_extended_attribute_impl(
604 "security.selinux".into(),
605 context.clone(),
606 SetExtendedAttributeMode::Set,
607 transaction,
608 )
609 .await?;
610 } else {
611 return Err(anyhow!(FxfsError::InvalidArgs)
612 .context("Only set SELinux context with `data` member."));
613 }
614 }
615 self.store()
616 .update_attributes(transaction, self.object_id, node_attributes, change_time)
617 .await
618 }
619
620 pub async fn zero(
622 &self,
623 transaction: &mut Transaction<'_>,
624 attribute_id: AttributeId,
625 range: Range<u64>,
626 ) -> Result<(), Error> {
627 let deallocated =
628 self.deallocate_old_extents(transaction, attribute_id, range.clone()).await?;
629 if deallocated > 0 {
630 self.update_allocated_size(transaction, 0, deallocated).await?;
631 transaction.add(
632 self.store().store_object_id,
633 Mutation::merge_object(
634 ObjectKey::extent(self.object_id(), attribute_id, range),
635 ObjectValue::Extent(ExtentValue::deleted_extent()),
636 ),
637 );
638 }
639 Ok(())
640 }
641
642 pub async fn align_buffer(
645 &self,
646 attribute_id: AttributeId,
647 offset: u64,
648 buf: BufferRef<'_>,
649 ) -> Result<(std::ops::Range<u64>, Buffer<'_>), Error> {
650 let block_size = self.block_size();
651 let end = offset + buf.len() as u64;
652 let aligned = block_size.align_range_outwards(offset..end).ok_or(FxfsError::TooBig)?;
653
654 let mut aligned_buf =
655 self.store().device.allocate_buffer((aligned.end - aligned.start) as usize).await;
656
657 if aligned.start < offset {
659 let mut head_block = aligned_buf.subslice_mut(0..block_size.get() as usize);
660 let read =
661 self.read_aligned(attribute_id, aligned.start, head_block.reborrow()).await?;
662 let len = head_block.len();
663 head_block.subslice_mut(read..len).fill(0);
664 }
665
666 if aligned.end > end {
668 let end_block_offset = aligned.end - block_size;
669 if offset <= end_block_offset {
671 let mut tail_block =
672 aligned_buf.subslice_mut(aligned_buf.len() - block_size.get() as usize..);
673 let read = self
674 .read_aligned(attribute_id, end_block_offset, tail_block.reborrow())
675 .await?;
676 let len = tail_block.len();
677 tail_block.subslice_mut(read..len).fill(0);
678 }
679 }
680
681 aligned_buf
682 .subslice_mut((offset - aligned.start) as usize..(end - aligned.start) as usize)
683 .copy_from_buffer(buf);
684
685 Ok((aligned, aligned_buf))
686 }
687
688 pub async fn shrink(
694 &self,
695 transaction: &mut Transaction<'_>,
696 attribute_id: AttributeId,
697 size: u64,
698 ) -> Result<NeedsTrim, Error> {
699 let store = self.store();
700 let needs_trim = matches!(
701 store
702 .trim_some(transaction, self.object_id(), attribute_id, TrimMode::FromOffset(size))
703 .await?,
704 TrimResult::Incomplete
705 );
706 if needs_trim {
707 let graveyard_id = store.graveyard_directory_object_id();
710 let in_graveyard = store
712 .tree
713 .find_map(&ObjectKey::graveyard_entry(graveyard_id, self.object_id()), |item| {
714 matches!(item.value, ObjectValue::Some | ObjectValue::Trim)
715 })
716 .await?
717 .unwrap_or(false);
718 if !in_graveyard {
719 transaction.add(
720 store.store_object_id,
721 Mutation::replace_or_insert_object(
722 ObjectKey::graveyard_entry(graveyard_id, self.object_id()),
723 ObjectValue::Trim,
724 ),
725 );
726 }
727 }
728 Ok(NeedsTrim(needs_trim))
729 }
730
731 pub async fn read_and_decrypt(
733 &self,
734 attribute_id: AttributeId,
735 device_offset: u64,
736 file_offset: u64,
737 mut buffer: MutableBufferRef<'_>,
738 key_id: u64,
739 ) -> Result<(), Error> {
740 let store = self.store();
741 store.device_read_ops.fetch_add(1, Ordering::Relaxed);
742
743 let _watchdog = Watchdog::new(10, |count| {
744 warn!("Read I/O request blocked for {} seconds", count * 10);
745 });
746
747 let (_key_id, key) = self.get_key(Some(key_id)).await?;
748 if let Some(key) = key {
749 if let Some((dun, slot)) =
750 key.crypt_ctx(self.object_id, attribute_id.raw(), file_offset)
751 {
752 store
753 .device
754 .read_with_opts(
755 device_offset as u64,
756 buffer.reborrow(),
757 ReadOptions { inline_crypto: InlineCryptoOptions::enabled(slot, dun) },
758 )
759 .await?;
760 } else {
761 store.device.read(device_offset, buffer.reborrow()).await?;
762 key.decrypt(
763 self.object_id,
764 attribute_id.raw(),
765 device_offset,
766 file_offset,
767 buffer.as_mut_ptr_slice(),
768 )?;
769 }
770 } else {
771 store.device.read(device_offset, buffer.reborrow()).await?;
772 }
773
774 Ok(())
775 }
776
777 pub async fn get_key(
782 &self,
783 key_id: Option<u64>,
784 ) -> Result<(u64, Option<Arc<dyn Cipher>>), Error> {
785 let store = self.store();
786 let result = match self.encryption {
787 Encryption::None => (VOLUME_DATA_KEY_ID, None),
788 Encryption::CachedKeys => {
789 if let Some(key_id) = key_id {
790 (
791 key_id,
792 Some(
793 store
794 .key_manager
795 .get_key(
796 self.object_id,
797 store.crypt().ok_or_else(|| anyhow!("No crypt!"))?.as_ref(),
798 async || store.get_keys(self.object_id).await,
799 key_id,
800 )
801 .await?,
802 ),
803 )
804 } else {
805 let (key_id, key) = store
806 .key_manager
807 .get_fscrypt_key_if_present(
808 self.object_id,
809 store.crypt().ok_or_else(|| anyhow!("No crypt!"))?.as_ref(),
810 async || store.get_keys(self.object_id).await,
811 )
812 .await?;
813 (key_id, Some(key))
814 }
815 }
816 Encryption::PermanentKeys => {
817 (VOLUME_DATA_KEY_ID, Some(store.key_manager.get(self.object_id).await?.unwrap()))
818 }
819 };
820
821 if let Some(ref key) = result.1 {
823 if key.crypt_ctx(self.object_id, 0, 0).is_some() {
824 if !store.filesystem().options().barriers_enabled {
825 return Err(anyhow!(FxfsError::InvalidArgs)
826 .context("Barriers must be enabled for inline encrypted writes."));
827 }
828 }
829 }
830
831 Ok(result)
832 }
833
834 pub async fn extent_stream<'a, 'b>(
838 &'a self,
839 merger: &'b mut Merger<'a, ObjectKey, ObjectValue>,
840 attribute_id: AttributeId,
841 ) -> Result<impl Stream<Item = Result<FileExtent, Error>> + 'b, Error> {
842 let object_id = self.object_id();
843 let iter = merger
844 .query(Query::FullRange(&ObjectKey::attribute(
845 object_id,
846 attribute_id,
847 AttributeKey::Extent(Extent::search_key_from_offset(0)),
848 )))
849 .await?;
850 Ok(unfold(
851 (iter, object_id, attribute_id),
852 |(mut iter, object_id, attribute_id)| async move {
853 loop {
854 match iter.get() {
855 Some(ItemRef {
856 key:
857 ObjectKey {
858 object_id: id,
859 data:
860 ObjectKeyData::Attribute(
861 attr_id,
862 AttributeKey::Extent(extent_key),
863 ),
864 },
865 value: ObjectValue::Extent(extent_value),
866 ..
867 }) if *id == object_id && *attr_id == attribute_id => {
868 let logical_range = extent_key.0.clone();
869 let device_range = match extent_value {
870 ExtentValue::Some { device_offset, .. } => {
871 let len = logical_range.end - logical_range.start;
872 Some(*device_offset..*device_offset + len)
873 }
874 ExtentValue::None => None,
876 };
877
878 if let Err(e) = iter.advance().await {
881 return Some((Err(e.into()), (iter, object_id, attribute_id)));
882 }
883
884 if let Some(device_range) = device_range {
885 return Some((
886 Ok(FileExtent::new(logical_range.start, device_range).unwrap()),
887 (iter, object_id, attribute_id),
888 ));
889 } else {
890 continue;
893 }
894 }
895 _ => return None,
897 }
898 }
899 },
900 ))
901 }
902
903 async fn get_or_create_key(
907 &self,
908 transaction: &mut Transaction<'_>,
909 ) -> Result<Arc<dyn Cipher>, Error> {
910 let store = self.store();
911
912 if let Some(key) = store.key_manager.get(self.object_id).await.context("get failed")? {
914 return Ok(key);
915 }
916
917 let crypt = store.crypt().ok_or_else(|| anyhow!("No crypt!"))?;
918
919 let (mut encryption_keys, mut cipher_set) = if let Some(value) =
921 store.tree.find_value(&ObjectKey::keys(self.object_id)).await.context("find failed")?
922 {
923 if let ObjectValue::Keys(encryption_keys) = value {
924 let cipher_set = store
925 .key_manager
926 .get_keys(
927 self.object_id,
928 crypt.as_ref(),
929 &mut Some(async || Ok(encryption_keys.clone())),
930 false,
931 false,
932 )
933 .await
934 .context("get_keys failed")?;
935 match cipher_set.find_key(VOLUME_DATA_KEY_ID) {
936 FindKeyResult::NotFound => {}
937 FindKeyResult::Unavailable => return Err(FxfsError::NoKey.into()),
938 FindKeyResult::Key(key) => return Ok(key),
939 }
940 (encryption_keys, (*cipher_set).clone())
941 } else {
942 return Err(anyhow!(FxfsError::Inconsistent));
943 }
944 } else {
945 Default::default()
946 };
947
948 let (key, unwrapped_key) = crypt.create_key(self.object_id, KeyPurpose::Data).await?;
950 let cipher: Arc<dyn Cipher> = Arc::new(FxfsCipher::new(&unwrapped_key));
951
952 cipher_set.add_key(VOLUME_DATA_KEY_ID, CipherHolder::Cipher(cipher.clone()));
955 let cipher_set = Arc::new(cipher_set);
956
957 struct UnwrappedKeys {
960 object_id: u64,
961 new_keys: Arc<CipherSet>,
962 }
963
964 impl AssociatedObject for UnwrappedKeys {
965 fn will_apply_mutation(
966 &self,
967 _mutation: &Mutation,
968 object_id: u64,
969 manager: &ObjectManager,
970 ) {
971 manager.store(object_id).unwrap().key_manager.insert(
972 self.object_id,
973 self.new_keys.clone(),
974 false,
975 );
976 }
977 }
978
979 encryption_keys.insert(VOLUME_DATA_KEY_ID, EncryptionKey::Fxfs(key).into());
980
981 transaction.add_with_object(
982 store.store_object_id(),
983 Mutation::replace_or_insert_object(
984 ObjectKey::keys(self.object_id),
985 ObjectValue::keys(encryption_keys),
986 ),
987 AssocObj::Owned(Box::new(UnwrappedKeys {
988 object_id: self.object_id,
989 new_keys: cipher_set,
990 })),
991 );
992
993 Ok(cipher)
994 }
995
996 pub async fn read_aligned(
1004 &self,
1005 attribute_id: AttributeId,
1006 offset: u64,
1007 mut buf: MutableBufferRef<'_>,
1008 ) -> Result<usize, Error> {
1009 let block_size = self.block_size();
1010 ensure!(block_size.is_aligned(offset), FxfsError::InvalidArgs);
1011 ensure!(block_size.is_aligned(buf.len() as u64), FxfsError::InvalidArgs);
1012 let fs = self.store().filesystem();
1013 let guard = fs
1014 .lock_manager()
1015 .read_lock(lock_keys![LockKey::object_attribute(
1016 self.store().store_object_id(),
1017 self.object_id(),
1018 attribute_id,
1019 )])
1020 .await;
1021
1022 let key = ObjectKey::attribute(self.object_id(), attribute_id, AttributeKey::Attribute);
1023 let size = self
1024 .store()
1025 .tree()
1026 .find_map(&key, |item| match item.value {
1027 ObjectValue::Attribute { size, .. } => Ok(*size),
1028 _ => Err(anyhow!(FxfsError::Inconsistent)),
1029 })
1030 .await?
1031 .transpose()?
1032 .unwrap_or(0);
1033 if offset >= size {
1034 return Ok(0);
1035 }
1036 let length = min(buf.len() as u64, size - offset) as usize;
1037 let aligned_length =
1038 block_size.align_up(length as u64).ok_or(FxfsError::Inconsistent)? as usize;
1039 buf = buf.subslice_mut(0..aligned_length);
1040 self.read_aligned_unchecked(attribute_id, offset, buf, &guard).await?;
1041 Ok(length)
1042 }
1043
1044 pub async fn read_aligned_unchecked(
1055 &self,
1056 attribute_id: AttributeId,
1057 offset: u64,
1058 buf: MutableBufferRef<'_>,
1059 _guard: &ReadGuard<'_>,
1060 ) -> Result<(), Error> {
1061 if buf.is_empty() {
1062 return Ok(());
1063 }
1064 let end_offset = offset + buf.len() as u64;
1065 let tree = &self.store().tree;
1066 let layer_set = tree.layer_set();
1067 let mut merger = layer_set.merger();
1068 let iter = merger
1069 .query(Query::LimitedRange(&ObjectKey::extent(
1070 self.object_id(),
1071 attribute_id,
1072 offset..end_offset,
1073 )))
1074 .await?;
1075 self.read_aligned_impl(attribute_id, offset, buf, iter).await
1076 }
1077
1078 pub async fn read_attr(&self, attribute_id: AttributeId) -> Result<Option<Box<[u8]>>, Error> {
1080 let store = self.store();
1081 let tree = &store.tree;
1082 let layer_set = tree.layer_set();
1083 let mut merger = layer_set.merger();
1084 let key = ObjectKey::attribute(self.object_id(), attribute_id, AttributeKey::Attribute);
1085 let iter = merger.query(Query::FullRange(&key)).await?;
1086 match iter.get() {
1087 Some(item) if item.key == &key => match item.value {
1088 ObjectValue::Attribute { .. } => Ok(Some(self.read_attr_from_iter(iter).await?)),
1089 ObjectValue::None => Ok(None),
1091 _ => Err(FxfsError::Inconsistent.into()),
1092 },
1093 _ => Ok(None),
1094 }
1095 }
1096
1097 pub async fn read_attr_from_iter(
1100 &self,
1101 mut iter: MergerIterator<'_, '_, ObjectKey, ObjectValue>,
1102 ) -> Result<Box<[u8]>, Error> {
1103 let (mut buffer, size, attribute_id) = match iter.get() {
1104 Some(ItemRef {
1105 key:
1106 ObjectKey {
1107 object_id,
1108 data: ObjectKeyData::Attribute(attribute_id, AttributeKey::Attribute),
1109 },
1110 value: ObjectValue::Attribute { size, .. },
1111 ..
1112 }) if *object_id == self.object_id => {
1113 (
1115 self.store()
1116 .device
1117 .allocate_buffer(self.block_size().align_up(*size).unwrap() as usize)
1118 .await,
1119 *size as usize,
1120 *attribute_id,
1121 )
1122 }
1123 _ => bail!(FxfsError::InvalidArgs),
1124 };
1125 iter.advance().await?;
1126 self.read_aligned_impl(attribute_id, 0, buffer.as_mut(), iter).await?;
1127 Ok(buffer.as_ref().subslice(0..size).to_vec().into_boxed_slice())
1128 }
1129
1130 async fn read_aligned_impl(
1138 &self,
1139 attribute_id: AttributeId,
1140 mut offset: u64,
1141 mut buf: MutableBufferRef<'_>,
1142 mut iter: MergerIterator<'_, '_, ObjectKey, ObjectValue>,
1143 ) -> Result<(), Error> {
1144 if buf.is_empty() {
1145 return Ok(());
1146 }
1147 let block_size = self.block_size();
1148 debug_assert!(block_size.is_aligned(offset));
1149 debug_assert!(block_size.is_aligned(buf.len() as u64));
1150 let device_block_size = self.store().device.block_size() as u64;
1151 debug_assert!(buf.range().start as u64 % device_block_size == 0);
1152
1153 self.store().logical_read_ops.fetch_add(1, Ordering::Relaxed);
1154
1155 let mut reads = Vec::new();
1156 loop {
1157 let (mut extent_key, mut device_offset, key_id, mode) = match iter.get() {
1158 Some(ItemRef {
1159 key:
1160 ObjectKey {
1161 object_id,
1162 data:
1163 ObjectKeyData::Attribute(attr_id, AttributeKey::Extent(extent_key)),
1164 },
1165 value: ObjectValue::Extent(extent_value),
1166 ..
1167 }) if *object_id == self.object_id && *attr_id == attribute_id => {
1168 match extent_value {
1169 ExtentValue::Some { device_offset, mode, key_id } => {
1170 ensure!(extent_key.is_valid(), FxfsError::Inconsistent);
1171 ensure!(block_size.is_aligned(extent_key), FxfsError::Inconsistent);
1172 ensure!(
1173 device_offset % device_block_size == 0,
1174 FxfsError::Inconsistent
1175 );
1176 ensure!(extent_key.end > offset, FxfsError::Inconsistent);
1178 (extent_key.0.clone(), *device_offset, *key_id, mode)
1179 }
1180 ExtentValue::None => {
1181 iter.advance().await?;
1183 continue;
1184 }
1185 }
1186 }
1187 _ => break,
1188 };
1189 if extent_key.start > offset {
1190 let split = std::cmp::min((extent_key.start - offset) as usize, buf.len());
1192 let (mut to_zero, remaining) = buf.split_at_mut(split);
1193 to_zero.fill(0);
1194 buf = remaining;
1195 if buf.is_empty() {
1196 break;
1197 }
1198 offset = extent_key.start;
1199 }
1200
1201 let mut bitmap_offset = 0;
1202 if extent_key.start < offset {
1203 let diff = offset - extent_key.start;
1204 extent_key.start = offset;
1205 device_offset += diff;
1206 bitmap_offset = diff;
1207 }
1208 if (extent_key.end - extent_key.start) > buf.len() as u64 {
1209 extent_key.end = extent_key.start + buf.len() as u64;
1210 }
1211 let (mut to_read, remaining) =
1212 buf.split_at_mut((extent_key.end - extent_key.start) as usize);
1213 buf = remaining;
1214 let maybe_bitmap = match mode {
1215 ExtentMode::OverwritePartial(bitmap) => {
1216 let mut read_bitmap =
1217 bitmap.clone().split_off((bitmap_offset / block_size) as usize);
1218 read_bitmap.truncate(((to_read.len() as u64) / block_size) as usize);
1219 Some(read_bitmap)
1220 }
1221 _ => None,
1222 };
1223 reads.push(async move {
1224 self.read_and_decrypt(
1225 attribute_id,
1226 device_offset,
1227 extent_key.start,
1228 to_read.reborrow(),
1229 key_id,
1230 )
1231 .await?;
1232 if let Some(bitmap) = maybe_bitmap {
1233 apply_bitmap_zeroing(block_size, &bitmap, to_read);
1234 }
1235 Ok::<(), Error>(())
1236 });
1237 if buf.is_empty() {
1238 break;
1239 }
1240 offset = extent_key.end;
1241 iter.advance().await?;
1242 }
1243 match reads.len() {
1244 0 => {}
1245 1 => reads.pop().unwrap().await?,
1246 _ => {
1247 try_join_all(reads).await?;
1248 }
1249 }
1250 buf.fill(0);
1251 Ok(())
1252 }
1253
1254 pub async fn write_at(
1262 &self,
1263 attribute_id: AttributeId,
1264 offset: u64,
1265 buf: MutableBufferRef<'_>,
1266 key_id: Option<u64>,
1267 device_offset: u64,
1268 ) -> Result<MaybeChecksums, Error> {
1269 self.write_at_with_flags(
1270 attribute_id,
1271 offset,
1272 buf,
1273 key_id,
1274 device_offset,
1275 WriteFlags::empty(),
1276 )
1277 .await
1278 }
1279
1280 pub async fn write_at_with_flags(
1283 &self,
1284 attribute_id: AttributeId,
1285 offset: u64,
1286 buf: MutableBufferRef<'_>,
1287 key_id: Option<u64>,
1288 mut device_offset: u64,
1289 flags: WriteFlags,
1290 ) -> Result<MaybeChecksums, Error> {
1291 let mut transfer_buf;
1292 let block_size = self.block_size();
1293 let (range, mut transfer_buf_ref) =
1294 if offset % block_size == 0 && buf.len() as u64 % block_size == 0 {
1295 (offset..offset + buf.len() as u64, buf)
1296 } else {
1297 let (range, buf) = self.align_buffer(attribute_id, offset, buf.as_ref()).await?;
1298 transfer_buf = buf;
1299 device_offset -= offset - range.start;
1300 (range, transfer_buf.as_mut())
1301 };
1302
1303 let mut crypt_ctx = None;
1304 if let (_, Some(key)) = self.get_key(key_id).await? {
1305 if let Some(ctx) = key.crypt_ctx(self.object_id, attribute_id.raw(), range.start) {
1306 crypt_ctx = Some(ctx);
1307 } else {
1308 key.encrypt(
1309 self.object_id,
1310 attribute_id.raw(),
1311 device_offset,
1312 range.start,
1313 transfer_buf_ref.as_mut_ptr_slice(),
1314 )?;
1315 }
1316 }
1317 self.write_aligned(transfer_buf_ref.as_ref(), device_offset, crypt_ctx, flags).await
1318 }
1319
1320 #[cfg(feature = "migration")]
1325 pub async fn raw_multi_write(
1326 &self,
1327 transaction: &mut Transaction<'_>,
1328 attribute_id: AttributeId,
1329 key_id: Option<u64>,
1330 ranges: &[Range<u64>],
1331 buf: MutableBufferRef<'_>,
1332 ) -> Result<(), Error> {
1333 self.multi_write_internal(transaction, attribute_id, key_id, ranges, buf).await?;
1334 Ok(())
1335 }
1336
1337 async fn multi_write_internal(
1343 &self,
1344 transaction: &mut Transaction<'_>,
1345 attribute_id: AttributeId,
1346 key_id: Option<u64>,
1347 ranges: &[Range<u64>],
1348 mut buf: MutableBufferRef<'_>,
1349 ) -> Result<(u64, u64), Error> {
1350 if buf.is_empty() {
1351 return Ok((0, 0));
1352 }
1353 let block_size = self.block_size();
1354 let store = self.store();
1355 let store_id = store.store_object_id();
1356
1357 let (key_id, key) = if key_id == Some(VOLUME_DATA_KEY_ID)
1360 && matches!(self.encryption, Encryption::CachedKeys)
1361 {
1362 (
1363 VOLUME_DATA_KEY_ID,
1364 Some(
1365 self.get_or_create_key(transaction)
1366 .await
1367 .context("get_or_create_key failed")?,
1368 ),
1369 )
1370 } else {
1371 self.get_key(key_id).await?
1372 };
1373 if let Some(key) = &key {
1374 if !key.supports_inline_encryption() {
1375 let mut slice = buf.as_mut_ptr_slice();
1376 for r in ranges {
1377 let l = r.end - r.start;
1378 let (head, tail) = slice.split_at_mut(l as usize);
1379 key.encrypt(
1380 self.object_id,
1381 attribute_id.raw(),
1382 0, r.start,
1384 MutPtrByteSlice::from(head),
1385 )?;
1386 slice = tail;
1387 }
1388 }
1389 }
1390
1391 let mut allocated = 0;
1392 let allocator = store.allocator();
1393 let trace = self.trace();
1394 let mut writes = FuturesOrdered::new();
1395
1396 let mut logical_ranges = ranges.iter();
1397 let mut current_range = logical_ranges.next().unwrap().clone();
1398
1399 while !buf.is_empty() {
1400 let mut device_range = allocator
1401 .allocate(transaction, store_id, buf.len() as u64)
1402 .await
1403 .context("allocation failed")?;
1404 if trace {
1405 info!(
1406 store_id,
1407 oid = self.object_id(),
1408 device_range:?,
1409 len = device_range.end - device_range.start;
1410 "A",
1411 );
1412 }
1413 let mut device_range_len = device_range.end - device_range.start;
1414 allocated += device_range_len;
1415 while device_range_len > 0 {
1417 if current_range.end <= current_range.start {
1418 current_range = logical_ranges.next().unwrap().clone();
1419 }
1420 let (crypt_ctx, split) = if let Some(key) = &key {
1421 if key.supports_inline_encryption() {
1422 let split = std::cmp::min(
1423 current_range.end - current_range.start,
1424 device_range_len,
1425 );
1426 let crypt_ctx =
1427 key.crypt_ctx(self.object_id, attribute_id.raw(), current_range.start);
1428 current_range.start += split;
1429 (crypt_ctx, split)
1430 } else {
1431 (None, device_range_len)
1432 }
1433 } else {
1434 (None, device_range_len)
1435 };
1436
1437 let (head, tail) = buf.split_at_mut(split as usize);
1438 buf = tail;
1439
1440 writes.push_back(async move {
1441 let len = head.len() as u64;
1442 Result::<_, Error>::Ok((
1443 device_range.start,
1444 len,
1445 self.write_aligned(
1446 head.as_ref(),
1447 device_range.start,
1448 crypt_ctx,
1449 WriteFlags::empty(),
1450 )
1451 .await?,
1452 ))
1453 });
1454 device_range.start += split;
1455 device_range_len = device_range.end - device_range.start;
1456 }
1457 }
1458
1459 self.store().logical_write_ops.fetch_add(1, Ordering::Relaxed);
1460 let ((mutations, checksums), deallocated) = try_join!(
1461 async {
1462 let mut current_range = 0..0;
1463 let mut mutations = Vec::new();
1464 let mut out_checksums = Vec::new();
1465 let mut ranges = ranges.iter();
1466 while let Some((mut device_offset, mut len, mut checksums)) =
1467 writes.try_next().await?
1468 {
1469 while len > 0 {
1470 if current_range.end <= current_range.start {
1471 current_range = ranges.next().unwrap().clone();
1472 }
1473 let chunk_len = std::cmp::min(len, current_range.end - current_range.start);
1474 let tail = checksums.split_off((chunk_len / block_size) as usize);
1475 if let Some(checksums) = checksums.maybe_as_ref() {
1476 out_checksums.push((
1477 device_offset..device_offset + chunk_len,
1478 checksums.to_owned(),
1479 ));
1480 }
1481 mutations.push(Mutation::merge_object(
1482 ObjectKey::extent(
1483 self.object_id(),
1484 attribute_id,
1485 current_range.start..current_range.start + chunk_len,
1486 ),
1487 ObjectValue::Extent(ExtentValue::new(
1488 device_offset,
1489 checksums.to_mode(),
1490 key_id,
1491 )),
1492 ));
1493 checksums = tail;
1494 device_offset += chunk_len;
1495 len -= chunk_len;
1496 current_range.start += chunk_len;
1497 }
1498 }
1499 Result::<_, Error>::Ok((mutations, out_checksums))
1500 },
1501 async {
1502 let mut deallocated = 0;
1503 for r in ranges {
1504 deallocated +=
1505 self.deallocate_old_extents(transaction, attribute_id, r.clone()).await?;
1506 }
1507 Result::<_, Error>::Ok(deallocated)
1508 }
1509 )?;
1510
1511 for m in mutations {
1512 transaction.add(store_id, m);
1513 }
1514
1515 if !store.filesystem().options().barriers_enabled {
1517 for (r, c) in checksums {
1518 transaction.add_checksum(r, c, true);
1519 }
1520 }
1521 Ok((allocated, deallocated))
1522 }
1523
1524 pub async fn multi_write(
1530 &self,
1531 transaction: &mut Transaction<'_>,
1532 attribute_id: AttributeId,
1533 key_id: Option<u64>,
1534 ranges: &[Range<u64>],
1535 buf: MutableBufferRef<'_>,
1536 ) -> Result<(), Error> {
1537 let (allocated, deallocated) =
1538 self.multi_write_internal(transaction, attribute_id, key_id, ranges, buf).await?;
1539 if allocated == 0 && deallocated == 0 {
1540 return Ok(());
1541 }
1542 self.update_allocated_size(transaction, allocated, deallocated).await
1543 }
1544
1545 pub async fn multi_overwrite<'a>(
1550 &'a self,
1551 transaction: &mut Transaction<'a>,
1552 attr_id: AttributeId,
1553 ranges: &[Range<u64>],
1554 mut buf: MutableBufferRef<'_>,
1555 ) -> Result<(), Error> {
1556 if buf.is_empty() {
1557 return Ok(());
1558 }
1559 let block_size = self.block_size();
1560 let store = self.store();
1561 let tree = store.tree();
1562 let store_id = store.store_object_id();
1563
1564 let (key_id, key) = self.get_key(None).await?;
1565 if let Some(key) = &key {
1566 if !key.supports_inline_encryption() {
1567 let mut slice = buf.as_mut_ptr_slice();
1568 for r in ranges {
1569 let l = r.end - r.start;
1570 let (head, tail) = slice.split_at_mut(l as usize);
1571 key.encrypt(
1572 self.object_id,
1573 attr_id.raw(),
1574 0, r.start,
1576 MutPtrByteSlice::from(head),
1577 )?;
1578 slice = tail;
1579 }
1580 }
1581 }
1582
1583 let mut range_iter = ranges.iter();
1584 let mut target_range = range_iter.next().unwrap().clone();
1586 let mut mutations = Vec::new();
1587 let writes = FuturesUnordered::new();
1588
1589 let layer_set = tree.layer_set();
1590 let mut merger = layer_set.merger();
1591 let mut iter = merger
1592 .query(Query::FullRange(&ObjectKey::attribute(
1593 self.object_id(),
1594 attr_id,
1595 AttributeKey::Extent(Extent::search_key_from_offset(target_range.start)),
1596 )))
1597 .await?;
1598
1599 loop {
1600 match iter.get() {
1601 Some(ItemRef {
1602 key:
1603 ObjectKey {
1604 object_id,
1605 data:
1606 ObjectKeyData::Attribute(attribute_id, AttributeKey::Extent(extent)),
1607 },
1608 value: ObjectValue::Extent(extent_value),
1609 ..
1610 }) if *object_id == self.object_id() && *attribute_id == attr_id => {
1611 if extent.end <= target_range.start {
1615 iter.advance().await?;
1616 continue;
1617 }
1618 let (device_offset, mode) = match extent_value {
1619 ExtentValue::None => {
1620 return Err(anyhow!(FxfsError::Inconsistent)).with_context(|| {
1621 format!(
1622 "multi_overwrite failed: target_range ({}, {}) overlaps with \
1623 deleted extent found at ({}, {})",
1624 target_range.start, target_range.end, extent.start, extent.end,
1625 )
1626 });
1627 }
1628 ExtentValue::Some { device_offset, mode, .. } => (device_offset, mode),
1629 };
1630 if extent.start > target_range.start {
1633 return Err(anyhow!(FxfsError::Inconsistent)).with_context(|| {
1634 format!(
1635 "multi_overwrite failed: target range ({}, {}) starts before first \
1636 extent found at ({}, {})",
1637 target_range.start, target_range.end, extent.start, extent.end,
1638 )
1639 });
1640 }
1641 let mut bitmap = match mode {
1642 ExtentMode::Raw | ExtentMode::Cow(_) => {
1643 return Err(anyhow!(FxfsError::Inconsistent)).with_context(|| {
1644 format!(
1645 "multi_overwrite failed: \
1646 extent from ({}, {}) which overlaps target range ({}, {}) had the \
1647 wrong extent mode",
1648 extent.start, extent.end, target_range.start, target_range.end,
1649 )
1650 });
1651 }
1652 ExtentMode::OverwritePartial(bitmap) => {
1653 OverwriteBitmaps::new(bitmap.clone())
1654 }
1655 ExtentMode::Overwrite => OverwriteBitmaps::None,
1656 };
1657 loop {
1658 let offset_within_extent = target_range.start - extent.start;
1659 let bitmap_offset = offset_within_extent / block_size;
1660 let write_device_offset = *device_offset + offset_within_extent;
1661 let write_end = min(extent.end, target_range.end);
1662 let write_len = write_end - target_range.start;
1663 let write_device_range =
1664 write_device_offset..write_device_offset + write_len;
1665 let (current_buf, remaining_buf) = buf.split_at_mut(write_len as usize);
1666
1667 bitmap.set_offset(bitmap_offset as usize);
1668 let checksum_ranges = ChecksumRangeChunk::group_first_write_ranges(
1669 &mut bitmap,
1670 block_size,
1671 write_device_range,
1672 );
1673
1674 let crypt_ctx = if let Some(key) = &key {
1675 key.crypt_ctx(self.object_id, attr_id.raw(), target_range.start)
1676 } else {
1677 None
1678 };
1679
1680 writes.push(async move {
1681 let maybe_checksums = self
1682 .write_aligned(
1683 current_buf.as_ref(),
1684 write_device_offset,
1685 crypt_ctx,
1686 WriteFlags::empty(),
1687 )
1688 .await?;
1689 Ok::<_, Error>(match maybe_checksums {
1690 MaybeChecksums::None => Vec::new(),
1691 MaybeChecksums::Fletcher(checksums) => checksum_ranges
1692 .into_iter()
1693 .map(
1694 |ChecksumRangeChunk {
1695 checksum_range,
1696 device_range,
1697 is_first_write,
1698 }| {
1699 (
1700 device_range,
1701 checksums[checksum_range].to_vec(),
1702 is_first_write,
1703 )
1704 },
1705 )
1706 .collect(),
1707 })
1708 });
1709 buf = remaining_buf;
1710 target_range.start += write_len;
1711 if target_range.start == target_range.end {
1712 match range_iter.next() {
1713 None => break,
1714 Some(next_range) => target_range = next_range.clone(),
1715 }
1716 }
1717 if extent.end <= target_range.start {
1718 break;
1719 }
1720 }
1721 if let Some((mut bitmap, write_bitmap)) = bitmap.take_bitmaps() {
1722 if bitmap.or(&write_bitmap) {
1723 let mode = if bitmap.all() {
1724 ExtentMode::Overwrite
1725 } else {
1726 ExtentMode::OverwritePartial(bitmap)
1727 };
1728 mutations.push(Mutation::merge_object(
1729 ObjectKey::extent(self.object_id(), attr_id, extent.clone().into()),
1730 ObjectValue::Extent(ExtentValue::new(*device_offset, mode, key_id)),
1731 ))
1732 }
1733 }
1734 if target_range.start == target_range.end {
1735 break;
1736 }
1737 iter.advance().await?;
1738 }
1739 _ => bail!(anyhow!(FxfsError::Internal).context(
1743 "found a non-extent object record while there were still ranges to process"
1744 )),
1745 }
1746 }
1747
1748 let checksums = writes.try_collect::<Vec<_>>().await?;
1749 if !store.filesystem().options().barriers_enabled {
1751 for (r, c, first_write) in checksums.into_iter().flatten() {
1752 transaction.add_checksum(r, c, first_write);
1753 }
1754 }
1755
1756 for m in mutations {
1757 transaction.add(store_id, m);
1758 }
1759
1760 Ok(())
1761 }
1762
1763 #[trace]
1770 pub async fn write_new_attr_in_batches<'a>(
1771 &'a self,
1772 transaction: &mut Transaction<'a>,
1773 attribute_id: AttributeId,
1774 data: &[u8],
1775 batch_size: usize,
1776 ) -> Result<(), Error> {
1777 transaction.add(
1778 self.store().store_object_id,
1779 Mutation::replace_or_insert_object(
1780 ObjectKey::attribute(self.object_id(), attribute_id, AttributeKey::Attribute),
1781 ObjectValue::attribute(data.len() as u64, false),
1782 ),
1783 );
1784 let chunks = data.chunks(batch_size);
1785 let num_chunks = chunks.len();
1786 if num_chunks > 1 {
1787 transaction.add(
1788 self.store().store_object_id,
1789 Mutation::replace_or_insert_object(
1790 ObjectKey::graveyard_attribute_entry(
1791 self.store().graveyard_directory_object_id(),
1792 self.object_id(),
1793 attribute_id,
1794 ),
1795 ObjectValue::Some,
1796 ),
1797 );
1798 }
1799 let mut start_offset = 0;
1800 for (i, chunk) in chunks.enumerate() {
1801 let rounded_len = self.block_size().align_up(chunk.len() as u64).unwrap();
1802 let mut buffer = self.store().device.allocate_buffer(rounded_len as usize).await;
1803 let mut slice = buffer.as_mut_ptr_slice();
1804 slice.subslice_mut(0..chunk.len()).copy_from_slice(chunk);
1805 slice.subslice_mut(chunk.len()..slice.len()).fill(0);
1806 self.multi_write(
1807 transaction,
1808 attribute_id,
1809 Some(VOLUME_DATA_KEY_ID),
1810 &[start_offset..start_offset + rounded_len],
1811 buffer.as_mut(),
1812 )
1813 .await?;
1814 start_offset += rounded_len;
1815 if i < num_chunks - 1 {
1817 transaction.commit_and_continue().await?;
1818 }
1819 }
1820 Ok(())
1821 }
1822
1823 pub async fn write_attr(
1830 &self,
1831 transaction: &mut Transaction<'_>,
1832 attribute_id: AttributeId,
1833 data: &[u8],
1834 ) -> Result<NeedsTrim, Error> {
1835 let rounded_len = self.block_size().align_up(data.len() as u64).unwrap();
1836 let store = self.store();
1837 let tree = store.tree();
1838 let should_trim = tree
1839 .find_map(
1840 &ObjectKey::attribute(self.object_id(), attribute_id, AttributeKey::Attribute),
1841 |item| match item.value {
1842 ObjectValue::Attribute { size: _, has_overwrite_extents: true } => {
1843 Err(anyhow!(FxfsError::Inconsistent)
1844 .context("write_attr on an attribute with overwrite extents"))
1845 }
1846 ObjectValue::Attribute { size, .. } => Ok((data.len() as u64) < *size),
1847 _ => Err(FxfsError::Inconsistent.into()),
1848 },
1849 )
1850 .await?
1851 .transpose()?
1852 .unwrap_or(false);
1853 let mut buffer = self.store().device.allocate_buffer(rounded_len as usize).await;
1854 let mut slice = buffer.as_mut_ptr_slice();
1855 slice.subslice_mut(0..data.len()).copy_from_slice(data);
1856 slice.subslice_mut(data.len()..slice.len()).fill(0);
1857 self.multi_write(
1858 transaction,
1859 attribute_id,
1860 Some(VOLUME_DATA_KEY_ID),
1861 &[0..rounded_len],
1862 buffer.as_mut(),
1863 )
1864 .await?;
1865 transaction.add(
1866 self.store().store_object_id,
1867 Mutation::replace_or_insert_object(
1868 ObjectKey::attribute(self.object_id(), attribute_id, AttributeKey::Attribute),
1869 ObjectValue::attribute(data.len() as u64, false),
1870 ),
1871 );
1872 if should_trim {
1873 self.shrink(transaction, attribute_id, data.len() as u64).await
1874 } else {
1875 Ok(NeedsTrim(false))
1876 }
1877 }
1878
1879 pub async fn list_extended_attributes(&self) -> Result<Vec<Vec<u8>>, Error> {
1880 let layer_set = self.store().tree().layer_set();
1881 let mut merger = layer_set.merger();
1882 let mut iter = merger
1884 .query(Query::FullRange(&ObjectKey::extended_attribute(self.object_id(), Vec::new())))
1885 .await?;
1886 let mut out = Vec::new();
1887 while let Some(item) = iter.get() {
1888 if item.value != &ObjectValue::None {
1890 match item.key {
1891 ObjectKey { object_id, data: ObjectKeyData::ExtendedAttribute { name } }
1892 if *object_id == self.object_id() =>
1893 {
1894 out.push(name.clone());
1895 }
1896 _ => break,
1901 }
1902 }
1903 iter.advance().await?;
1904 }
1905 Ok(out)
1906 }
1907
1908 pub async fn get_inline_selinux_context(&self) -> Result<Option<fio::SelinuxContext>, Error> {
1912 const_assert!(fio::MAX_SELINUX_CONTEXT_ATTRIBUTE_LEN as usize <= MAX_INLINE_XATTR_SIZE);
1915 self.store()
1916 .tree()
1917 .find_map(
1918 &ObjectKey::extended_attribute(self.object_id(), fio::SELINUX_CONTEXT_NAME.into()),
1919 |item| match item.value {
1920 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::Inline(value)) => {
1921 Ok(fio::SelinuxContext::Data(value.clone()))
1922 }
1923 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::AttributeId(_)) => {
1924 Ok(fio::SelinuxContext::UseExtendedAttributes(fio::EmptyStruct {}))
1925 }
1926 _ => Err(anyhow!(FxfsError::Inconsistent).context(
1927 "get_inline_extended_attribute: Expected ExtendedAttribute value",
1928 )),
1929 },
1930 )
1931 .await?
1932 .transpose()
1933 }
1934
1935 pub async fn get_extended_attribute(&self, name: Vec<u8>) -> Result<Vec<u8>, Error> {
1936 let value = self
1937 .store()
1938 .tree()
1939 .find_value(&ObjectKey::extended_attribute(self.object_id(), name))
1940 .await?
1941 .ok_or(FxfsError::NotFound)?;
1942 match value {
1943 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::Inline(value)) => Ok(value),
1944 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::AttributeId(id)) => {
1945 ensure!(id.is_xattr(), FxfsError::Inconsistent);
1946 Ok(self.read_attr(id).await?.ok_or(FxfsError::Inconsistent)?.into_vec())
1947 }
1948 _ => {
1949 bail!(
1950 anyhow!(FxfsError::Inconsistent)
1951 .context("get_extended_attribute: Expected ExtendedAttribute value")
1952 )
1953 }
1954 }
1955 }
1956
1957 pub async fn set_extended_attribute(
1958 &self,
1959 name: Vec<u8>,
1960 value: Vec<u8>,
1961 mode: SetExtendedAttributeMode,
1962 ) -> Result<(), Error> {
1963 let store = self.store();
1964 let keys = lock_keys![LockKey::object(store.store_object_id(), self.object_id())];
1967 let mut transaction = store.new_transaction(keys, Options::default()).await?;
1968 self.set_extended_attribute_impl(name, value, mode, &mut transaction).await?;
1969 transaction.commit().await?;
1970 Ok(())
1971 }
1972
1973 async fn set_extended_attribute_impl(
1974 &self,
1975 name: Vec<u8>,
1976 value: Vec<u8>,
1977 mode: SetExtendedAttributeMode,
1978 transaction: &mut Transaction<'_>,
1979 ) -> Result<(), Error> {
1980 ensure!(name.len() <= MAX_XATTR_NAME_SIZE, FxfsError::TooBig);
1981 ensure!(value.len() <= MAX_XATTR_VALUE_SIZE, FxfsError::TooBig);
1982 let tree = self.store().tree();
1983 let object_key = ObjectKey::extended_attribute(self.object_id(), name);
1984
1985 let existing_attribute_id = {
1986 let find_result = tree
1987 .find_map(&object_key, |item| match item.value {
1988 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::Inline(..)) => Ok(None),
1989 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::AttributeId(id)) => {
1990 ensure!(id.is_xattr(), FxfsError::Inconsistent);
1991 Ok(Some(*id))
1992 }
1993 _ => Err(anyhow!(FxfsError::Inconsistent)
1994 .context("expected extended attribute value")),
1995 })
1996 .await?;
1997 let (found, existing_attribute_id) = match find_result {
1998 Some(id) => (true, id?),
1999 None => (false, None),
2000 };
2001 match mode {
2002 SetExtendedAttributeMode::Create if found => {
2003 bail!(FxfsError::AlreadyExists)
2004 }
2005 SetExtendedAttributeMode::Replace if !found => {
2006 bail!(FxfsError::NotFound)
2007 }
2008 _ => (),
2009 }
2010 existing_attribute_id
2011 };
2012
2013 if let Some(attribute_id) = existing_attribute_id {
2014 let _ = self.write_attr(transaction, attribute_id, &value).await?;
2020 } else if value.len() <= MAX_INLINE_XATTR_SIZE {
2021 transaction.add(
2022 self.store().store_object_id(),
2023 Mutation::replace_or_insert_object(
2024 object_key,
2025 ObjectValue::inline_extended_attribute(value),
2026 ),
2027 );
2028 } else {
2029 let mut attribute_id = AttributeId::XATTR_RANGE_START;
2035 let layer_set = tree.layer_set();
2036 let mut merger = layer_set.merger();
2037 let key = ObjectKey::attribute(self.object_id(), attribute_id, AttributeKey::Attribute);
2038 let mut iter = merger.query(Query::FullRange(&key)).await?;
2039 loop {
2040 match iter.get() {
2041 None => break,
2044 Some(ItemRef {
2045 key: ObjectKey { object_id, data: ObjectKeyData::Attribute(attr_id, _) },
2046 value,
2047 ..
2048 }) if *object_id == self.object_id() => {
2049 if matches!(value, ObjectValue::None) {
2050 break;
2053 }
2054 if attribute_id < *attr_id {
2055 break;
2057 } else if attribute_id == *attr_id {
2058 attribute_id = attribute_id.next();
2060 if attribute_id == AttributeId::XATTR_RANGE_END {
2061 bail!(FxfsError::NoSpace);
2062 }
2063 }
2064 }
2068 _ => break,
2072 }
2073 iter.advance().await?;
2074 }
2075
2076 let _ = self.write_attr(transaction, attribute_id, &value).await?;
2078 transaction.add(
2079 self.store().store_object_id(),
2080 Mutation::replace_or_insert_object(
2081 object_key,
2082 ObjectValue::extended_attribute(attribute_id),
2083 ),
2084 );
2085 }
2086
2087 Ok(())
2088 }
2089
2090 pub async fn remove_extended_attribute(&self, name: Vec<u8>) -> Result<(), Error> {
2091 let store = self.store();
2092 let tree = store.tree();
2093 let object_key = ObjectKey::extended_attribute(self.object_id(), name);
2094
2095 let keys = lock_keys![LockKey::object(store.store_object_id(), self.object_id())];
2100 let mut transaction = store.new_transaction(keys, Options::default()).await?;
2101
2102 let attribute_to_delete = tree
2103 .find_map(&object_key, |item| match item.value {
2104 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::AttributeId(id)) => {
2105 ensure!(id.is_xattr(), FxfsError::Inconsistent);
2106 Ok(Some(*id))
2107 }
2108 ObjectValue::ExtendedAttribute(ExtendedAttributeValue::Inline(..)) => Ok(None),
2109 _ => Err(anyhow!(FxfsError::Inconsistent)
2110 .context("remove_extended_attribute: Expected ExtendedAttribute value")),
2111 })
2112 .await?
2113 .ok_or(FxfsError::NotFound)??;
2114
2115 transaction.add(
2116 store.store_object_id(),
2117 Mutation::replace_or_insert_object(object_key, ObjectValue::None),
2118 );
2119
2120 if let Some(attribute_id) = attribute_to_delete {
2127 let trim_result = store
2128 .trim_some(
2129 &mut transaction,
2130 self.object_id(),
2131 attribute_id,
2132 TrimMode::FromOffset(0),
2133 )
2134 .await?;
2135 assert_matches!(trim_result, TrimResult::Done(_));
2138 transaction.add(
2139 store.store_object_id(),
2140 Mutation::replace_or_insert_object(
2141 ObjectKey::attribute(self.object_id, attribute_id, AttributeKey::Attribute),
2142 ObjectValue::None,
2143 ),
2144 );
2145 }
2146
2147 transaction.commit().await?;
2148 Ok(())
2149 }
2150
2151 pub fn pre_fetch_keys(&self) -> Option<impl Future<Output = ()> + use<S>> {
2154 if let Encryption::CachedKeys = self.encryption {
2155 let owner = self.owner.clone();
2156 let object_id = self.object_id;
2157 Some(async move {
2158 let store = owner.as_ref().as_ref();
2159 if let Some(crypt) = store.crypt() {
2160 let _ = store
2161 .key_manager
2162 .get_keys(
2163 object_id,
2164 crypt.as_ref(),
2165 &mut Some(async || store.get_keys(object_id).await),
2166 false,
2167 false,
2168 )
2169 .await;
2170 }
2171 })
2172 } else {
2173 None
2174 }
2175 }
2176}
2177
2178impl<S: HandleOwner> Drop for StoreObjectHandle<S> {
2179 fn drop(&mut self) {
2180 if self.is_encrypted() {
2181 let _ = self.store().key_manager.remove(self.object_id);
2182 }
2183 }
2184}
2185
2186#[must_use]
2190pub struct NeedsTrim(pub bool);
2191
2192#[cfg(test)]
2193mod tests {
2194 use super::{ChecksumRangeChunk, OverwriteBitmaps};
2195 use crate::errors::FxfsError;
2196 use crate::filesystem::{FxFilesystem, OpenFxFilesystem};
2197 use crate::object_handle::{ObjectHandle, WriteObjectHandle};
2198 use crate::object_store::data_object_handle::WRITE_ATTR_BATCH_SIZE;
2199 use crate::object_store::transaction::{Mutation, Options, lock_keys};
2200 use crate::object_store::{
2201 AttributeId, AttributeKey, DataObjectHandle, Directory, HandleOptions, LockKey, ObjectKey,
2202 ObjectStore, ObjectValue, SetExtendedAttributeMode, StoreObjectHandle,
2203 };
2204 use bit_vec::BitVec;
2205 use fuchsia_async as fasync;
2206 use futures::{TryStreamExt, join};
2207 use std::sync::Arc;
2208 use storage_device::DeviceHolder;
2209 use storage_device::fake_device::FakeDevice;
2210 use storage_units::BlockSize;
2211
2212 const TEST_DEVICE_BLOCK_SIZE: u32 = 512;
2213 const TEST_OBJECT_NAME: &str = "foo";
2214
2215 fn is_error(actual: anyhow::Error, expected: FxfsError) {
2216 assert_eq!(*actual.root_cause().downcast_ref::<FxfsError>().unwrap(), expected)
2217 }
2218
2219 async fn test_filesystem() -> OpenFxFilesystem {
2220 let device = DeviceHolder::new(FakeDevice::new(16384, TEST_DEVICE_BLOCK_SIZE));
2221 FxFilesystem::new_empty(device).await.expect("new_empty failed")
2222 }
2223
2224 async fn test_filesystem_and_empty_object() -> (OpenFxFilesystem, DataObjectHandle<ObjectStore>)
2225 {
2226 let fs = test_filesystem().await;
2227 let store = fs.root_store();
2228
2229 let mut transaction = fs
2230 .root_store()
2231 .new_transaction(
2232 lock_keys![LockKey::object(
2233 store.store_object_id(),
2234 store.root_directory_object_id()
2235 )],
2236 Options::default(),
2237 )
2238 .await
2239 .expect("new_transaction failed");
2240
2241 let object =
2242 ObjectStore::create_object(&store, &mut transaction, HandleOptions::default(), None)
2243 .await
2244 .expect("create_object failed");
2245
2246 let root_directory =
2247 Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
2248 root_directory
2249 .add_child_file(&mut transaction, TEST_OBJECT_NAME, &object)
2250 .await
2251 .expect("add_child_file failed");
2252
2253 transaction.commit().await.expect("commit failed");
2254
2255 (fs, object)
2256 }
2257
2258 #[fuchsia::test]
2259 async fn test_extent_stream() {
2260 let (fs, object) = test_filesystem_and_empty_object().await;
2261
2262 let mut buf = object.allocate_buffer(8192).await;
2264 buf.as_mut_ptr_slice().fill(0xaa);
2265 object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
2266 object.write_or_append(Some(16384), buf.as_ref()).await.expect("write failed");
2267
2268 let basic = StoreObjectHandle::new(
2269 object.owner().clone(),
2270 object.object_id(),
2271 false,
2272 HandleOptions::default(),
2273 false,
2274 );
2275
2276 let store = fs.root_store();
2277 let layer_set = store.tree.layer_set();
2278 let mut merger = layer_set.merger();
2279
2280 let stream = basic
2281 .extent_stream(&mut merger, AttributeId::DATA)
2282 .await
2283 .expect("extent_stream failed");
2284
2285 let result: Vec<_> = stream.try_collect().await.expect("stream item error");
2286
2287 assert_eq!(result.len(), 2);
2289 assert_eq!(result[0].logical_offset(), 0);
2290 assert_eq!(result[0].length(), 8192);
2291 assert_eq!(result[1].logical_offset(), 16384);
2292 assert_eq!(result[1].length(), 8192);
2293 }
2294
2295 #[fuchsia::test(threads = 3)]
2296 async fn extended_attribute_double_remove() {
2297 let (fs, object) = test_filesystem_and_empty_object().await;
2302 let basic = Arc::new(StoreObjectHandle::new(
2303 object.owner().clone(),
2304 object.object_id(),
2305 false,
2306 HandleOptions::default(),
2307 false,
2308 ));
2309 let basic_a = basic.clone();
2310 let basic_b = basic.clone();
2311
2312 basic
2313 .set_extended_attribute(
2314 b"security.selinux".to_vec(),
2315 b"bar".to_vec(),
2316 SetExtendedAttributeMode::Set,
2317 )
2318 .await
2319 .expect("failed to set attribute");
2320
2321 let a_task = fasync::Task::spawn(async move {
2324 basic_a.remove_extended_attribute(b"security.selinux".to_vec()).await
2325 });
2326 let b_task = fasync::Task::spawn(async move {
2327 basic_b.remove_extended_attribute(b"security.selinux".to_vec()).await
2328 });
2329 match join!(a_task, b_task) {
2330 (Ok(()), Ok(())) => panic!("both remove calls succeeded"),
2331 (Err(_), Err(_)) => panic!("both remove calls failed"),
2332
2333 (Ok(()), Err(e)) => is_error(e, FxfsError::NotFound),
2334 (Err(e), Ok(())) => is_error(e, FxfsError::NotFound),
2335 }
2336
2337 fs.close().await.expect("Close failed");
2338 }
2339
2340 #[fuchsia::test(threads = 3)]
2341 async fn extended_attribute_double_create() {
2342 let (fs, object) = test_filesystem_and_empty_object().await;
2347 let basic = Arc::new(StoreObjectHandle::new(
2348 object.owner().clone(),
2349 object.object_id(),
2350 false,
2351 HandleOptions::default(),
2352 false,
2353 ));
2354 let basic_a = basic.clone();
2355 let basic_b = basic.clone();
2356
2357 let a_task = fasync::Task::spawn(async move {
2360 basic_a
2361 .set_extended_attribute(
2362 b"security.selinux".to_vec(),
2363 b"one".to_vec(),
2364 SetExtendedAttributeMode::Create,
2365 )
2366 .await
2367 });
2368 let b_task = fasync::Task::spawn(async move {
2369 basic_b
2370 .set_extended_attribute(
2371 b"security.selinux".to_vec(),
2372 b"two".to_vec(),
2373 SetExtendedAttributeMode::Create,
2374 )
2375 .await
2376 });
2377 match join!(a_task, b_task) {
2378 (Ok(()), Ok(())) => panic!("both set calls succeeded"),
2379 (Err(_), Err(_)) => panic!("both set calls failed"),
2380
2381 (Ok(()), Err(e)) => {
2382 assert_eq!(
2383 basic
2384 .get_extended_attribute(b"security.selinux".to_vec())
2385 .await
2386 .expect("failed to get xattr"),
2387 b"one"
2388 );
2389 is_error(e, FxfsError::AlreadyExists);
2390 }
2391 (Err(e), Ok(())) => {
2392 assert_eq!(
2393 basic
2394 .get_extended_attribute(b"security.selinux".to_vec())
2395 .await
2396 .expect("failed to get xattr"),
2397 b"two"
2398 );
2399 is_error(e, FxfsError::AlreadyExists);
2400 }
2401 }
2402
2403 fs.close().await.expect("Close failed");
2404 }
2405
2406 struct TestAttr {
2407 name: Vec<u8>,
2408 value: Vec<u8>,
2409 }
2410
2411 impl TestAttr {
2412 fn new(name: impl AsRef<[u8]>, value: impl AsRef<[u8]>) -> Self {
2413 Self { name: name.as_ref().to_vec(), value: value.as_ref().to_vec() }
2414 }
2415 fn name(&self) -> Vec<u8> {
2416 self.name.clone()
2417 }
2418 fn value(&self) -> Vec<u8> {
2419 self.value.clone()
2420 }
2421 }
2422
2423 #[fuchsia::test]
2424 async fn extended_attributes() {
2425 let (fs, object) = test_filesystem_and_empty_object().await;
2426
2427 let test_attr = TestAttr::new(b"security.selinux", b"foo");
2428
2429 assert_eq!(object.list_extended_attributes().await.unwrap(), Vec::<Vec<u8>>::new());
2430 is_error(
2431 object.get_extended_attribute(test_attr.name()).await.unwrap_err(),
2432 FxfsError::NotFound,
2433 );
2434
2435 object
2436 .set_extended_attribute(
2437 test_attr.name(),
2438 test_attr.value(),
2439 SetExtendedAttributeMode::Set,
2440 )
2441 .await
2442 .unwrap();
2443 assert_eq!(object.list_extended_attributes().await.unwrap(), vec![test_attr.name()]);
2444 assert_eq!(
2445 object.get_extended_attribute(test_attr.name()).await.unwrap(),
2446 test_attr.value()
2447 );
2448
2449 object.remove_extended_attribute(test_attr.name()).await.unwrap();
2450 assert_eq!(object.list_extended_attributes().await.unwrap(), Vec::<Vec<u8>>::new());
2451 is_error(
2452 object.get_extended_attribute(test_attr.name()).await.unwrap_err(),
2453 FxfsError::NotFound,
2454 );
2455
2456 object
2458 .set_extended_attribute(
2459 test_attr.name(),
2460 test_attr.value(),
2461 SetExtendedAttributeMode::Set,
2462 )
2463 .await
2464 .unwrap();
2465 assert_eq!(object.list_extended_attributes().await.unwrap(), vec![test_attr.name()]);
2466 assert_eq!(
2467 object.get_extended_attribute(test_attr.name()).await.unwrap(),
2468 test_attr.value()
2469 );
2470
2471 object.remove_extended_attribute(test_attr.name()).await.unwrap();
2472 assert_eq!(object.list_extended_attributes().await.unwrap(), Vec::<Vec<u8>>::new());
2473 is_error(
2474 object.get_extended_attribute(test_attr.name()).await.unwrap_err(),
2475 FxfsError::NotFound,
2476 );
2477
2478 fs.close().await.expect("close failed");
2479 }
2480
2481 #[fuchsia::test]
2482 async fn extended_attribute_invalid_id_range() {
2483 let (fs, object) = test_filesystem_and_empty_object().await;
2484
2485 let store = object.owner();
2486 let mut transaction = store
2487 .new_transaction(
2488 lock_keys![LockKey::object(store.store_object_id(), object.object_id())],
2489 Options::default(),
2490 )
2491 .await
2492 .unwrap();
2493 transaction.add(
2494 store.store_object_id(),
2495 Mutation::replace_or_insert_object(
2496 ObjectKey::extended_attribute(object.object_id(), b"bad_attr".to_vec()),
2497 ObjectValue::extended_attribute(AttributeId::DATA),
2498 ),
2499 );
2500 transaction.commit().await.unwrap();
2501
2502 is_error(
2503 object.get_extended_attribute(b"bad_attr".to_vec()).await.unwrap_err(),
2504 FxfsError::Inconsistent,
2505 );
2506 is_error(
2507 object
2508 .set_extended_attribute(
2509 b"bad_attr".to_vec(),
2510 b"new_val".to_vec(),
2511 SetExtendedAttributeMode::Set,
2512 )
2513 .await
2514 .unwrap_err(),
2515 FxfsError::Inconsistent,
2516 );
2517 is_error(
2518 object.remove_extended_attribute(b"bad_attr".to_vec()).await.unwrap_err(),
2519 FxfsError::Inconsistent,
2520 );
2521
2522 fs.close().await.expect("close failed");
2523 }
2524
2525 #[fuchsia::test]
2526 async fn large_extended_attribute() {
2527 let (fs, object) = test_filesystem_and_empty_object().await;
2528
2529 let test_attr = TestAttr::new(b"security.selinux", vec![3u8; 300]);
2530
2531 object
2532 .set_extended_attribute(
2533 test_attr.name(),
2534 test_attr.value(),
2535 SetExtendedAttributeMode::Set,
2536 )
2537 .await
2538 .unwrap();
2539 assert_eq!(
2540 object.get_extended_attribute(test_attr.name()).await.unwrap(),
2541 test_attr.value()
2542 );
2543
2544 assert_eq!(
2547 object
2548 .read_attr(AttributeId::XATTR_RANGE_START)
2549 .await
2550 .expect("read_attr failed")
2551 .expect("read_attr returned none")
2552 .into_vec(),
2553 test_attr.value()
2554 );
2555
2556 object.remove_extended_attribute(test_attr.name()).await.unwrap();
2557 is_error(
2558 object.get_extended_attribute(test_attr.name()).await.unwrap_err(),
2559 FxfsError::NotFound,
2560 );
2561
2562 object
2564 .set_extended_attribute(
2565 test_attr.name(),
2566 test_attr.value(),
2567 SetExtendedAttributeMode::Set,
2568 )
2569 .await
2570 .unwrap();
2571 assert_eq!(
2572 object.get_extended_attribute(test_attr.name()).await.unwrap(),
2573 test_attr.value()
2574 );
2575 object.remove_extended_attribute(test_attr.name()).await.unwrap();
2576 is_error(
2577 object.get_extended_attribute(test_attr.name()).await.unwrap_err(),
2578 FxfsError::NotFound,
2579 );
2580
2581 fs.close().await.expect("close failed");
2582 }
2583
2584 #[fuchsia::test]
2585 async fn multiple_extended_attributes() {
2586 let (fs, object) = test_filesystem_and_empty_object().await;
2587
2588 let attrs = [
2589 TestAttr::new(b"security.selinux", b"foo"),
2590 TestAttr::new(b"large.attribute", vec![3u8; 300]),
2591 TestAttr::new(b"an.attribute", b"asdf"),
2592 TestAttr::new(b"user.big", vec![5u8; 288]),
2593 TestAttr::new(b"user.tiny", b"smol"),
2594 TestAttr::new(b"this string doesn't matter", b"the quick brown fox etc"),
2595 TestAttr::new(b"also big", vec![7u8; 500]),
2596 TestAttr::new(b"all.ones", vec![1u8; 11111]),
2597 ];
2598
2599 for i in 0..attrs.len() {
2600 object
2601 .set_extended_attribute(
2602 attrs[i].name(),
2603 attrs[i].value(),
2604 SetExtendedAttributeMode::Set,
2605 )
2606 .await
2607 .unwrap();
2608 assert_eq!(
2609 object.get_extended_attribute(attrs[i].name()).await.unwrap(),
2610 attrs[i].value()
2611 );
2612 }
2613
2614 for i in 0..attrs.len() {
2615 let mut found_attrs = object.list_extended_attributes().await.unwrap();
2617 let mut expected_attrs: Vec<Vec<u8>> = attrs.iter().skip(i).map(|a| a.name()).collect();
2618 found_attrs.sort();
2619 expected_attrs.sort();
2620 assert_eq!(found_attrs, expected_attrs);
2621 for j in i..attrs.len() {
2622 assert_eq!(
2623 object.get_extended_attribute(attrs[j].name()).await.unwrap(),
2624 attrs[j].value()
2625 );
2626 }
2627
2628 object.remove_extended_attribute(attrs[i].name()).await.expect("failed to remove");
2629 is_error(
2630 object.get_extended_attribute(attrs[i].name()).await.unwrap_err(),
2631 FxfsError::NotFound,
2632 );
2633 }
2634
2635 fs.close().await.expect("close failed");
2636 }
2637
2638 #[fuchsia::test]
2639 async fn multiple_extended_attributes_delete() {
2640 let (fs, object) = test_filesystem_and_empty_object().await;
2641 let store = object.owner().clone();
2642
2643 let attrs = [
2644 TestAttr::new(b"security.selinux", b"foo"),
2645 TestAttr::new(b"large.attribute", vec![3u8; 300]),
2646 TestAttr::new(b"an.attribute", b"asdf"),
2647 TestAttr::new(b"user.big", vec![5u8; 288]),
2648 TestAttr::new(b"user.tiny", b"smol"),
2649 TestAttr::new(b"this string doesn't matter", b"the quick brown fox etc"),
2650 TestAttr::new(b"also big", vec![7u8; 500]),
2651 TestAttr::new(b"all.ones", vec![1u8; 11111]),
2652 ];
2653
2654 for i in 0..attrs.len() {
2655 object
2656 .set_extended_attribute(
2657 attrs[i].name(),
2658 attrs[i].value(),
2659 SetExtendedAttributeMode::Set,
2660 )
2661 .await
2662 .unwrap();
2663 assert_eq!(
2664 object.get_extended_attribute(attrs[i].name()).await.unwrap(),
2665 attrs[i].value()
2666 );
2667 }
2668
2669 let root_directory =
2671 Directory::open(object.owner(), object.store().root_directory_object_id())
2672 .await
2673 .expect("open failed");
2674 let mut transaction = fs
2675 .root_store()
2676 .new_transaction(
2677 lock_keys![
2678 LockKey::object(store.store_object_id(), store.root_directory_object_id()),
2679 LockKey::object(store.store_object_id(), object.object_id()),
2680 ],
2681 Options::default(),
2682 )
2683 .await
2684 .expect("new_transaction failed");
2685 crate::object_store::directory::replace_child(
2686 &mut transaction,
2687 None,
2688 (&root_directory, TEST_OBJECT_NAME),
2689 )
2690 .await
2691 .expect("replace_child failed");
2692 transaction.commit().await.unwrap();
2693 store.tombstone_object(object.object_id(), Options::default(), None).await.unwrap();
2694
2695 crate::fsck::fsck(fs.clone()).await.unwrap();
2696
2697 fs.close().await.expect("close failed");
2698 }
2699
2700 #[fuchsia::test]
2701 async fn extended_attribute_changing_sizes() {
2702 let (fs, object) = test_filesystem_and_empty_object().await;
2703
2704 let test_name = b"security.selinux";
2705 let test_small_attr = TestAttr::new(test_name, b"smol");
2706 let test_large_attr = TestAttr::new(test_name, vec![3u8; 300]);
2707
2708 object
2709 .set_extended_attribute(
2710 test_small_attr.name(),
2711 test_small_attr.value(),
2712 SetExtendedAttributeMode::Set,
2713 )
2714 .await
2715 .unwrap();
2716 assert_eq!(
2717 object.get_extended_attribute(test_small_attr.name()).await.unwrap(),
2718 test_small_attr.value()
2719 );
2720
2721 assert!(
2723 object
2724 .read_attr(AttributeId::XATTR_RANGE_START)
2725 .await
2726 .expect("read_attr failed")
2727 .is_none()
2728 );
2729
2730 crate::fsck::fsck(fs.clone()).await.unwrap();
2731
2732 object
2733 .set_extended_attribute(
2734 test_large_attr.name(),
2735 test_large_attr.value(),
2736 SetExtendedAttributeMode::Set,
2737 )
2738 .await
2739 .unwrap();
2740 assert_eq!(
2741 object.get_extended_attribute(test_large_attr.name()).await.unwrap(),
2742 test_large_attr.value()
2743 );
2744
2745 assert_eq!(
2748 object
2749 .read_attr(AttributeId::XATTR_RANGE_START)
2750 .await
2751 .expect("read_attr failed")
2752 .expect("read_attr returned none")
2753 .into_vec(),
2754 test_large_attr.value()
2755 );
2756
2757 crate::fsck::fsck(fs.clone()).await.unwrap();
2758
2759 object
2760 .set_extended_attribute(
2761 test_small_attr.name(),
2762 test_small_attr.value(),
2763 SetExtendedAttributeMode::Set,
2764 )
2765 .await
2766 .unwrap();
2767 assert_eq!(
2768 object.get_extended_attribute(test_small_attr.name()).await.unwrap(),
2769 test_small_attr.value()
2770 );
2771
2772 assert_eq!(
2775 object
2776 .read_attr(AttributeId::XATTR_RANGE_START)
2777 .await
2778 .expect("read_attr failed")
2779 .expect("read_attr returned none")
2780 .into_vec(),
2781 test_small_attr.value()
2782 );
2783
2784 crate::fsck::fsck(fs.clone()).await.unwrap();
2785
2786 object.remove_extended_attribute(test_small_attr.name()).await.expect("failed to remove");
2787
2788 crate::fsck::fsck(fs.clone()).await.unwrap();
2789
2790 fs.close().await.expect("close failed");
2791 }
2792
2793 #[fuchsia::test]
2794 async fn extended_attribute_max_size() {
2795 let (fs, object) = test_filesystem_and_empty_object().await;
2796
2797 let test_attr = TestAttr::new(
2798 vec![3u8; super::MAX_XATTR_NAME_SIZE],
2799 vec![1u8; super::MAX_XATTR_VALUE_SIZE],
2800 );
2801
2802 object
2803 .set_extended_attribute(
2804 test_attr.name(),
2805 test_attr.value(),
2806 SetExtendedAttributeMode::Set,
2807 )
2808 .await
2809 .unwrap();
2810 assert_eq!(
2811 object.get_extended_attribute(test_attr.name()).await.unwrap(),
2812 test_attr.value()
2813 );
2814 assert_eq!(object.list_extended_attributes().await.unwrap(), vec![test_attr.name()]);
2815 object.remove_extended_attribute(test_attr.name()).await.unwrap();
2816
2817 fs.close().await.expect("close failed");
2818 }
2819
2820 #[fuchsia::test]
2821 async fn extended_attribute_remove_then_create() {
2822 let (fs, object) = test_filesystem_and_empty_object().await;
2823
2824 let test_attr = TestAttr::new(
2825 vec![3u8; super::MAX_XATTR_NAME_SIZE],
2826 vec![1u8; super::MAX_XATTR_VALUE_SIZE],
2827 );
2828
2829 object
2830 .set_extended_attribute(
2831 test_attr.name(),
2832 test_attr.value(),
2833 SetExtendedAttributeMode::Create,
2834 )
2835 .await
2836 .unwrap();
2837 fs.journal().force_compact().await.unwrap();
2838 object.remove_extended_attribute(test_attr.name()).await.unwrap();
2839 object
2840 .set_extended_attribute(
2841 test_attr.name(),
2842 test_attr.value(),
2843 SetExtendedAttributeMode::Create,
2844 )
2845 .await
2846 .unwrap();
2847
2848 assert_eq!(
2849 object.get_extended_attribute(test_attr.name()).await.unwrap(),
2850 test_attr.value()
2851 );
2852
2853 fs.close().await.expect("close failed");
2854 }
2855
2856 #[fuchsia::test]
2857 async fn large_extended_attribute_max_number() {
2858 let (fs, object) = test_filesystem_and_empty_object().await;
2859
2860 let max_xattrs = AttributeId::XATTR_RANGE_END.raw() - AttributeId::XATTR_RANGE_START.raw();
2861 for i in 0..max_xattrs {
2862 let test_attr = TestAttr::new(format!("{}", i).as_bytes(), vec![0x3; 300]);
2863 object
2864 .set_extended_attribute(
2865 test_attr.name(),
2866 test_attr.value(),
2867 SetExtendedAttributeMode::Set,
2868 )
2869 .await
2870 .unwrap_or_else(|_| panic!("failed to set xattr number {}", i));
2871 }
2872
2873 match object
2876 .set_extended_attribute(
2877 b"one.too.many".to_vec(),
2878 vec![0x3; 300],
2879 SetExtendedAttributeMode::Set,
2880 )
2881 .await
2882 {
2883 Ok(()) => panic!("set should not succeed"),
2884 Err(e) => is_error(e, FxfsError::NoSpace),
2885 }
2886
2887 object
2889 .set_extended_attribute(
2890 b"this.is.okay".to_vec(),
2891 b"small value".to_vec(),
2892 SetExtendedAttributeMode::Set,
2893 )
2894 .await
2895 .unwrap();
2896
2897 object
2899 .set_extended_attribute(b"11".to_vec(), vec![0x4; 300], SetExtendedAttributeMode::Set)
2900 .await
2901 .unwrap();
2902 object
2903 .set_extended_attribute(
2904 b"12".to_vec(),
2905 vec![0x1; 300],
2906 SetExtendedAttributeMode::Replace,
2907 )
2908 .await
2909 .unwrap();
2910
2911 object.remove_extended_attribute(b"5".to_vec()).await.unwrap();
2913 object
2914 .set_extended_attribute(
2915 b"new attr".to_vec(),
2916 vec![0x3; 300],
2917 SetExtendedAttributeMode::Set,
2918 )
2919 .await
2920 .unwrap();
2921
2922 fs.close().await.expect("close failed");
2923 }
2924
2925 #[fuchsia::test]
2926 async fn write_attr_trims_beyond_new_end() {
2927 let (fs, object) = test_filesystem_and_empty_object().await;
2932
2933 let block_size = fs.block_size();
2934 let buf_size = block_size * 2;
2935 let attribute_id = AttributeId::TEST_ID;
2936
2937 let mut transaction = (*object).new_transaction(attribute_id).await.unwrap();
2938 let mut buffer = object.allocate_buffer(buf_size as usize).await;
2939 buffer.fill(3);
2940 object
2943 .multi_write(
2944 &mut transaction,
2945 attribute_id,
2946 &[0..block_size.get(), block_size.get()..block_size * 2],
2947 buffer.as_mut(),
2948 )
2949 .await
2950 .unwrap();
2951 transaction.add(
2952 object.store().store_object_id,
2953 Mutation::replace_or_insert_object(
2954 ObjectKey::attribute(object.object_id(), attribute_id, AttributeKey::Attribute),
2955 ObjectValue::attribute(block_size * 2, false),
2956 ),
2957 );
2958 transaction.commit().await.unwrap();
2959
2960 crate::fsck::fsck(fs.clone()).await.unwrap();
2961
2962 let mut transaction = (*object).new_transaction(attribute_id).await.unwrap();
2963 let needs_trim = (*object)
2964 .write_attr(&mut transaction, attribute_id, &vec![3u8; block_size.get() as usize])
2965 .await
2966 .unwrap();
2967 assert!(!needs_trim.0);
2968 transaction.commit().await.unwrap();
2969
2970 crate::fsck::fsck(fs.clone()).await.unwrap();
2971
2972 fs.close().await.expect("close failed");
2973 }
2974
2975 #[fuchsia::test]
2976 async fn write_new_attr_in_batches_multiple_txns() {
2977 let (fs, object) = test_filesystem_and_empty_object().await;
2978 let merkle_tree = vec![1; 3 * WRITE_ATTR_BATCH_SIZE];
2979 let mut transaction = (*object).new_transaction(AttributeId::TEST_ID).await.unwrap();
2980 object
2981 .write_new_attr_in_batches(
2982 &mut transaction,
2983 AttributeId::TEST_ID,
2984 &merkle_tree,
2985 WRITE_ATTR_BATCH_SIZE,
2986 )
2987 .await
2988 .expect("failed to write merkle attribute");
2989
2990 transaction.add(
2991 object.store().store_object_id,
2992 Mutation::replace_or_insert_object(
2993 ObjectKey::graveyard_attribute_entry(
2994 object.store().graveyard_directory_object_id(),
2995 object.object_id(),
2996 AttributeId::TEST_ID,
2997 ),
2998 ObjectValue::None,
2999 ),
3000 );
3001 transaction.commit().await.unwrap();
3002 assert_eq!(
3003 object.read_attr(AttributeId::TEST_ID).await.expect("read_attr failed"),
3004 Some(merkle_tree.into())
3005 );
3006
3007 fs.close().await.expect("close failed");
3008 }
3009
3010 #[cfg(target_os = "fuchsia")]
3012 #[fuchsia::test(allow_stalls = false)]
3013 async fn test_watchdog() {
3014 use super::Watchdog;
3015 use fuchsia_async::{MonotonicDuration, MonotonicInstant, TestExecutor};
3016 use std::sync::mpsc::channel;
3017
3018 TestExecutor::advance_to(make_time(0)).await;
3019 let (sender, receiver) = channel();
3020
3021 fn make_time(time_secs: i64) -> MonotonicInstant {
3022 MonotonicInstant::from_nanos(0) + MonotonicDuration::from_seconds(time_secs)
3023 }
3024
3025 {
3026 let _watchdog = Watchdog::new(10, move |count| {
3027 sender.send(count).expect("Sending value");
3028 });
3029
3030 TestExecutor::advance_to(make_time(5)).await;
3032 receiver.try_recv().expect_err("Should not have message");
3033
3034 TestExecutor::advance_to(make_time(10)).await;
3036 assert_eq!(1, receiver.recv().expect("Receiving"));
3037
3038 TestExecutor::advance_to(make_time(15)).await;
3040 receiver.try_recv().expect_err("Should not have message");
3041
3042 TestExecutor::advance_to(make_time(30)).await;
3044 assert_eq!(2, receiver.recv().expect("Receiving"));
3045 assert_eq!(3, receiver.recv().expect("Receiving"));
3046 }
3047
3048 TestExecutor::advance_to(make_time(100)).await;
3050 receiver.recv().expect_err("Watchdog should be gone");
3051 }
3052
3053 #[fuchsia::test]
3054 fn test_checksum_range_chunk() {
3055 let block_size = BlockSize::SIZE_4KIB;
3056
3057 assert_eq!(
3059 ChecksumRangeChunk::group_first_write_ranges(
3060 &mut OverwriteBitmaps::None,
3061 block_size,
3062 block_size * 2..block_size * 5,
3063 ),
3064 vec![ChecksumRangeChunk {
3065 checksum_range: 0..3,
3066 device_range: block_size * 2..block_size * 5,
3067 is_first_write: false,
3068 }],
3069 );
3070
3071 let mut bitmaps = OverwriteBitmaps::new(BitVec::from_bytes(&[0b11110000]));
3072 assert_eq!(
3073 ChecksumRangeChunk::group_first_write_ranges(
3074 &mut bitmaps,
3075 block_size,
3076 block_size * 2..block_size * 5,
3077 ),
3078 vec![ChecksumRangeChunk {
3079 checksum_range: 0..3,
3080 device_range: block_size * 2..block_size * 5,
3081 is_first_write: false,
3082 }],
3083 );
3084 assert_eq!(
3085 bitmaps.take_bitmaps(),
3086 Some((BitVec::from_bytes(&[0b11110000]), BitVec::from_bytes(&[0b11100000])))
3087 );
3088
3089 let mut bitmaps = OverwriteBitmaps::new(BitVec::from_bytes(&[0b11110000]));
3090 bitmaps.set_offset(2);
3091 assert_eq!(
3092 ChecksumRangeChunk::group_first_write_ranges(
3093 &mut bitmaps,
3094 block_size,
3095 block_size * 2..block_size * 5,
3096 ),
3097 vec![
3098 ChecksumRangeChunk {
3099 checksum_range: 0..2,
3100 device_range: block_size * 2..block_size * 4,
3101 is_first_write: false,
3102 },
3103 ChecksumRangeChunk {
3104 checksum_range: 2..3,
3105 device_range: block_size * 4..block_size * 5,
3106 is_first_write: true,
3107 },
3108 ],
3109 );
3110 assert_eq!(
3111 bitmaps.take_bitmaps(),
3112 Some((BitVec::from_bytes(&[0b11110000]), BitVec::from_bytes(&[0b00111000])))
3113 );
3114
3115 let mut bitmaps = OverwriteBitmaps::new(BitVec::from_bytes(&[0b11110000]));
3116 bitmaps.set_offset(4);
3117 assert_eq!(
3118 ChecksumRangeChunk::group_first_write_ranges(
3119 &mut bitmaps,
3120 block_size,
3121 block_size * 2..block_size * 5,
3122 ),
3123 vec![ChecksumRangeChunk {
3124 checksum_range: 0..3,
3125 device_range: block_size * 2..block_size * 5,
3126 is_first_write: true,
3127 }],
3128 );
3129 assert_eq!(
3130 bitmaps.take_bitmaps(),
3131 Some((BitVec::from_bytes(&[0b11110000]), BitVec::from_bytes(&[0b00001110])))
3132 );
3133
3134 let mut bitmaps = OverwriteBitmaps::new(BitVec::from_bytes(&[0b01010101]));
3135 assert_eq!(
3136 ChecksumRangeChunk::group_first_write_ranges(
3137 &mut bitmaps,
3138 block_size,
3139 block_size * 2..block_size * 10,
3140 ),
3141 vec![
3142 ChecksumRangeChunk {
3143 checksum_range: 0..1,
3144 device_range: block_size * 2..block_size * 3,
3145 is_first_write: true,
3146 },
3147 ChecksumRangeChunk {
3148 checksum_range: 1..2,
3149 device_range: block_size * 3..block_size * 4,
3150 is_first_write: false,
3151 },
3152 ChecksumRangeChunk {
3153 checksum_range: 2..3,
3154 device_range: block_size * 4..block_size * 5,
3155 is_first_write: true,
3156 },
3157 ChecksumRangeChunk {
3158 checksum_range: 3..4,
3159 device_range: block_size * 5..block_size * 6,
3160 is_first_write: false,
3161 },
3162 ChecksumRangeChunk {
3163 checksum_range: 4..5,
3164 device_range: block_size * 6..block_size * 7,
3165 is_first_write: true,
3166 },
3167 ChecksumRangeChunk {
3168 checksum_range: 5..6,
3169 device_range: block_size * 7..block_size * 8,
3170 is_first_write: false,
3171 },
3172 ChecksumRangeChunk {
3173 checksum_range: 6..7,
3174 device_range: block_size * 8..block_size * 9,
3175 is_first_write: true,
3176 },
3177 ChecksumRangeChunk {
3178 checksum_range: 7..8,
3179 device_range: block_size * 9..block_size * 10,
3180 is_first_write: false,
3181 },
3182 ],
3183 );
3184 assert_eq!(
3185 bitmaps.take_bitmaps(),
3186 Some((BitVec::from_bytes(&[0b01010101]), BitVec::from_bytes(&[0b11111111])))
3187 );
3188 }
3189
3190 #[fuchsia::test]
3191 async fn test_read_large_buffer_excessive_partitions() {
3192 let (_fs, object) = test_filesystem_and_empty_object().await;
3193
3194 let size = 5 * 1024 * 1024;
3196 let mut buf = object.allocate_buffer(size).await;
3197 buf.as_mut_ptr_slice().fill(0xab);
3198 object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
3199
3200 object.owner().flush().await.expect("flush failed");
3202
3203 let mut read_buf = object.allocate_buffer(size).await;
3205 assert_eq!(object.read_aligned(0, read_buf.as_mut()).await.expect("read failed"), size);
3206 assert_eq!(&read_buf.as_ptr_slice().to_vec()[..], &buf.as_ptr_slice().to_vec()[..]);
3207 }
3208}