diff --git a/zerofs/src/cli/debug.rs b/zerofs/src/cli/debug.rs index bbb4bef8..cc6c6127 100644 --- a/zerofs/src/cli/debug.rs +++ b/zerofs/src/cli/debug.rs @@ -227,9 +227,9 @@ mod tests { fn test_describe_key_decodes_every_kind() { let codec = KeyCodec::new(); let cases: Vec<(bytes::Bytes, KeyPrefix, &str)> = vec![ - (codec.inode_key(42), KeyPrefix::Inode, "inode_id=42"), + (codec.inode_key(42).into(), KeyPrefix::Inode, "inode_id=42"), ( - codec.extent_key(7, 99), + codec.extent_key(7, 99).into(), KeyPrefix::Extent, "inode_id=7, extent_index=99", ), @@ -289,7 +289,7 @@ mod tests { // A recognized kind with a truncated payload still classifies (and // counts) but falls back to the raw rendering instead of misdecoding. let inode_key = codec.inode_key(1); - let truncated = &inode_key[..inode_key.len() - 1]; + let truncated = &inode_key.as_ref()[..inode_key.as_ref().len() - 1]; let (prefix, detail) = describe_key(&codec, truncated).unwrap(); assert_eq!(prefix, KeyPrefix::Inode); assert!(detail.starts_with("raw=")); diff --git a/zerofs/src/cli/server.rs b/zerofs/src/cli/server.rs index 16bb01f4..699641d9 100644 --- a/zerofs/src/cli/server.rs +++ b/zerofs/src/cli/server.rs @@ -1394,7 +1394,7 @@ mod tests { let mut id = 0; while id < INODES { let v = db - .get_bytes(&codec.inode_key(id)) + .get_bytes(codec.inode_key(id).as_ref()) .await .expect("get") .expect("inode present"); @@ -1430,7 +1430,8 @@ mod tests { let mut batch = WriteBatch::new(); for i in 0..(INODES / 4) { let id = extent * (INODES / 4) + i; - batch.put_bytes(codec.inode_key(id), Bytes::from(vec![id as u8; 64])); + batch + .put_bytes(codec.inode_key(id).into(), Bytes::from(vec![id as u8; 64])); } raw.write_with_options(batch, &WriteOptions::default()) .await diff --git a/zerofs/src/db.rs b/zerofs/src/db.rs index 20ec9ade..37823e2f 100644 --- a/zerofs/src/db.rs +++ b/zerofs/src/db.rs @@ -540,7 +540,7 @@ impl Db { self.inner.is_read_only() } - pub async fn get_bytes(&self, key: &Bytes) -> Result> { + pub async fn get_bytes(&self, key: &[u8]) -> Result> { self.get_bytes_at(key, DurabilityLevel::Memory).await } @@ -549,7 +549,7 @@ impl Db { /// Unlike a serving read, this waits through a recoverable lease /// suspension. Commit preparation uses it before write admission so a /// temporary authority gap does not turn a safe retry into an I/O error. - pub(crate) async fn get_bytes_internal(&self, key: &Bytes) -> Result> { + pub(crate) async fn get_bytes_internal(&self, key: &[u8]) -> Result> { self.check_internal_lease().await?; self.check_closing()?; let result = self @@ -561,13 +561,13 @@ impl Db { } /// Point read seeing only object-storage-durable data. - pub async fn get_bytes_durable(&self, key: &Bytes) -> Result> { + pub async fn get_bytes_durable(&self, key: &[u8]) -> Result> { self.get_bytes_at(key, DurabilityLevel::Remote).await } async fn get_bytes_at( &self, - key: &Bytes, + key: &[u8], durability_filter: DurabilityLevel, ) -> Result> { self.check_lease()?; @@ -581,7 +581,7 @@ impl Db { async fn get_bytes_at_unchecked( &self, - key: &Bytes, + key: &[u8], durability_filter: DurabilityLevel, ) -> Result> { let read_options = ReadOptions { diff --git a/zerofs/src/fs/boot.rs b/zerofs/src/fs/boot.rs index aae27938..139b3031 100644 --- a/zerofs/src/fs/boot.rs +++ b/zerofs/src/fs/boot.rs @@ -108,7 +108,7 @@ impl ZeroFS { }; let root_inode_key = key_codec.inode_key(0); - if db.get_bytes(&root_inode_key).await?.is_none() { + if db.get_bytes(root_inode_key.as_ref()).await?.is_none() { if db.is_read_only() { return Err(anyhow::anyhow!( "Cannot initialize filesystem in read-only mode. Root inode does not exist." @@ -134,7 +134,7 @@ impl ZeroFS { }; let serialized = bincode::serialize(&Inode::Directory(root_dir))?; db.put_with_options( - &root_inode_key, + &root_inode_key.into(), &serialized, &PutOptions::default(), &WriteOptions::default(), diff --git a/zerofs/src/fs/filter_policy.rs b/zerofs/src/fs/filter_policy.rs index 5ec63432..2b4fabd5 100644 --- a/zerofs/src/fs/filter_policy.rs +++ b/zerofs/src/fs/filter_policy.rs @@ -99,12 +99,12 @@ mod tests { let p = ZerofsPrefixExtractor; let codec = KeyCodec::new(); for k in [ - codec.inode_key(1), + codec.inode_key(1).into(), codec.dir_cookie_counter_key(1), codec.stats_shard_key(0), codec.system_counter_key(), codec.tombstone_key(123, 1), - codec.extent_key(42, 7), + codec.extent_key(42, 7).into(), ] { assert_eq!(extract(&p, k), None); } diff --git a/zerofs/src/fs/flush_coordinator.rs b/zerofs/src/fs/flush_coordinator.rs index 58e778fe..11058252 100644 --- a/zerofs/src/fs/flush_coordinator.rs +++ b/zerofs/src/fs/flush_coordinator.rs @@ -510,8 +510,8 @@ mod tests { fn dirty_batch() -> WriteBatch { let codec = crate::fs::key_codec::KeyCodec::new(); let mut batch = WriteBatch::new(); - batch.put_bytes(codec.inode_key(1), Bytes::from_static(b"inode")); - batch.put_bytes(codec.extent_key(1, 0), Bytes::from_static(b"extent")); + batch.put_bytes(codec.inode_key(1).into(), Bytes::from_static(b"inode")); + batch.put_bytes(codec.extent_key(1, 0).into(), Bytes::from_static(b"extent")); batch } diff --git a/zerofs/src/fs/key_codec.rs b/zerofs/src/fs/key_codec.rs index 053a0d0f..c5d48887 100644 --- a/zerofs/src/fs/key_codec.rs +++ b/zerofs/src/fs/key_codec.rs @@ -43,25 +43,25 @@ const PREFIX_ORPHAN: u8 = 0x08; const PREFIX_SEGCOUNT: u8 = 0x09; const PREFIX_EXTENT: u8 = 0xFE; -const SYSTEM_COUNTER_SUBTYPE: u8 = 0x01; +const SYSTEM_COUNTER_KEY: &[u8; 6] = b"meta\x06\x01"; // HA: the highest shipped replication batch seqno (with its writer epoch) that // has been flushed into this data db. Written atomically with each shipped // batch so a promoted standby can prune its tail to exactly what the db already // holds (see write_coordinator + takeover replay). -const SYSTEM_HA_SEQNO_SUBTYPE: u8 = 0x02; +const SYSTEM_HA_SEQNO_KEY: &[u8; 6] = b"meta\x06\x02"; // Durability lineage token (see fsync-honesty / ZeroFS::lineage_token). A single // u64 identifying the current unbroken durable lineage. Regenerated at a cold // bootstrap or a Solo-tainted takeover; carried forward unchanged at an untainted // takeover (so a clean failover keeps a client's fsync transparent). -const SYSTEM_LINEAGE_SUBTYPE: u8 = 0x03; +const SYSTEM_LINEAGE_KEY: &[u8; 6] = b"meta\x06\x03"; // Solo taint: set to the lineage token that was live when the leader first // downgraded to Solo replication. A takeover reads it to decide keep-vs-regenerate // the lineage token (taint == stored lineage => the lineage may be missing acked // Solo writes => regenerate, so those writes' fsync fails instead of reporting success). -const SYSTEM_TAINT_SUBTYPE: u8 = 0x04; +const SYSTEM_TAINT_KEY: &[u8; 6] = b"meta\x06\x04"; // Wall-clock epoch-seconds (u64 LE via encode_u64) of the last completed slow // orphan sweep. -const SYSTEM_ORPHAN_SWEEP_SUBTYPE: u8 = 0x05; +const SYSTEM_ORPHAN_SWEEP_KEY: &[u8; 6] = b"meta\x06\x05"; const U64_SIZE: usize = std::mem::size_of::(); @@ -70,6 +70,39 @@ pub const META_DOMAIN: &[u8] = b"meta"; /// Domain prefix for bulk extent data. pub const EXTENT_DOMAIN: &[u8] = b"extent"; +const INODE_KEY_SIZE: usize = META_DOMAIN.len() + 1 + U64_SIZE; +const EXTENT_KEY_SIZE: usize = EXTENT_DOMAIN.len() + 1 + U64_SIZE * 2; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct InodeKey([u8; INODE_KEY_SIZE]); + +impl AsRef<[u8]> for InodeKey { + fn as_ref(&self) -> &[u8] { + &self.0 + } +} + +impl From for Bytes { + fn from(key: InodeKey) -> Self { + Self::copy_from_slice(key.as_ref()) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct ExtentKey([u8; EXTENT_KEY_SIZE]); + +impl AsRef<[u8]> for ExtentKey { + fn as_ref(&self) -> &[u8] { + &self.0 + } +} + +impl From for Bytes { + fn from(key: ExtentKey) -> Self { + Self::copy_from_slice(key.as_ref()) + } +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum KeyPrefix { Inode, @@ -189,12 +222,12 @@ impl KeyCodec { /// Total bytes in a complete inode key. pub fn inode_key_size(&self) -> usize { - self.id_offset(KeyPrefix::Inode) + U64_SIZE + INODE_KEY_SIZE } /// Total bytes in a complete extent key. pub fn extent_key_size(&self) -> usize { - self.id_offset(KeyPrefix::Extent) + U64_SIZE * 2 + EXTENT_KEY_SIZE } /// Total bytes in a complete tombstone key. @@ -207,19 +240,22 @@ impl KeyCodec { self.id_offset(KeyPrefix::Orphan) + U64_SIZE } - pub fn inode_key(&self, inode_id: InodeId) -> Bytes { - let mut key = Vec::with_capacity(self.inode_key_size()); - self.push_prefix(&mut key, KeyPrefix::Inode); - key.extend_from_slice(&inode_id.to_be_bytes()); - Bytes::from(key) + pub fn inode_key(&self, inode_id: InodeId) -> InodeKey { + let mut key = [0; INODE_KEY_SIZE]; + key[..META_DOMAIN.len()].copy_from_slice(META_DOMAIN); + key[self.kind_offset(KeyPrefix::Inode)] = PREFIX_INODE; + key[self.id_offset(KeyPrefix::Inode)..].copy_from_slice(&inode_id.to_be_bytes()); + InodeKey(key) } - pub fn extent_key(&self, inode_id: InodeId, extent_index: u64) -> Bytes { - let mut key = Vec::with_capacity(self.extent_key_size()); - self.push_prefix(&mut key, KeyPrefix::Extent); - key.extend_from_slice(&inode_id.to_be_bytes()); - key.extend_from_slice(&extent_index.to_be_bytes()); - Bytes::from(key) + pub fn extent_key(&self, inode_id: InodeId, extent_index: u64) -> ExtentKey { + let mut key = [0; EXTENT_KEY_SIZE]; + key[..EXTENT_DOMAIN.len()].copy_from_slice(EXTENT_DOMAIN); + key[self.kind_offset(KeyPrefix::Extent)] = PREFIX_EXTENT; + let id_offset = self.id_offset(KeyPrefix::Extent); + key[id_offset..id_offset + U64_SIZE].copy_from_slice(&inode_id.to_be_bytes()); + key[id_offset + U64_SIZE..].copy_from_slice(&extent_index.to_be_bytes()); + ExtentKey(key) } pub fn parse_extent_key(&self, key: &[u8]) -> Option { @@ -354,10 +390,7 @@ impl KeyCodec { } pub fn system_counter_key(&self) -> Bytes { - let mut key = Vec::with_capacity(self.id_offset(KeyPrefix::System) + 1); - self.push_prefix(&mut key, KeyPrefix::System); - key.push(SYSTEM_COUNTER_SUBTYPE); - Bytes::from(key) + Bytes::from_static(SYSTEM_COUNTER_KEY) } /// Key for HA provenance flushed atomically with each replicated leader @@ -365,10 +398,7 @@ impl KeyCodec { /// history, and highest locally applied ship attempt; takeover validates the /// volatile tail and exact-result ledger against it. pub fn ha_seqno_key(&self) -> Bytes { - let mut key = Vec::with_capacity(self.id_offset(KeyPrefix::System) + 1); - self.push_prefix(&mut key, KeyPrefix::System); - key.push(SYSTEM_HA_SEQNO_SUBTYPE); - Bytes::from(key) + Bytes::from_static(SYSTEM_HA_SEQNO_KEY) } pub(crate) fn encode_ha_stamp(stamp: &HaStamp) -> Bytes { @@ -388,27 +418,18 @@ impl KeyCodec { /// Key for the durability lineage token (a single u64). pub fn lineage_key(&self) -> Bytes { - let mut key = Vec::with_capacity(self.id_offset(KeyPrefix::System) + 1); - self.push_prefix(&mut key, KeyPrefix::System); - key.push(SYSTEM_LINEAGE_SUBTYPE); - Bytes::from(key) + Bytes::from_static(SYSTEM_LINEAGE_KEY) } /// Key for the Solo taint (the lineage token that went Solo). pub fn taint_key(&self) -> Bytes { - let mut key = Vec::with_capacity(self.id_offset(KeyPrefix::System) + 1); - self.push_prefix(&mut key, KeyPrefix::System); - key.push(SYSTEM_TAINT_SUBTYPE); - Bytes::from(key) + Bytes::from_static(SYSTEM_TAINT_KEY) } /// Key for the last-orphan-sweep wall-clock timestamp (epoch seconds, a u64 via /// [`Self::encode_u64`]). pub fn last_orphan_sweep_key(&self) -> Bytes { - let mut key = Vec::with_capacity(self.id_offset(KeyPrefix::System) + 1); - self.push_prefix(&mut key, KeyPrefix::System); - key.push(SYSTEM_ORPHAN_SWEEP_SUBTYPE); - Bytes::from(key) + Bytes::from_static(SYSTEM_ORPHAN_SWEEP_KEY) } pub fn encode_u64(value: u64) -> Bytes { @@ -668,25 +689,25 @@ mod tests { let inode_id = 7u64; let extent_index = 99u64; let key = codec.extent_key(inode_id, extent_index); - assert_eq!(codec.parse_extent_key(&key), Some(extent_index)); + assert_eq!(codec.parse_extent_key(key.as_ref()), Some(extent_index)); } #[test] fn test_layout_routing() { let codec = KeyCodec::new(); let inode_key = codec.inode_key(0); - assert!(inode_key.starts_with(META_DOMAIN)); - assert_eq!(inode_key[META_DOMAIN.len()], PREFIX_INODE); + assert!(inode_key.as_ref().starts_with(META_DOMAIN)); + assert_eq!(inode_key.as_ref()[META_DOMAIN.len()], PREFIX_INODE); let extent_key = codec.extent_key(0, 0); - assert!(extent_key.starts_with(EXTENT_DOMAIN)); - assert_eq!(extent_key[EXTENT_DOMAIN.len()], PREFIX_EXTENT); + assert!(extent_key.as_ref().starts_with(EXTENT_DOMAIN)); + assert_eq!(extent_key.as_ref()[EXTENT_DOMAIN.len()], PREFIX_EXTENT); let tombstone = codec.tombstone_key(0, 0); assert!(tombstone.starts_with(META_DOMAIN)); // No metadata key should be misrouted into the extent domain. - assert!(!inode_key.starts_with(EXTENT_DOMAIN)); + assert!(!inode_key.as_ref().starts_with(EXTENT_DOMAIN)); assert!(!tombstone.starts_with(EXTENT_DOMAIN)); } @@ -709,7 +730,7 @@ mod tests { assert!(sc.as_ref() >= start.as_ref() && sc.as_ref() < end.as_ref()); let ino = codec.inode_key(9); assert!(!(ino.as_ref() >= start.as_ref() && ino.as_ref() < end.as_ref())); - assert_eq!(codec.parse_segcount_key(&ino), None); + assert_eq!(codec.parse_segcount_key(ino.as_ref()), None); } #[test] @@ -862,7 +883,10 @@ mod tests { ParsedKey::Unknown )); let inode_key = codec.inode_key(1); - assert!(matches!(codec.parse_key(&inode_key), ParsedKey::Unknown)); + assert!(matches!( + codec.parse_key(inode_key.as_ref()), + ParsedKey::Unknown + )); } } @@ -885,8 +909,8 @@ mod prop_tests { let codec = KeyCodec::new(); let ka = codec.extent_key(a.0, a.1); let kb = codec.extent_key(b.0, b.1); - prop_assert_eq!(codec.parse_extent_key(&ka), Some(a.1)); - prop_assert_eq!(codec.parse_extent_key(&kb), Some(b.1)); + prop_assert_eq!(codec.parse_extent_key(ka.as_ref()), Some(a.1)); + prop_assert_eq!(codec.parse_extent_key(kb.as_ref()), Some(b.1)); prop_assert_eq!(ka.as_ref().cmp(kb.as_ref()), a.cmp(&b)); } @@ -950,7 +974,7 @@ mod prop_tests { name in prop::collection::vec(any::(), 0..40), ) { let codec = KeyCodec::new(); - prop_assert_eq!(codec.parse_extent_key(&codec.inode_key(ino)), None); + prop_assert_eq!(codec.parse_extent_key(codec.inode_key(ino).as_ref()), None); prop_assert_eq!(codec.parse_extent_key(&codec.tombstone_key(x, ino)), None); prop_assert_eq!(codec.parse_extent_key(&codec.orphan_key(ino)), None); prop_assert_eq!(codec.parse_extent_key(&codec.dir_scan_key(ino, x)), None); diff --git a/zerofs/src/fs/mod.rs b/zerofs/src/fs/mod.rs index 50007446..2028e401 100644 --- a/zerofs/src/fs/mod.rs +++ b/zerofs/src/fs/mod.rs @@ -223,6 +223,7 @@ mod tests { let io = codec.id_offset(KeyPrefix::Inode); for id in [0u64, 42, 999] { let key = codec.inode_key(id); + let key = key.as_ref(); assert_eq!(key[ko], u8::from(KeyPrefix::Inode)); assert_eq!(&key[io..io + 8], &id.to_be_bytes()); } @@ -237,6 +238,7 @@ mod tests { let io = codec.id_offset(KeyPrefix::Extent); for (ino, idx) in [(1u64, 0u64), (42, 10), (999, 999)] { let key = codec.extent_key(ino, idx); + let key = key.as_ref(); assert_eq!(key[ko], u8::from(KeyPrefix::Extent)); assert_eq!(&key[io..io + 8], &ino.to_be_bytes()); assert_eq!(&key[io + 8..io + 16], &idx.to_be_bytes()); diff --git a/zerofs/src/fs/store/directory.rs b/zerofs/src/fs/store/directory.rs index a35a047d..d315bc72 100644 --- a/zerofs/src/fs/store/directory.rs +++ b/zerofs/src/fs/store/directory.rs @@ -587,7 +587,7 @@ mod tests { let inode = test_file_inode(10); db.put_with_options( - &codec.inode_key(inode_id), + &codec.inode_key(inode_id).into(), &bincode::serialize(&inode).unwrap(), &PutOptions::default(), &WriteOptions::default(), @@ -653,7 +653,7 @@ mod tests { let inode = test_file_inode(10); db.put_with_options( - &codec.inode_key(inode_id), + &codec.inode_key(inode_id).into(), &bincode::serialize(&inode).unwrap(), &PutOptions::default(), &WriteOptions::default(), diff --git a/zerofs/src/fs/store/extent/read.rs b/zerofs/src/fs/store/extent/read.rs index 452e93db..4702e8a2 100644 --- a/zerofs/src/fs/store/extent/read.rs +++ b/zerofs/src/fs/store/extent/read.rs @@ -60,7 +60,7 @@ impl ExtentStore { /// hole. Resolves the extent key's `FrameLoc` then fetches the frame. pub async fn get(&self, id: InodeId, extent_idx: u64) -> Result, FsError> { let key = self.key_codec.extent_key(id, extent_idx); - let encoded = match self.db.get_bytes(&key).await { + let encoded = match self.db.get_bytes(key.as_ref()).await { Ok(v) => v, Err(e) => { error!( @@ -102,7 +102,7 @@ impl ExtentStore { // a hole; otherwise the error is real. let reresolved = self .db - .get_bytes(&key) + .get_bytes(key.as_ref()) .await .map_err(|_| FsError::IoError)?; match Self::decode_reresolved_extent(id, extent_idx, reresolved.as_ref())? { @@ -337,7 +337,7 @@ impl ExtentStore { async fn segment_at(&self, id: InodeId, extent: u64) -> Option { let key = self.key_codec.extent_key(id, extent); self.db - .get_bytes(&key) + .get_bytes(key.as_ref()) .await .ok() .flatten() @@ -613,7 +613,7 @@ mod tests { }; let key = store.key_codec.extent_key(1, 5); let mut txn = db.new_transaction().unwrap(); - txn.put_bytes(&key, Bytes::copy_from_slice(&bogus.encode())); + txn.put_bytes(&key.into(), Bytes::copy_from_slice(&bogus.encode())); commit(&store, txn).await; assert!(matches!(store.get(1, 5).await, Err(FsError::IoError))); @@ -632,7 +632,10 @@ mod tests { // extents 0..=2 so the ranged-scan path (not `get`) resolves it. let key = store.key_codec.extent_key(1, 1); let mut txn = db.new_transaction().unwrap(); - txn.put_bytes(&key, Bytes::from_static(&[0u8; FrameLoc::ENCODED_LEN - 1])); + txn.put_bytes( + &key.into(), + Bytes::from_static(&[0u8; FrameLoc::ENCODED_LEN - 1]), + ); commit(&store, txn).await; let r = store diff --git a/zerofs/src/fs/store/extent/reclaim/cycle.rs b/zerofs/src/fs/store/extent/reclaim/cycle.rs index ad1a8a5d..9fa3d4be 100644 --- a/zerofs/src/fs/store/extent/reclaim/cycle.rs +++ b/zerofs/src/fs/store/extent/reclaim/cycle.rs @@ -791,12 +791,12 @@ async fn verify_segment_reclaimable( let key = store.key_codec.extent_key(inode, extent); let current = store .db - .get_bytes(&key) + .get_bytes(key.as_ref()) .await .map_err(|_| FsError::IoError)?; let durable = store .db - .get_bytes_durable(&key) + .get_bytes_durable(key.as_ref()) .await .map_err(|_| FsError::IoError)?; let points_here = |encoded: Option, view: &str| match encoded { @@ -1073,7 +1073,7 @@ mod tests { KeyCodec::encode_segcount(0, total), ); txn.put_bytes( - &store.key_codec.extent_key(1, 0), + &store.key_codec.extent_key(1, 0).into(), Bytes::from_static(b"malformed FrameLoc"), ); commit(&store, txn).await; diff --git a/zerofs/src/fs/store/extent/reclaim/repack.rs b/zerofs/src/fs/store/extent/reclaim/repack.rs index ac8209e1..94821236 100644 --- a/zerofs/src/fs/store/extent/reclaim/repack.rs +++ b/zerofs/src/fs/store/extent/reclaim/repack.rs @@ -248,7 +248,7 @@ async fn plan_source(store: &ExtentStore, segid: Segid) -> Result( @@ -447,7 +447,7 @@ async fn repoint_inode( let key = store.key_codec.extent_key(inode, extent); if let Some(enc) = store .db - .get_bytes(&key) + .get_bytes(key.as_ref()) .await .map_err(|_| FsError::IoError)? && let Some(loc) = FrameLoc::decode(&enc) @@ -456,7 +456,7 @@ async fn repoint_inode( // frame between gather and swap. && loc == old_loc { - txn.put_bytes(&key, Bytes::copy_from_slice(&new_loc.encode())); + txn.put_bytes(&key.into(), Bytes::copy_from_slice(&new_loc.encode())); store.seg_delta(&mut txn, old_loc.segid, -(loc.byte_len as i64), 0); store.seg_delta(&mut txn, new_loc.segid, new_loc.byte_len as i64, 0); swapped += 1; diff --git a/zerofs/src/fs/store/extent/test_util.rs b/zerofs/src/fs/store/extent/test_util.rs index 90005d6e..806b3f83 100644 --- a/zerofs/src/fs/store/extent/test_util.rs +++ b/zerofs/src/fs/store/extent/test_util.rs @@ -380,7 +380,7 @@ pub(super) async fn frameloc_of( extent: u64, ) -> Option { let key = store.key_codec.extent_key(inode, extent); - db.get_bytes(&key) + db.get_bytes(key.as_ref()) .await .unwrap() .and_then(|b| FrameLoc::decode(&b)) diff --git a/zerofs/src/fs/store/extent/write.rs b/zerofs/src/fs/store/extent/write.rs index db763e99..e54c72f3 100644 --- a/zerofs/src/fs/store/extent/write.rs +++ b/zerofs/src/fs/store/extent/write.rs @@ -95,7 +95,7 @@ impl ExtentStore { /// caller's job (see [`Self::delete_range`], which debits as it scans). pub fn delete(&self, txn: &mut Transaction, id: InodeId, extent_idx: u64) { let key = self.key_codec.extent_key(id, extent_idx); - txn.delete_bytes(&key); + txn.delete_bytes(&key.into()); } /// Stage live/total byte deltas for `segid`'s counter onto the txn; the @@ -263,7 +263,7 @@ impl ExtentStore { byte_len: 4 + sealed_len, }; txn.put_bytes( - &self.key_codec.extent_key(id, *extent), + &self.key_codec.extent_key(id, *extent).into(), Bytes::copy_from_slice(&loc.encode()), ); // Credit the frame just appended: both live and total. @@ -443,7 +443,7 @@ impl ExtentStore { // Read the existing content of any partially-overwritten extent (full // overwrites and extents past EOF need no read). - let existing_extents: HashMap = stream::iter(start_extent..=end_extent) + let mut existing_extents: HashMap = stream::iter(start_extent..=end_extent) .map(|extent_idx| { let extent_start = extent_idx * EXTENT_SIZE as u64; let extent_end = extent_start + EXTENT_SIZE as u64; @@ -486,7 +486,10 @@ impl ExtentStore { let extent: Bytes = if write_start == 0 && write_end == EXTENT_SIZE { data.slice(data_offset..data_offset + write_len) } else { - let mut buf = BytesMut::from(existing_extents[&extent_idx].as_ref()); + // Consume the decoded extent so uniquely owned storage can be reused. + // Static zero extents and shared buffers are copied by BytesMut. + let existing = existing_extents.remove(&extent_idx).expect("extent loaded"); + let mut buf = BytesMut::from(existing); buf[write_start..write_end] .copy_from_slice(&data[data_offset..data_offset + write_len]); buf.freeze() @@ -881,7 +884,7 @@ mod tests { write_and_check(&store, &db, &mut model, 0, &vec![0u8; EXTENT_SIZE]).await; // The extent key for extent 0 must be gone. let key = store.key_codec.extent_key(1, 0); - assert!(db.get_bytes(&key).await.unwrap().is_none()); + assert!(db.get_bytes(key.as_ref()).await.unwrap().is_none()); } #[tokio::test] @@ -1016,7 +1019,7 @@ mod tests { let data = incompressible(7, 3 * EXTENT_SIZE); write_and_check(&store, &db, &mut model, 0, &data).await; let key = store.key_codec.extent_key(1, 1); - let encoded = db.get_bytes(&key).await.unwrap().unwrap(); + let encoded = db.get_bytes(key.as_ref()).await.unwrap().unwrap(); let loc = FrameLoc::decode(&encoded).unwrap(); let shipped = store.read_frame_for_ship(loc).unwrap(); assert_eq!(shipped.len(), loc.byte_len as usize); @@ -1041,9 +1044,9 @@ mod tests { &data[EXTENT_SIZE..2 * EXTENT_SIZE] ); assert!(store.unflushed_bytes() > 0); - let ops = store.enrich_repl_ops(vec![ReplOp::Put(key.clone(), encoded.clone())]); + let ops = store.enrich_repl_ops(vec![ReplOp::Put(key.into(), encoded.clone())]); assert!(matches!(&ops[..], [ReplOp::PutFrame(k, v, frame)] - if k == &key && v == &encoded && frame == &shipped)); + if k.as_ref() == key.as_ref() && v == &encoded && frame == &shipped)); assert!( store .frame_chunks_in_ram( @@ -1070,7 +1073,7 @@ mod tests { assert_eq!(directory.len(), 3); assert_eq!(directory[1].byte_offset, loc.byte_offset); assert!(matches!( - &store.enrich_repl_ops(vec![ReplOp::Put(key, encoded)])[..], + &store.enrich_repl_ops(vec![ReplOp::Put(key.into(), encoded)])[..], [ReplOp::Put(_, _)] )); } diff --git a/zerofs/src/fs/store/inode.rs b/zerofs/src/fs/store/inode.rs index a397bd34..d167157b 100644 --- a/zerofs/src/fs/store/inode.rs +++ b/zerofs/src/fs/store/inode.rs @@ -67,7 +67,7 @@ impl InodeStore { let data = self .db - .get_bytes(&key) + .get_bytes(key.as_ref()) .await .map_err(|e| { let error = FsError::from_db_error(&e); @@ -129,14 +129,14 @@ impl InodeStore { ) -> Result<(), Box> { let key = self.key_codec.inode_key(id); let data = Bytes::from(bincode::serialize(inode)?); - txn.put_bytes(&key, data); + txn.put_bytes(&key.into(), data); txn.invalidate_cached_inode(id); Ok(()) } pub fn delete(&self, txn: &mut Transaction, id: InodeId) { let key = self.key_codec.inode_key(id); - txn.delete_bytes(&key); + txn.delete_bytes(&key.into()); txn.invalidate_cached_inode(id); } diff --git a/zerofs/src/fs/write_coordinator.rs b/zerofs/src/fs/write_coordinator.rs index 7b61f95b..d1f321b4 100644 --- a/zerofs/src/fs/write_coordinator.rs +++ b/zerofs/src/fs/write_coordinator.rs @@ -679,7 +679,7 @@ mod tests { let key = codec().extent_key(41, 0); let locks = fs.lock_manager.acquire_multi(vec![41, 42]).await; let mut txn = Transaction::new(); - txn.put_bytes(&key, Bytes::from_static(b"applied")); + txn.put_bytes(&key.into(), Bytes::from_static(b"applied")); let mutation = LockedMutation::new(txn, locks); let write_barrier = fs.db.flush_barrier().write_owned().await; let coordinator = fs.write_coordinator.clone(); @@ -715,7 +715,7 @@ mod tests { .await .expect("multi-inode locks were not returned after apply"); assert_eq!( - fs.db.get_bytes(&key).await.unwrap(), + fs.db.get_bytes(key.as_ref()).await.unwrap(), Some(Bytes::from_static(b"applied")) ); } @@ -830,9 +830,9 @@ mod tests { let mut txn = Transaction::new(); // Use a real codec-built key so the segment extractor accepts it. let key = codec().extent_key(1, 0); - txn.put_bytes(&key, Bytes::from_static(b"value")); + txn.put_bytes(&key.into(), Bytes::from_static(b"value")); fs.write_coordinator.commit(txn).await.unwrap(); - let v = fs.db.get_bytes(&key).await.unwrap(); + let v = fs.db.get_bytes(key.as_ref()).await.unwrap(); assert_eq!(v.as_deref(), Some(&b"value"[..])); } @@ -847,7 +847,7 @@ mod tests { let k = codec.extent_key(1, i); handles.push(tokio::spawn(async move { let mut txn = Transaction::new(); - txn.put_bytes(&k, Bytes::from(vec![1u8; 8])); + txn.put_bytes(&k.into(), Bytes::from(vec![1u8; 8])); c.commit(txn).await })); } @@ -855,7 +855,11 @@ mod tests { h.await.unwrap().unwrap(); } for i in 0u64..32 { - let v = fs.db.get_bytes(&codec.extent_key(1, i)).await.unwrap(); + let v = fs + .db + .get_bytes(codec.extent_key(1, i).as_ref()) + .await + .unwrap(); assert!(v.is_some()); } } @@ -871,7 +875,7 @@ mod tests { let k = codec.extent_key(2, i); handles.push(tokio::spawn(async move { let mut txn = Transaction::new(); - txn.put_bytes(&k, Bytes::from(vec![2u8; 8])); + txn.put_bytes(&k.into(), Bytes::from(vec![2u8; 8])); c.commit(txn).await })); } @@ -879,7 +883,11 @@ mod tests { h.await.unwrap().unwrap(); } for i in 0u64..16 { - let v = fs.db.get_bytes(&codec.extent_key(2, i)).await.unwrap(); + let v = fs + .db + .get_bytes(codec.extent_key(2, i).as_ref()) + .await + .unwrap(); assert!( v.is_some(), "extent {i} not durable after sync_writes commit" @@ -960,7 +968,7 @@ mod tests { // A commit that doesn't allocate any inode. let mut txn = Transaction::new(); - txn.put_bytes(&codec().extent_key(3, 0), Bytes::from_static(b"v")); + txn.put_bytes(&codec().extent_key(3, 0).into(), Bytes::from_static(b"v")); fs.write_coordinator.commit(txn).await.unwrap(); let after = fs.db.get_bytes(&counter_key).await.unwrap(); @@ -972,7 +980,7 @@ mod tests { // Now allocate and commit; counter must advance on disk. let _id = fs.inode_store.allocate(); let mut txn = Transaction::new(); - txn.put_bytes(&codec().extent_key(3, 1), Bytes::from_static(b"v")); + txn.put_bytes(&codec().extent_key(3, 1).into(), Bytes::from_static(b"v")); fs.write_coordinator.commit(txn).await.unwrap(); let after_allocate = fs.db.get_bytes(&counter_key).await.unwrap(); @@ -999,7 +1007,7 @@ mod tests { .unwrap(); let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(7, 0), Bytes::from_static(b"v")); + txn.put_bytes(&codec.extent_key(7, 0).into(), Bytes::from_static(b"v")); txn.add_seg_delta(&seg_key, 5, 5); txn.delete_segcount(&codec.segcount_key(1, 2), 7, 13); fs.write_coordinator @@ -1030,7 +1038,7 @@ mod tests { // that batch (and the staged counter put) to the segcount abort. let id = fs.inode_store.allocate(); let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(id, 0), Bytes::from_static(b"v")); + txn.put_bytes(&codec.extent_key(id, 0).into(), Bytes::from_static(b"v")); txn.add_seg_delta(&seg_key, 5, 5); fs.write_coordinator.commit(txn).await.unwrap_err(); @@ -1038,7 +1046,7 @@ mod tests { // covering `id` (via the id burned on abort); otherwise a restart // would hand out `id` again over this batch's durable records. let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(id, 1), Bytes::from_static(b"w")); + txn.put_bytes(&codec.extent_key(id, 1).into(), Bytes::from_static(b"w")); fs.write_coordinator.commit(txn).await.unwrap(); let persisted = fs @@ -1074,7 +1082,7 @@ mod tests { let key = codec.extent_key(inode_id, 0); handles.push(tokio::spawn(async move { let mut txn = Transaction::new(); - txn.put_bytes(&key, Bytes::from_static(b"x")); + txn.put_bytes(&key.into(), Bytes::from_static(b"x")); txn.add_stats_delta(inode_id, ((k + 1) * 10) as i64, 1); c.commit(txn).await })); @@ -1122,7 +1130,7 @@ mod tests { let key = codec.extent_key(inode_id, 1); handles.push(tokio::spawn(async move { let mut txn = Transaction::new(); - txn.put_bytes(&key, Bytes::from_static(b"y")); + txn.put_bytes(&key.into(), Bytes::from_static(b"y")); txn.add_stats_delta(inode_id, -(((k + 1) * 10) as i64), 0); c.commit(txn).await })); @@ -1339,7 +1347,7 @@ mod tests { let key = codec().extent_key(7, 0); let flushes_before = fs.flush_coordinator.completed_flush_count(); let mut txn = Transaction::new(); - txn.put_bytes(&key, Bytes::from_static(b"never-applied")); + txn.put_bytes(&key.into(), Bytes::from_static(b"never-applied")); assert_eq!( fs.write_coordinator .commit(txn) @@ -1393,7 +1401,7 @@ mod tests { let key = codec.extent_key(1, 0); let flushes_before = fs.flush_coordinator.completed_flush_count(); let mut txn = Transaction::new(); - txn.put_bytes(&key, Bytes::from_static(b"v")); + txn.put_bytes(&key.into(), Bytes::from_static(b"v")); let error = coord .commit(txn) .await @@ -1405,13 +1413,13 @@ mod tests { "a writer proven stale must not flush before returning the clean failure" ); assert!( - fs.db.get_bytes(&key).await.unwrap().is_none(), + fs.db.get_bytes(key.as_ref()).await.unwrap().is_none(), "a deposed leader must not apply the rejected batch" ); // Deposal is terminal: later batches fail too. let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(1, 1), Bytes::from_static(b"w")); + txn.put_bytes(&codec.extent_key(1, 1).into(), Bytes::from_static(b"w")); assert_eq!( coord.commit(txn).await.expect_err("deposal must be sticky"), FsError::LeaderRejectedBeforeApply @@ -1443,10 +1451,9 @@ mod tests { let first_commit = { let coordinator = fs.write_coordinator.clone(); - let first_key = first_key.clone(); tokio::spawn(async move { let mut txn = Transaction::new(); - txn.put_bytes(&first_key, Bytes::from_static(b"first")); + txn.put_bytes(&first_key.into(), Bytes::from_static(b"first")); coordinator.commit(txn).await }) }; @@ -1535,7 +1542,7 @@ mod tests { // The Solo mutation follows the lineage-taint flush. let solo_key = codec().extent_key(91, 0); let mut solo = Transaction::new(); - solo.put_bytes(&solo_key, Bytes::from_static(b"solo")); + solo.put_bytes(&solo_key.into(), Bytes::from_static(b"solo")); fs.write_coordinator.commit(solo).await.unwrap(); let ha_stamp_key = codec().ha_seqno_key(); let solo_stamp = fs @@ -1554,10 +1561,9 @@ mod tests { let op_id = [0x91; 16]; let reconnect_commit = { let coordinator = fs.write_coordinator.clone(); - let reconnect_key = reconnect_key.clone(); tokio::spawn(async move { let mut txn = Transaction::new(); - txn.put_bytes(&reconnect_key, Bytes::from_static(b"reconnected")); + txn.put_bytes(&reconnect_key.into(), Bytes::from_static(b"reconnected")); txn.set_dedup_result(op_id, crate::dedup::DedupResult::Applied); coordinator.commit(txn).await }) @@ -1579,7 +1585,11 @@ mod tests { "the reconnect batch must wait behind the Solo-base barrier" ); assert!( - fs.db.get_bytes(&reconnect_key).await.unwrap().is_none(), + fs.db + .get_bytes(reconnect_key.as_ref()) + .await + .unwrap() + .is_none(), "the reconnect mutation must not apply before the base flush" ); assert!( @@ -1627,7 +1637,10 @@ mod tests { lease.renew(std::time::Duration::from_secs(30)); let mut later = Transaction::new(); - later.put_bytes(&codec().extent_key(91, 2), Bytes::from_static(b"later")); + later.put_bytes( + &codec().extent_key(91, 2).into(), + Bytes::from_static(b"later"), + ); fs.write_coordinator .commit(later) .await @@ -1647,14 +1660,14 @@ mod tests { let codec = codec(); let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(1, 0), Bytes::from_static(b"a")); + txn.put_bytes(&codec.extent_key(1, 0).into(), Bytes::from_static(b"a")); coord.commit(txn).await.unwrap(); let baseline = fs.flush_coordinator.completed_flush_count(); control.set_sender_for_tests(None).await; for i in 1..=2u64 { let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(1, i), Bytes::from_static(b"s")); + txn.put_bytes(&codec.extent_key(1, i).into(), Bytes::from_static(b"s")); coord.commit(txn).await.unwrap(); } assert_eq!( @@ -1668,7 +1681,7 @@ mod tests { .await; let reconnect_op_id = [0xa7; 16]; let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(1, 3), Bytes::from_static(b"c")); + txn.put_bytes(&codec.extent_key(1, 3).into(), Bytes::from_static(b"c")); txn.set_dedup_result(reconnect_op_id, crate::dedup::DedupResult::Applied); coord.commit(txn).await.unwrap(); assert_eq!( @@ -1682,7 +1695,7 @@ mod tests { )); let mut txn = Transaction::new(); - txn.put_bytes(&codec.extent_key(1, 4), Bytes::from_static(b"d")); + txn.put_bytes(&codec.extent_key(1, 4).into(), Bytes::from_static(b"d")); coord.commit(txn).await.unwrap(); assert_eq!( fs.flush_coordinator.completed_flush_count(), @@ -1702,7 +1715,10 @@ mod tests { let coord = replicating_coordinator(&fs, repl); let mut txn = Transaction::new(); - txn.put_bytes(&codec().extent_key(1, 0), Bytes::from_static(b"acked")); + txn.put_bytes( + &codec().extent_key(1, 0).into(), + Bytes::from_static(b"acked"), + ); coord.commit(txn).await.unwrap(); let baseline = fs.flush_coordinator.completed_flush_count(); assert_eq!( diff --git a/zerofs/src/nbd/server.rs b/zerofs/src/nbd/server.rs index d9d0de23..10139502 100644 --- a/zerofs/src/nbd/server.rs +++ b/zerofs/src/nbd/server.rs @@ -163,7 +163,8 @@ impl NBDSession { async fn perform_handshake(&mut self) -> Result<()> { let handshake = NBDServerHandshake::new(NBD_FLAG_FIXED_NEWSTYLE | NBD_FLAG_NO_ZEROES); - let handshake_bytes = handshake.to_bytes()?; + let mut handshake_bytes = [0; NBDServerHandshake::WIRE_SIZE]; + handshake.encode(&mut handshake_bytes)?; self.writer.write_all(&handshake_bytes).await?; self.writer.flush().await?; @@ -373,7 +374,8 @@ impl NBDSession { async fn send_option_reply(&mut self, option: u32, reply_type: u32, data: &[u8]) -> Result<()> { let reply = NBDOptionReply::new(option, reply_type, data.len() as u32); - let reply_bytes = reply.to_bytes()?; + let mut reply_bytes = [0; NBDOptionReply::WIRE_SIZE]; + reply.encode(&mut reply_bytes)?; self.writer.write_all(&reply_bytes).await?; if !data.is_empty() { self.writer.write_all(data).await?; @@ -562,7 +564,8 @@ impl NBDSession { async fn send_simple_reply(&mut self, cookie: u64, error: u32, data: &[u8]) -> Result<()> { let reply = NBDSimpleReply::new(cookie, error); - let reply_bytes = reply.to_bytes()?; + let mut reply_bytes = [0; NBDSimpleReply::WIRE_SIZE]; + reply.encode(&mut reply_bytes)?; self.writer.write_all(&reply_bytes).await?; if !data.is_empty() { self.writer.write_all(data).await?; diff --git a/zerofs/src/segment.rs b/zerofs/src/segment.rs index a7cffeef..0501d0f9 100644 --- a/zerofs/src/segment.rs +++ b/zerofs/src/segment.rs @@ -184,21 +184,21 @@ pub enum SegmentError { Codec(#[from] CodecError), } -fn frame_aad(segid: Segid, frame_index: u32, inode: u64, extent: u64) -> Vec { - let mut v = Vec::with_capacity(1 + 16 + 4 + 8 + 8); - v.push(b'F'); - v.extend_from_slice(&segid.to_le_bytes()); - v.extend_from_slice(&frame_index.to_le_bytes()); - v.extend_from_slice(&inode.to_le_bytes()); - v.extend_from_slice(&extent.to_le_bytes()); +fn frame_aad(segid: Segid, frame_index: u32, inode: u64, extent: u64) -> [u8; 37] { + let mut v = [0; 37]; + v[0] = b'F'; + v[1..17].copy_from_slice(&segid.to_le_bytes()); + v[17..21].copy_from_slice(&frame_index.to_le_bytes()); + v[21..29].copy_from_slice(&inode.to_le_bytes()); + v[29..37].copy_from_slice(&extent.to_le_bytes()); v } -fn dir_aad(segid: Segid, k: u32) -> Vec { - let mut v = Vec::with_capacity(1 + 16 + 4); - v.push(b'D'); - v.extend_from_slice(&segid.to_le_bytes()); - v.extend_from_slice(&k.to_le_bytes()); +fn dir_aad(segid: Segid, k: u32) -> [u8; 21] { + let mut v = [0; 21]; + v[0] = b'D'; + v[1..17].copy_from_slice(&segid.to_le_bytes()); + v[17..21].copy_from_slice(&k.to_le_bytes()); v } @@ -725,7 +725,7 @@ pub(crate) fn read_frames_from_chunks( frame_aad(segid, fi, inode, extent), )); } - let open = |(frame, aad): (bytes::Bytes, Vec)| { + let open = |(frame, aad): (bytes::Bytes, [u8; 37])| { codec.open_owned(frame, &aad).map_err(SegmentError::from) }; let runtime = tokio::runtime::Handle::try_current().ok(); diff --git a/zerofs/src/segment_extractor.rs b/zerofs/src/segment_extractor.rs index e83373e0..c4897f50 100644 --- a/zerofs/src/segment_extractor.rs +++ b/zerofs/src/segment_extractor.rs @@ -87,7 +87,7 @@ mod tests { let ex = ZeroFsSegmentExtractor; let codec = KeyCodec::new(); assert_eq!( - ex.prefix_len(&PrefixTarget::Point(codec.inode_key(1))), + ex.prefix_len(&PrefixTarget::Point(codec.inode_key(1).into())), Some(META_DOMAIN.len()) ); assert_eq!( @@ -95,7 +95,7 @@ mod tests { Some(META_DOMAIN.len()) ); assert_eq!( - ex.prefix_len(&PrefixTarget::Point(codec.extent_key(1, 0))), + ex.prefix_len(&PrefixTarget::Point(codec.extent_key(1, 0).into())), Some(EXTENT_DOMAIN.len()) ); } diff --git a/zerofs/tests/dst/reclamation_cases.rs b/zerofs/tests/dst/reclamation_cases.rs index 08c4051f..22861f30 100644 --- a/zerofs/tests/dst/reclamation_cases.rs +++ b/zerofs/tests/dst/reclamation_cases.rs @@ -59,7 +59,7 @@ async fn create_file(fs: &ZeroFS, name: &[u8], data: &[u8]) -> InodeId { async fn location(fs: &ZeroFS, id: InodeId, extent: u64) -> FrameLoc { FrameLoc::decode( &fs.db - .get_bytes(&KeyCodec::new().extent_key(id, extent)) + .get_bytes(KeyCodec::new().extent_key(id, extent).as_ref()) .await .unwrap() .unwrap(), @@ -169,14 +169,14 @@ impl RepackFixture { ); assert!( fs.db - .get_bytes(&KeyCodec::new().extent_key(self.victim, 0)) + .get_bytes(KeyCodec::new().extent_key(self.victim, 0).as_ref()) .await .unwrap() .is_none() ); assert!( fs.db - .get_bytes(&KeyCodec::new().extent_key(self.victim, 10_010)) + .get_bytes(KeyCodec::new().extent_key(self.victim, 10_010).as_ref()) .await .unwrap() .is_none()