bitrot streaming

Signed-off-by: junxiang Mu <1948535941@qq.com>
This commit is contained in:
junxiang Mu
2024-10-30 17:30:15 +08:00
parent e6cd184cd6
commit c0e6751893
6 changed files with 254 additions and 32 deletions
Generated
+65
View File
@@ -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"
+1
View File
@@ -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
+176 -23
View File
@@ -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<u8>,
_algorithm: BitrotAlgorithm,
_sum: Vec<u8>,
}
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<dyn Write>;
type BitrotWriter = Box<dyn Write + Send>;
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<BitrotWriter> {
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<dyn ReadAt>;
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<u8> {
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<Option<Vec<u8>>>,
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<Self> {
let hasher = algo.new();
let (tx, mut rx) = mpsc::channel::<Option<Vec<u8>>>(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<Self>) -> Result<()> {
let _ = self.tx.send(None).await?;
let _ = self.task.await;
Ok(())
}
}
struct StreamingBitrotReader {
disk: DiskStore,
_data: Vec<u8>,
volume: String,
file_path: String,
till_offset: usize,
curr_offset: usize,
hasher: Hasher,
shard_size: usize,
buf: Vec<u8>,
hash_bytes: Vec<u8>,
}
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<u8>, 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::<Vec<_>>();
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(())
}
+6 -6
View File
@@ -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;
+3
View File
@@ -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<Self>) -> Result<()> {
Ok(())
}
}
#[async_trait::async_trait]
+3 -3
View File
@@ -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)) {