Skip to main content

fxfs/object_store/
store_object_handle.rs

1// Copyright 2023 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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
51/// Maximum size for an extended attribute name.
52pub const MAX_XATTR_NAME_SIZE: usize = 255;
53/// Maximum size an extended attribute can be before it's stored in an object attribute instead of
54/// inside the record directly.
55pub const MAX_INLINE_XATTR_SIZE: usize = 256;
56/// Maximum size for an extended attribute value. NB: the maximum size for an extended attribute is
57/// 64kB, which we rely on for correctness when deleting attributes, so ensure it's always
58/// enforced.
59pub const MAX_XATTR_VALUE_SIZE: usize = 64000;
60
61/// Zeroes blocks in 'buffer' based on `bitmap`, one bit per block from start of buffer.
62fn 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/// When writing, often the logic should be generic over whether or not checksums are generated.
78/// This provides that and a handy way to convert to the more general ExtentMode that eventually
79/// stores it on disk.
80#[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/// The mode of operation when setting extended attributes. This is the same as the fidl definition
117/// but is replicated here so we don't have fuchsia.io structures in the api, so this can be used
118/// on host.
119#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub enum SetExtendedAttributeMode {
121    /// Create the extended attribute if it doesn't exist, replace the value if it does.
122    Set,
123    /// Create the extended attribute if it doesn't exist, fail if it does.
124    Create,
125    /// Replace the extended attribute value if it exists, fail if it doesn't.
126    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    /// The object doesn't use encryption.
141    None,
142
143    /// The object has keys that are cached (which means unwrapping occurs on-demand) with
144    /// KeyManager.
145    CachedKeys,
146
147    /// The object has permanent keys registered with KeyManager.
148    PermanentKeys,
149}
150
151#[derive(PartialEq, Debug)]
152enum OverwriteBitmaps {
153    None,
154    Some {
155        /// The block bitmap of a partial overwrite extent in the tree.
156        extent_bitmap: BitVec,
157        /// A bitmap of the blocks written to by the current overwrite.
158        write_bitmap: BitVec,
159        /// BitVec doesn't have a slice equivalent, so for a particular section of the write we
160        /// keep track of an offset in the bitmaps to operate on.
161        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/// When writing to Overwrite ranges, we need to emit whether a set of checksums for a device range
214/// is the first write to that region or not. This tracks one such range so we can use it after the
215/// write to break up the returned checksum list.
216#[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            // If there is no bitmap, then the overwrite range is fully written to. However, we
232            // could still be within the journal flush window where one of the blocks was written
233            // to for the first time to put it in this state, so we still need to emit the
234            // checksums in case replay needs them.
235            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                // bitmap.get returning true means the block is initialized and therefore has been
251                // written to before.
252                if working_range.is_first_write != bitmaps.get_from_extent_bitmap(i).unwrap() {
253                    // is_first_write is tracking opposite of what comes back from the bitmap, so
254                    // if the are still opposites we continue our current range.
255                    working_range.checksum_range.end += 1;
256                    working_range.device_range.end += block_size;
257                } else {
258                    // If they are the same, then we need to make a new chunk.
259                    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
275/// StoreObjectHandle is the lowest-level, untyped handle to an object with the id [`object_id`] in
276/// a particular store, [`owner`]. It provides functionality shared across all objects, such as
277/// reading and writing attributes and managing encryption keys.
278///
279/// Since it's untyped, it doesn't do any object kind validation, and is generally meant to
280/// implement higher-level typed handles.
281///
282/// For file-like objects with a data attribute, DataObjectHandle implements traits and helpers for
283/// doing more complex extent management and caches the content size.
284///
285/// For directory-like objects, Directory knows how to add and remove child objects and enumerate
286/// its children.
287pub 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 this isn't being scheduled this will purposely result in fast looping
329                        // when it does. This will be insightful about the state of the thread and
330                        // task scheduling.
331                        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    /// Make a new StoreObjectHandle for the object with id [`object_id`] in store [`owner`].
346    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    /// Get the default set of transaction options for this object. This is mostly the overall
380    /// default, modified by any [`HandleOptions`] held by this handle.
381    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    // If |transaction| has an impending mutation for the underlying object, returns that.
414    // Otherwise, looks up the object from the tree.
415    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    // Returns the amount deallocated.
423    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    // Writes aligned data (that should already be encrypted) to the given offset and computes
489    // checksums if requested. The aligned data must be from a single logical file range.
490    //
491    // `flags` are forwarded to the underlying device as `WriteOptions::flags` (e.g.
492    // `WriteFlags::PRE_BARRIER`).
493    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    /// Flushes the underlying device.  This is expensive and should be used sparingly.
546    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            // The only way for these to fail are if the volume is inconsistent.
566            *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                // The allocated and deallocated shouldn't exceed the max size of the file which is
576                // bound within i64.
577                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            // This can occur when the object mutation is created from an object in the tree which
586            // was corrupt.
587            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    /// Zeroes the given range.  The range must be aligned.  Returns the amount of data deallocated.
621    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    // Returns a new aligned buffer (reading the head and tail blocks if necessary) with a copy of
643    // the data from `buf`.
644    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        // Deal with head alignment.
658        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        // Deal with tail alignment.
667        if aligned.end > end {
668            let end_block_offset = aligned.end - block_size;
669            // There's no need to read the tail block if we read it as part of the head block.
670            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    /// Trim an attribute's extents, potentially adding a graveyard trim entry if more trimming is
689    /// needed, so the transaction can be committed without worrying about leaking data.
690    ///
691    /// This doesn't update the size stored in the attribute value - the caller is responsible for
692    /// doing that to keep the size up to date.
693    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            // Add the object to the graveyard in case the following transactions don't get
708            // replayed.
709            let graveyard_id = store.graveyard_directory_object_id();
710            // Check if the object is already in the graveyard.
711            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    /// Reads and decrypts a singular logical range.
732    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    /// Returns the specified key. If `key_id` is None, it will try and return the fscrypt key if
778    /// it is present, or the volume key if it isn't. If the fscrypt key is present, but the key
779    /// cannot be unwrapped, then this will return `FxfsError::NoKey`. If the volume is not
780    /// encrypted, this returns None.
781    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        // Ensure that if the key we receive uses inline encryption, barriers should be enabled.
822        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    /// Returns a stream of extents for the given attribute ID, mapping logical offsets to device
835    /// offsets. Extents are sorted by logical offset and non-overlapping. Extents with no physical
836    /// backing are skipped.
837    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                                // Extent with no physical backing (e.g. deleted extent).
875                                ExtentValue::None => None,
876                            };
877
878                            // Advance the iterator before returning so the next invocation sees the
879                            // following entry.
880                            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                                // Skip extents with no physical backing; loop to check the next
891                                // entry.
892                                continue;
893                            }
894                        }
895                        // No more entries matching this object_id and attribute_id.
896                        _ => return None,
897                    }
898                }
899            },
900        ))
901    }
902
903    /// This will only work for a non-permanent volume data key. This is designed to be used with
904    /// extended attributes where we'll only create the key on demand for directories and encrypted
905    /// files.
906    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        // Fast path: try and get keys from the cache.
913        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        // Next, see if the keys are already created.
920        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                        /* permanent= */ false,
931                        /* force= */ 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        // Proceed to create the key.  The transaction holds the required locks.
949        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        // Add new cipher to cloned cipher set. This will replace existing one
953        // if transaction is successful.
954        cipher_set.add_key(VOLUME_DATA_KEY_ID, CipherHolder::Cipher(cipher.clone()));
955        let cipher_set = Arc::new(cipher_set);
956
957        // Arrange for the CipherSet to be added to the cache when (and if) the transaction
958        // commits.
959        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                    /* permanent= */ 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    /// Reads up to `buf.len()` bytes from attribute `attribute_id` starting at `offset`.
997    ///
998    /// Both `offset` and `buf.len()` must be aligned to the object's `block_size()`.
999    ///
1000    /// Returns the number of bytes read. If `offset >= size`, returns 0. Holes/sparse extents
1001    /// within the read range are zero-filled. Callers should not make any assumptions about the
1002    /// contents of the buffer past the returned read amount.
1003    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    /// Read `buf.len()` bytes from the attribute `attribute_id`, starting at `offset`, into `buf`.
1045    /// It's required that a read lock on this attribute id is taken before this is called.
1046    ///
1047    /// Both `offset` and `buf.len()` must be aligned to the object's `block_size()`.
1048    ///
1049    /// This function doesn't do any size checking - any portion of `buf` past the end of the file
1050    /// will be filled with zeros. The caller is responsible for enforcing the file size on reads.
1051    /// This is because, just looking at the extents, we can't tell the difference between the file
1052    /// actually ending and there just being a section at the end with no data (since attributes
1053    /// are sparse).
1054    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    /// Reads an entire attribute.
1079    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                // Attribute was deleted.
1090                ObjectValue::None => Ok(None),
1091                _ => Err(FxfsError::Inconsistent.into()),
1092            },
1093            _ => Ok(None),
1094        }
1095    }
1096
1097    /// Reads an entire attribute pointed to by `iter`. `iter` must be pointing to the
1098    /// `AttributeKey::Attribute` of the attribute.
1099    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                // TODO(https://fxbug.dev/42073113): size > max buffer size
1114                (
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    /// Reads block-aligned data from `offset` for attribute `attribute_id` into `buf` using `iter`.
1131    ///
1132    /// Both `offset` and `buf.len()` must be aligned to the object's `block_size()`, and `buf`'s
1133    /// memory buffer range start must be aligned to the device block size.
1134    ///
1135    /// Extents from `iter` are read and decrypted in parallel. Holes and any unallocated portions
1136    /// within `buf` are zero-filled.
1137    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                            // The iterator shouldn't start entirely before the offset.
1177                            ensure!(extent_key.end > offset, FxfsError::Inconsistent);
1178                            (extent_key.0.clone(), *device_offset, *key_id, mode)
1179                        }
1180                        ExtentValue::None => {
1181                            // Treat deleted extents like holes.
1182                            iter.advance().await?;
1183                            continue;
1184                        }
1185                    }
1186                }
1187                _ => break,
1188            };
1189            if extent_key.start > offset {
1190                // Zero fill holes in the attribute.
1191                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    /// Writes potentially unaligned data at `device_offset` and returns checksums if requested.
1255    /// The data will be encrypted if necessary.  `buf` is mutable as an optimization, since the
1256    /// write may require encryption, we can encrypt the buffer in-place rather than copying to
1257    /// another buffer if the write is already aligned.
1258    ///
1259    /// NOTE: This will not create keys if they are missing (it will fail with an error if that
1260    /// happens to be the case).
1261    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    /// Same as `write_at`, but allows passing `WriteFlags` (such as `WriteFlags::PRE_BARRIER`) to
1281    /// the underlying device write.
1282    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    /// Writes to multiple ranges with data provided in `buf`. This function is specifically
1321    /// designed for migration purposes, allowing raw writes to the device without updating
1322    /// object metadata like allocated size or mtime. It's essential for scenarios where
1323    /// data needs to be transferred directly without triggering standard filesystem operations.
1324    #[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    /// This is a low-level write function that writes to multiple ranges. Users should generally
1338    /// use `multi_write` instead of this function as this does not update the object's allocated
1339    /// size, mtime, atime, etc.
1340    ///
1341    /// Returns (allocated, deallocated) bytes on success.
1342    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        // The only key we allow to be created on-the-fly is a non permanent key wrapped with the
1358        // volume data key.
1359        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, /* TODO(https://fxbug.dev/421269588): plumb through device_offset. */
1383                        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            // If inline encryption is NOT supported, this loop should only happen once.
1416            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        // Only store checksums in the journal if barriers are not enabled.
1516        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    /// Writes to multiple ranges with data provided in `buf`.  The buffer can be modified in place
1525    /// if encryption takes place.  The ranges must all be aligned and no change to content size is
1526    /// applied; the caller is responsible for updating size if required.  If `key_id` is None, it
1527    /// means pick the default key for the object which is the fscrypt key if present, or the volume
1528    /// data key, or no key if it's an unencrypted file.
1529    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    /// Write data to overwrite extents with the provided set of ranges. This makes a strong
1546    /// assumption that the ranges are actually going to be already allocated overwrite extents and
1547    /// will error out or do something wrong if they aren't. It also assumes the ranges passed to
1548    /// it are sorted.
1549    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, /* TODO(https://fxbug.dev/421269588): plumb through device_offset. */
1575                        r.start,
1576                        MutPtrByteSlice::from(head),
1577                    )?;
1578                    slice = tail;
1579                }
1580            }
1581        }
1582
1583        let mut range_iter = ranges.iter();
1584        // There should be at least one range if the buffer has data in it
1585        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 this extent ends before the target range starts (not possible on the
1612                    // first loop because of the query parameters but possible on further loops),
1613                    // advance until we find a the next one we care about.
1614                    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                    // The ranges passed to this function should already by allocated, so
1631                    // extent records should exist for them.
1632                    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                // We've either run past the end of the existing extents or something is wrong with
1740                // the tree. The main section should break if it finishes the ranges, so either
1741                // case, this is an error.
1742                _ => 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        // Only store checksums in the journal if barriers are not enabled.
1750        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    /// Writes an attribute that should not already exist and therefore does not require trimming.
1764    /// Breaks up the write into multiple transactions if `data.len()` is larger than `batch_size`.
1765    /// If writing the attribute requires multiple transactions, adds the attribute to the
1766    /// graveyard. The caller is responsible for removing the attribute from the graveyard when it
1767    /// commits the last transaction.  This always writes using a key wrapped with the volume data
1768    /// key.
1769    #[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            // Do not commit the last chunk.
1816            if i < num_chunks - 1 {
1817                transaction.commit_and_continue().await?;
1818            }
1819        }
1820        Ok(())
1821    }
1822
1823    /// Writes an entire attribute. Returns whether or not the attribute needs to continue being
1824    /// trimmed - if the new data is shorter than the old data, this will trim any extents beyond
1825    /// the end of the new size, but if there were too many for a single transaction, a commit
1826    /// needs to be made before trimming again, so the responsibility is left to the caller so as
1827    /// to not accidentally split the transaction when it's not in a consistent state.  This will
1828    /// write using the volume data key; the fscrypt key is not supported.
1829    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        // Seek to the first extended attribute key for this object.
1883        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            // Skip deleted extended attributes.
1889            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                    // Once we hit a record belonging to another object, or one that is not an
1897                    // extended attribute key, we have reached the end of this object's extended
1898                    // attributes. Subsequent objects' records will start with lower variants
1899                    // (e.g. ObjectKeyData::Object) which trigger this break.
1900                    _ => break,
1901                }
1902            }
1903            iter.advance().await?;
1904        }
1905        Ok(out)
1906    }
1907
1908    /// Looks up the values for the extended attribute `fio::SELINUX_CONTEXT_NAME`, returning it
1909    /// if it is found inline. If it is not inline, it will request use of the
1910    /// `get_extended_attributes` method. If the entry doesn't exist at all, returns None.
1911    pub async fn get_inline_selinux_context(&self) -> Result<Option<fio::SelinuxContext>, Error> {
1912        // This optimization is only useful as long as the attribute is smaller than inline sizes.
1913        // Avoid reading the data out of the attributes.
1914        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        // NB: We need to take this lock before we potentially look up the value to prevent racing
1965        // with another set.
1966        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            // If we already have an attribute id allocated for this extended attribute, we always
2015            // use it, even if the value has shrunk enough to be stored inline. We don't need to
2016            // worry about trimming here for the same reason we don't need to worry about it when
2017            // we delete xattrs - they simply aren't large enough to ever need more than one
2018            // transaction.
2019            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            // If there isn't an existing attribute id and we are going to store the value in
2030            // an attribute, find the next empty attribute id in the range. We search for fxfs
2031            // attribute records specifically, instead of the extended attribute records, because
2032            // even if the extended attribute record is removed the attribute may not be fully
2033            // trimmed yet.
2034            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 means the key passed to seek wasn't found. That means the first
2042                    // attribute is available and we can just stop right away.
2043                    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                            // This attribute was once used but is now deleted, so it's safe to use
2051                            // again.
2052                            break;
2053                        }
2054                        if attribute_id < *attr_id {
2055                            // We found a gap - use it.
2056                            break;
2057                        } else if attribute_id == *attr_id {
2058                            // This attribute id is in use, try the next one.
2059                            attribute_id = attribute_id.next();
2060                            if attribute_id == AttributeId::XATTR_RANGE_END {
2061                                bail!(FxfsError::NoSpace);
2062                            }
2063                        }
2064                        // If we don't hit either of those cases, we are still moving through the
2065                        // extent keys for the current attribute, so just keep advancing until the
2066                        // attribute id changes.
2067                    }
2068                    // As we are working our way through the iterator, if we hit anything that
2069                    // doesn't have our object id or attribute key data, we've gone past the end of
2070                    // this section and can stop.
2071                    _ => break,
2072                }
2073                iter.advance().await?;
2074            }
2075
2076            // We know this won't need trimming because it's a new attribute.
2077            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        // NB: The API says we have to return an error if the attribute doesn't exist, so we have
2096        // to look it up first to make sure we have a record of it before we delete it. Make sure
2097        // we take a lock and make a transaction before we do so we don't race with other
2098        // operations.
2099        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 the attribute wasn't stored inline, we need to deallocate all the extents too. This
2121        // would normally need to interact with the graveyard for correctness - if there are too
2122        // many extents to delete to fit in a single transaction then we could potentially have
2123        // consistency issues. However, the maximum size of an extended attribute is small enough
2124        // that it will never come close to that limit even in the worst case, so we just delete
2125        // everything in one shot.
2126        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            // In case you didn't read the comment above - this should not be used to delete
2136            // arbitrary attributes!
2137            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    /// Returns a future that will pre-fetches the keys so as to avoid paying the performance
2152    /// penalty later. Must ensure that the object is not removed before the future completes.
2153    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                            /* permanent= */ false,
2167                            /* force= */ 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/// When truncating an object, sometimes it might not be possible to complete the transaction in a
2187/// single transaction, in which case the caller needs to finish trimming the object in subsequent
2188/// transactions (by calling ObjectStore::trim).
2189#[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        // Write 8KiB at offset 0 and 8KiB at offset 16KiB, leaving a hole in between.
2263        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            /* permanent_keys: */ 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        // Expect 2 physical mappings
2288        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        // This test is intended to trip a potential race condition in remove. Removing an
2298        // attribute that doesn't exist is an error, so we need to check before we remove, but if
2299        // we aren't careful, two parallel removes might both succeed in the check and then both
2300        // remove the value.
2301        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            /* permanent_keys: */ 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        // Try to remove the attribute twice at the same time. One should succeed in the race and
2322        // return Ok, and the other should fail the race and return NOT_FOUND.
2323        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        // This test is intended to trip a potential race in set when using the create flag,
2343        // similar to above. If the create mode is set, we need to check that the attribute isn't
2344        // already created, but if two parallel creates both succeed in that check, and we aren't
2345        // careful with locking, they will both succeed and one will overwrite the other.
2346        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            /* permanent_keys: */ false,
2351            HandleOptions::default(),
2352            false,
2353        ));
2354        let basic_a = basic.clone();
2355        let basic_b = basic.clone();
2356
2357        // Try to set the attribute twice at the same time. One should succeed in the race and
2358        // return Ok, and the other should fail the race and return ALREADY_EXISTS.
2359        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        // Make sure we can object the same attribute being set again.
2457        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        // Probe the fxfs attributes to make sure it did the expected thing. This relies on inside
2545        // knowledge of how the attribute id is chosen.
2546        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        // Make sure we can object the same attribute being set again.
2563        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            // Make sure expected attributes are still available.
2616            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        // Unlink the file
2670        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        // With a small attribute, we don't expect it to write to an fxfs attribute.
2722        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        // Once the value is above the threshold, we expect it to get upgraded to an fxfs
2746        // attribute.
2747        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        // Even though we are back under the threshold, we still expect it to be stored in an fxfs
2773        // attribute, because we don't downgrade to inline once we've allocated one.
2774        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        // That should have taken up all the attributes we've allocated to extended attributes, so
2874        // this one should return ERR_NO_SPACE.
2875        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        // But inline attributes don't need an attribute number, so it should work fine.
2888        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        // And updating existing ones should be okay.
2898        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        // And we should be able to remove an attribute and set another one.
2912        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        // When writing, multi_write will deallocate old extents that overlap with the new data,
2928        // but it doesn't trim anything beyond that, since it doesn't know what the total size will
2929        // be. write_attr does know, because it writes the whole attribute at once, so we need to
2930        // make sure it cleans up properly.
2931        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        // Writing two separate ranges, even if they are contiguous, forces them to be separate
2941        // extent records.
2942        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    // Running on target only, to use fake time features in the executor.
3011    #[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            // Too early.
3031            TestExecutor::advance_to(make_time(5)).await;
3032            receiver.try_recv().expect_err("Should not have message");
3033
3034            // First message.
3035            TestExecutor::advance_to(make_time(10)).await;
3036            assert_eq!(1, receiver.recv().expect("Receiving"));
3037
3038            // Too early for the next.
3039            TestExecutor::advance_to(make_time(15)).await;
3040            receiver.try_recv().expect_err("Should not have message");
3041
3042            // Missed one. They'll be spooled up.
3043            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        // Watchdog is dropped, nothing should trigger.
3049        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        // No bitmap means one chunk that covers the whole range
3058        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        // Write 5 MiB of data so that reading it spans > MAX_HASH_PARTITIONS (5 > 4 partitions).
3195        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        // Flush to create persistent layers with bloom filters.
3201        object.owner().flush().await.expect("flush failed");
3202
3203        // Read the entire 5 MiB buffer back.
3204        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}