merge versioning, fix buff bitrot

This commit is contained in:
weisd
2024-11-02 15:12:06 +08:00
parent 09ea11c13d
commit 52a81f8fd5
9 changed files with 183 additions and 70 deletions
+1 -1
View File
@@ -1,7 +1,7 @@
use crate::error::ErrorCode;
use crate::Result as LocalResult;
use axum::{extract::State, Json};
use axum::Json;
use serde::Serialize;
use time::OffsetDateTime;
-1
View File
@@ -2,7 +2,6 @@ pub mod error;
pub mod handlers;
use axum::{extract::Request, response::Response, routing::get, BoxError, Router};
use ecstore::store::ECStore;
use error::ErrorCode;
use handlers::list_pools;
use tower::Service;
+118 -3
View File
@@ -1,11 +1,12 @@
use crate::{
disk::{error::DiskError, DiskStore},
disk::{error::DiskError, DiskStore, FileReader, FileWriter},
erasure::{ReadAt, Write},
error::{Error, Result},
store_api::BitrotAlgorithm,
};
use blake2::Blake2b512;
use blake2::Digest as _;
use bytes::Bytes;
use highway::{HighwayHash, HighwayHasher, Key};
use lazy_static::lazy_static;
use sha2::{digest::core_api::BlockSizeUser, Digest, Sha256};
@@ -14,9 +15,14 @@ use std::{
collections::HashMap,
io::{Cursor, Read},
};
use tokio::{
io::AsyncWriteExt,
spawn,
sync::mpsc::{self, Sender},
sync::{
mpsc::{self, Sender},
RwLock,
},
task::JoinHandle,
};
@@ -145,7 +151,7 @@ pub fn bitrot_algorithm_from_string(s: &str) -> BitrotAlgorithm {
BitrotAlgorithm::HighwayHash256S
}
pub type BitrotWriter = Box<dyn Write + Send>;
pub type BitrotWriter = Box<dyn Write + Send + 'static>;
pub async fn new_bitrot_writer(
disk: DiskStore,
@@ -476,6 +482,115 @@ impl ReadAt for StreamingBitrotReader {
}
}
pub struct BitrotFileWriter {
pub inner: FileWriter,
hasher: Hasher,
_shard_size: usize,
}
impl BitrotFileWriter {
pub fn new(inner: FileWriter, algo: BitrotAlgorithm, _shard_size: usize) -> Self {
let hasher = algo.new();
Self {
inner,
hasher,
_shard_size,
}
}
pub fn writer(&self) -> &FileWriter {
&self.inner
}
}
#[async_trait::async_trait]
impl Write for BitrotFileWriter {
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();
let _ = self.inner.write(&hash_bytes).await?;
let _ = self.inner.write(buf).await?;
Ok(())
}
}
pub fn new_bitrot_filewriter(inner: FileWriter, algo: BitrotAlgorithm, shard_size: usize) -> Result<BitrotWriter> {
Ok(Box::new(BitrotFileWriter::new(inner, algo, shard_size)))
}
#[derive(Debug)]
struct BitrotFileReader {
pub inner: FileReader,
till_offset: usize,
curr_offset: usize,
hasher: Hasher,
shard_size: usize,
buf: Vec<u8>,
hash_bytes: Vec<u8>,
}
impl BitrotFileReader {
pub fn new(inner: FileReader, algo: BitrotAlgorithm, till_offset: usize, shard_size: usize) -> Self {
let hasher = algo.new();
Self {
inner,
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 BitrotFileReader {
async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec<u8>, usize)> {
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 (buf, _) = self.inner.read_at(stream_offset, buf_len).await?;
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<_>>();
self.hasher.reset();
self.hasher.update(&buf);
let actual = self.hasher.clone().finalize();
if actual != self.hash_bytes {
return Err(Error::new(DiskError::FileCorrupt));
}
let readed_len = buf.len();
self.curr_offset += readed_len;
Ok((buf, readed_len))
}
}
pub fn new_bitrot_filereader(inner: FileReader, till_offset: usize, algo: BitrotAlgorithm, shard_size: usize) -> BitrotReader {
Box::new(BitrotFileReader::new(inner, algo, till_offset, shard_size))
}
#[cfg(test)]
mod test {
use std::{collections::HashMap, fs};
+56 -14
View File
@@ -644,7 +644,7 @@ pub struct ReadOptions {
pub enum FileWriter {
Local(LocalFileWriter),
Remote(RemoteFileWriter),
Buffer(Vec<u8>),
Buffer(BufferWriter),
}
#[async_trait::async_trait]
@@ -655,15 +655,41 @@ impl Write for FileWriter {
async fn write(&mut self, buf: &[u8]) -> Result<()> {
match self {
Self::Local(local_file_writer) => local_file_writer.write(buf).await,
Self::Remote(remote_file_writer) => remote_file_writer.write(buf).await,
Self::Buffer(buffer) => {
buffer.extend_from_slice(buf);
Ok(())
}
Self::Local(writer) => writer.write(buf).await,
Self::Remote(writter) => writter.write(buf).await,
Self::Buffer(writer) => writer.write(buf).await,
}
}
}
#[derive(Debug)]
pub struct BufferWriter {
pub inner: Vec<u8>,
}
impl BufferWriter {
pub fn new(inner: Vec<u8>) -> Self {
Self { inner }
}
pub fn as_ref(&self) -> &[u8] {
self.inner.as_ref()
}
}
#[async_trait::async_trait]
impl Write for BufferWriter {
fn as_any(&self) -> &dyn Any {
self
}
async fn write(&mut self, buf: &[u8]) -> Result<()> {
let _ = self.inner.write(buf).await?;
self.inner.flush().await?;
Ok(())
}
}
#[derive(Debug)]
pub struct LocalFileWriter {
pub inner: File,
@@ -778,23 +804,39 @@ impl Write for RemoteFileWriter {
pub enum FileReader {
Local(LocalFileReader),
Remote(RemoteFileReader),
Buffer(Vec<u8>),
Buffer(BufferReader),
}
#[async_trait::async_trait]
impl ReadAt for FileReader {
async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec<u8>, usize)> {
match self {
Self::Local(local_file_writer) => local_file_writer.read_at(offset, length).await,
Self::Remote(remote_file_writer) => remote_file_writer.read_at(offset, length).await,
Self::Buffer(buffer) => {
let s = &buffer[offset..offset + length];
Ok((s.to_vec(), s.len()))
}
Self::Local(reader) => reader.read_at(offset, length).await,
Self::Remote(reader) => reader.read_at(offset, length).await,
Self::Buffer(reader) => reader.read_at(offset, length).await,
}
}
}
#[derive(Debug)]
pub struct BufferReader {
pub inner: Vec<u8>,
}
impl BufferReader {
pub fn new(inner: Vec<u8>) -> Self {
Self { inner }
}
}
#[async_trait::async_trait]
impl ReadAt for BufferReader {
async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec<u8>, usize)> {
let s = &self.inner[offset..offset + length];
Ok((s.to_vec(), s.len()))
}
}
#[derive(Debug)]
pub struct LocalFileReader {
pub inner: File,
+1 -1
View File
@@ -70,7 +70,7 @@ impl InlineData {
Ok(None)
}
fn validate(&self) -> Result<()> {
pub fn validate(&self) -> Result<()> {
if self.0.is_empty() {
return Ok(());
}
-45
View File
@@ -1,45 +0,0 @@
use sha2::{
digest::{Reset, Update},
Digest, Sha256 as sha_sha256,
};
trait Hasher {
fn write(&mut self, bytes: &[u8]);
fn reset(&mut self);
fn sum(&mut self) -> impl AsRef<[u8]>;
fn size(&self) -> usize;
fn block_size(&self) -> usize;
}
struct Sha256 {
hasher: sha_sha256,
}
impl Sha256 {
pub fn new() -> Self {
Self {
hasher: sha_sha256::new(),
}
}
}
impl Hasher for Sha256 {
fn write(&mut self, bytes: &[u8]) {
Update::update(&mut self.hasher, bytes);
}
fn reset(&mut self) {
Reset::reset(&mut self.hasher);
}
fn sum(&mut self) -> impl AsRef<[u8]> {
self.hasher.clone().finalize()
}
fn size(&self) -> usize {
32
}
fn block_size(&self) -> usize {
64
}
}
-1
View File
@@ -2,7 +2,6 @@ pub mod crypto;
pub mod ellipses;
pub mod fs;
pub mod hash;
pub mod hasher;
pub mod net;
pub mod os;
pub mod path;
-1
View File
@@ -6,7 +6,6 @@ mod storage;
use clap::Parser;
use common::error::{Error, Result};
use ecstore::{
config::GLOBAL_ConfigSys,
endpoints::EndpointServerPools,
set_global_endpoints,
store::{init_local_disks, ECStore},
+7 -3
View File
@@ -512,6 +512,13 @@ impl S3 for FS {
Ok(S3Response::new(output))
}
async fn list_object_versions(
&self,
_req: S3Request<ListObjectVersionsInput>,
) -> S3Result<S3Response<ListObjectVersionsOutput>> {
Err(s3_error!(NotImplemented, "ListObjectVersions is not implemented yet"))
}
#[tracing::instrument(level = "debug", skip(self, req))]
async fn put_object(&self, req: S3Request<PutObjectInput>) -> S3Result<S3Response<PutObjectOutput>> {
let input = req.input;
@@ -529,9 +536,6 @@ impl S3 for FS {
key,
metadata,
content_length,
content_type,
checksum_sha256,
content_md5,
..
} = input;