From c0e6751893562d5334f8c32953cde429703c4cdc Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Wed, 30 Oct 2024 17:30:15 +0800 Subject: [PATCH] bitrot streaming Signed-off-by: junxiang Mu <1948535941@qq.com> --- Cargo.lock | 65 ++++++++++++ ecstore/Cargo.toml | 1 + ecstore/src/bitrot.rs | 199 +++++++++++++++++++++++++++++++---- ecstore/src/disk/local.rs | 12 +-- ecstore/src/erasure.rs | 3 + ecstore/src/heal/heal_ops.rs | 6 +- 6 files changed, 254 insertions(+), 32 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 2c884b2b6..5b1c17ecc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -502,6 +502,7 @@ dependencies = [ "lock", "netif", "nix", + "num", "num_cpus", "openssl", "path-absolutize", @@ -1172,12 +1173,76 @@ dependencies = [ "simdutf8", ] +[[package]] +name = "num" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35bd024e8b2ff75562e5f34e7f4905839deb4b22955ef5e73d2fea1b9813cb23" +dependencies = [ + "num-bigint", + "num-complex", + "num-integer", + "num-iter", + "num-rational", + "num-traits", +] + +[[package]] +name = "num-bigint" +version = "0.4.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a5e44f723f1133c9deac646763579fdb3ac745e418f2a7af9cd0c431da1f20b9" +dependencies = [ + "num-integer", + "num-traits", +] + +[[package]] +name = "num-complex" +version = "0.4.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73f88a1307638156682bada9d7604135552957b7818057dcef22705b4d509495" +dependencies = [ + "num-traits", +] + [[package]] name = "num-conv" version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "51d515d32fb182ee37cda2ccdcb92950d6a3c2893aa280e540671c2cd0f3b1d9" +[[package]] +name = "num-integer" +version = "0.1.46" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7969661fd2958a5cb096e56c8e1ad0444ac2bbcd0061bd28660485a44879858f" +dependencies = [ + "num-traits", +] + +[[package]] +name = "num-iter" +version = "0.1.45" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1429034a0490724d0075ebb2bc9e875d6503c3cf69e235a8941aa757d83ef5bf" +dependencies = [ + "autocfg", + "num-integer", + "num-traits", +] + +[[package]] +name = "num-rational" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f83d14da390562dca69fc84082e73e548e1ad308d24accdedd2720017cb37824" +dependencies = [ + "num-bigint", + "num-integer", + "num-traits", +] + [[package]] name = "num-traits" version = "0.2.19" diff --git a/ecstore/Cargo.toml b/ecstore/Cargo.toml index a7ac29104..7dbd0afff 100644 --- a/ecstore/Cargo.toml +++ b/ecstore/Cargo.toml @@ -51,6 +51,7 @@ tower.workspace = true rmp = "0.8.14" byteorder = "1.5.0" xxhash-rust = { version = "0.8.12", features = ["xxh64"] } +num = "0.4.3" num_cpus = "1.16" s3s-policy.workspace = true diff --git a/ecstore/src/bitrot.rs b/ecstore/src/bitrot.rs index dd6ab8c01..ee8ddba63 100644 --- a/ecstore/src/bitrot.rs +++ b/ecstore/src/bitrot.rs @@ -4,6 +4,10 @@ use blake2::Blake2b512; use highway::{HighwayHash, HighwayHasher, Key}; use lazy_static::lazy_static; use sha2::{digest::core_api::BlockSizeUser, Digest, Sha256}; +use tokio::{ + spawn, + sync::mpsc::{self, Sender}, task::JoinHandle, +}; use crate::{ disk::{error::DiskError, DiskStore}, @@ -113,15 +117,15 @@ impl BitrotAlgorithm { } pub struct BitrotVerifier { - algorithm: BitrotAlgorithm, - sum: Vec, + _algorithm: BitrotAlgorithm, + _sum: Vec, } impl BitrotVerifier { pub fn new(algorithm: BitrotAlgorithm, checksum: &[u8]) -> BitrotVerifier { BitrotVerifier { - algorithm, - sum: checksum.to_vec(), + _algorithm: algorithm, + _sum: checksum.to_vec(), } } } @@ -136,32 +140,40 @@ pub fn bitrot_algorithm_from_string(s: &str) -> BitrotAlgorithm { BitrotAlgorithm::HighwayHash256S } -type BitrotWriter = Box; +type BitrotWriter = Box; -pub fn new_bitrot_writer( +pub async fn new_bitrot_writer( disk: DiskStore, - _orig_volume: &str, + orig_volume: &str, volume: &str, file_path: &str, - _length: usize, + length: usize, algo: BitrotAlgorithm, shard_size: usize, -) -> BitrotWriter { - Box::new(WholeBitrotWriter::new(disk, volume, file_path, algo, shard_size)) +) -> Result { + if algo == BitrotAlgorithm::HighwayHash256S { + return Ok(Box::new( + StreamingBitrotWriter::new(disk, orig_volume, volume, file_path, length, algo, shard_size).await?, + )); + } + Ok(Box::new(WholeBitrotWriter::new(disk, volume, file_path, algo, shard_size))) } type BitrotReader = Box; pub fn new_bitrot_reader( disk: DiskStore, - _data: &[u8], + data: &[u8], bucket: &str, file_path: &str, till_offset: usize, algo: BitrotAlgorithm, sum: &[u8], - _shard_size: usize, + shard_size: usize, ) -> BitrotReader { + if algo == BitrotAlgorithm::HighwayHash256S { + return Box::new(StreamingBitrotReader::new(disk, data, bucket, file_path, algo, till_offset, shard_size)); + } Box::new(WholeBitrotReader::new(disk, bucket, file_path, algo, till_offset, sum)) } @@ -173,12 +185,11 @@ pub fn bitrot_writer_sum(w: &BitrotWriter) -> Vec { Vec::new() } -pub fn bitrot_shard_file_size(size: i64, _shard_size: i64, algo: BitrotAlgorithm) -> i64 { +pub fn bitrot_shard_file_size(size: usize, shard_size: usize, algo: BitrotAlgorithm) -> usize { if algo != BitrotAlgorithm::HighwayHash256S { return size; } - todo!() - // ceil_frac(size, shard_size) * algo.new().size() + size + size.div_ceil(shard_size) * algo.new().size() + size } pub fn bitrot_verify( @@ -278,6 +289,148 @@ impl ReadAt for WholeBitrotReader { } } +struct StreamingBitrotWriter { + hasher: Hasher, + tx: Sender>>, + task: JoinHandle<()>, +} + +impl StreamingBitrotWriter { + pub async fn new( + disk: DiskStore, + orig_volume: &str, + volume: &str, + file_path: &str, + length: usize, + algo: BitrotAlgorithm, + shard_size: usize, + ) -> Result { + let hasher = algo.new(); + let (tx, mut rx) = mpsc::channel::>>(10); + + let total_file_size = length.div_ceil(shard_size) * hasher.size() + length; + let mut writer = disk.create_file(orig_volume, volume, file_path, total_file_size).await?; + + let task = spawn(async move { + loop { + if let Some(Some(buf)) = rx.recv().await { + let _ = writer.write(&buf).await.unwrap(); + continue; + } + + break; + } + }); + + Ok(StreamingBitrotWriter { hasher, tx, task }) + } +} + +#[async_trait::async_trait] +impl Write for StreamingBitrotWriter { + fn as_any(&self) -> &dyn Any { + self + } + + async fn write(&mut self, buf: &[u8]) -> Result<()> { + if buf.is_empty() { + return Ok(()); + } + self.hasher.reset(); + self.hasher.update(&buf); + let hash_bytes = self.hasher.clone().finalize(); + println!("hash_bytes len: {}, buf len: {}", hash_bytes.len(), buf.len()); + let _ = self.tx.send(Some(hash_bytes)).await?; + let _ = self.tx.send(Some(buf.to_vec())).await?; + + Ok(()) + } + + async fn close(self: Box) -> Result<()> { + let _ = self.tx.send(None).await?; + let _ = self.task.await; + Ok(()) + } +} + +struct StreamingBitrotReader { + disk: DiskStore, + _data: Vec, + volume: String, + file_path: String, + till_offset: usize, + curr_offset: usize, + hasher: Hasher, + shard_size: usize, + buf: Vec, + hash_bytes: Vec, +} + +impl StreamingBitrotReader { + pub fn new( + disk: DiskStore, + data: &[u8], + volume: &str, + file_path: &str, + algo: BitrotAlgorithm, + till_offset: usize, + shard_size: usize, + ) -> Self { + let hasher = algo.new(); + Self { + disk, + _data: data.to_vec(), + volume: volume.to_string(), + file_path: file_path.to_string(), + till_offset: till_offset.div_ceil(shard_size) * hasher.size() + till_offset, + curr_offset: 0, + hash_bytes: Vec::with_capacity(hasher.size()), + hasher, + shard_size, + buf: Vec::new(), + } + } +} + +#[async_trait::async_trait] +impl ReadAt for StreamingBitrotReader { + async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec, usize)> { + println!("in read_at"); + if offset % self.shard_size != 0 { + return Err(Error::new(DiskError::Unexpected)); + } + if self.buf.is_empty() { + self.curr_offset = offset; + let stream_offset = (offset / self.shard_size) * self.hasher.size() + offset; + let buf_len = self.till_offset - stream_offset; + let mut file = self.disk.read_file(&self.volume, &self.file_path).await?; + println!("stream_offset: {}, buf_len: {}", stream_offset, buf_len); + let (buf, _) = file.read_at(stream_offset, buf_len).await?; + println!("buf: {:?}, len: {}", buf, buf.len()); + self.buf = buf; + } + if offset != self.curr_offset { + return Err(Error::new(DiskError::Unexpected)); + } + + self.hash_bytes = self.buf.drain(0..self.hash_bytes.capacity()).collect(); + let buf = self.buf.drain(0..length).collect::>(); + println!("self.buf: {:?}", self.buf); + self.hasher.reset(); + self.hasher.update(&buf); + let actual = self.hasher.clone().finalize(); + if actual != self.hash_bytes { + println!("except: {:?}, actual: {:?}", self.hash_bytes, actual); + return Err(Error::new(DiskError::FileCorrupt)); + } + + let readed_len = buf.len(); + self.curr_offset += readed_len; + + Ok((buf, readed_len)) + } +} + #[cfg(test)] mod test { use std::{collections::HashMap, fs}; @@ -347,9 +500,6 @@ mod test { #[tokio::test] async fn test_all_bitrot_algorithms() -> Result<()> { for algo in BITROT_ALGORITHMS.keys() { - if *algo == BitrotAlgorithm::HighwayHash256S { - continue; - } test_bitrot_reader_writer_algo(algo.clone()).await?; } @@ -366,19 +516,22 @@ mod test { let opt = DiskOption::default(); let disk = new_disk(&ep, &opt).await?; let _ = disk.make_volume(volume).await?; - let mut writer = new_bitrot_writer(disk.clone(), "", volume, file_path, 35, algo.clone(), 10); + let mut writer = new_bitrot_writer(disk.clone(), "", volume, file_path, 35, algo.clone(), 10).await?; let _ = writer.write(b"aaaaaaaaaa").await?; let _ = writer.write(b"aaaaaaaaaa").await?; let _ = writer.write(b"aaaaaaaaaa").await?; let _ = writer.write(b"aaaaa").await?; - let mut reader = new_bitrot_reader(disk, b"", volume, file_path, 35, algo, &bitrot_writer_sum(&writer), 10); + let sum = bitrot_writer_sum(&writer); + let _ = writer.close().await?; + + let mut reader = new_bitrot_reader(disk, b"", volume, file_path, 35, algo, &sum, 10); let read_len = 10; (_, _) = reader.read_at(0, read_len).await?; - (_, _) = reader.read_at(0, read_len).await?; - (_, _) = reader.read_at(0, read_len).await?; - (_, _) = reader.read_at(0, read_len / 2).await?; + (_, _) = reader.read_at(10, read_len).await?; + (_, _) = reader.read_at(20, read_len).await?; + (_, _) = reader.read_at(30, read_len / 2).await?; Ok(()) } diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index d37346aee..31bb69a9c 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -147,9 +147,9 @@ impl LocalDisk { id: disk_id.to_string(), ..Default::default() }; - if root { - return Err(Error::new(DiskError::DriveIsRoot)); - } + // if root { + // return Err(Error::new(DiskError::DriveIsRoot)); + // } // disk_info.healing = Ok(disk_info) @@ -185,9 +185,9 @@ impl LocalDisk { disk.minor = info.minor; disk.fstype = info.fstype; - if root { - return Err(Error::new(DiskError::DriveIsRoot)); - } + // if root { + // return Err(Error::new(DiskError::DriveIsRoot)); + // } if info.nrrequests > 0 { disk.nrrequests = info.nrrequests; diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index 4aac5c0ff..5eecdd696 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -328,6 +328,9 @@ impl Erasure { pub trait Write { fn as_any(&self) -> &dyn Any; async fn write(&mut self, buf: &[u8]) -> Result<()>; + async fn close(self: Box) -> Result<()> { + Ok(()) + } } #[async_trait::async_trait] diff --git a/ecstore/src/heal/heal_ops.rs b/ecstore/src/heal/heal_ops.rs index 87ae312f2..d0ef9e555 100644 --- a/ecstore/src/heal/heal_ops.rs +++ b/ecstore/src/heal/heal_ops.rs @@ -483,7 +483,7 @@ impl AllHealState { hsp.client_address = he.client_address.clone(); hsp.start_time = he.start_time; - he.stop(); + he.stop().await; loop { if he.has_ended().await { @@ -493,7 +493,7 @@ impl AllHealState { sleep(Duration::from_secs(1)).await; } - self.mu.write().await; + let _ = self.mu.write().await; self.heal_seq_map.remove(path); } else { hsp.client_token = "unknown".to_string(); @@ -526,7 +526,7 @@ impl AllHealState { } } - self.mu.write().await; + let _ = self.mu.write().await; for (k, v) in self.heal_seq_map.iter() { if !v.has_ended().await && (has_profix(&k, path_s) || has_profix(path_s, &k)) {