use filereader as asyncread

This commit is contained in:
weisd
2025-02-19 17:41:18 +08:00
parent 937a0c7dee
commit 7a7aee2049
14 changed files with 1067 additions and 499 deletions
+3
View File
@@ -1,5 +1,6 @@
pub mod handlers;
pub mod router;
mod rpc;
pub mod utils;
use common::error::Result;
@@ -11,6 +12,7 @@ use handlers::{
};
use hyper::Method;
use router::{AdminOperation, S3Router};
use rpc::regist_rpc_route;
use s3s::route::S3Route;
const ADMIN_PREFIX: &str = "/rustfs/admin";
@@ -21,6 +23,7 @@ pub fn make_admin_route() -> Result<impl S3Route> {
// 1
r.insert(Method::POST, "/", AdminOperation(&sts::AssumeRoleHandle {}))?;
regist_rpc_route(&mut r)?;
regist_user_route(&mut r)?;
r.insert(
+6 -1
View File
@@ -14,6 +14,7 @@ use s3s::S3Request;
use s3s::S3Response;
use s3s::S3Result;
use super::rpc::RPC_PREFIX;
use super::ADMIN_PREFIX;
pub struct S3Router<T> {
@@ -63,7 +64,7 @@ where
}
}
uri.path().starts_with(ADMIN_PREFIX)
uri.path().starts_with(ADMIN_PREFIX) || uri.path().starts_with(RPC_PREFIX)
}
async fn call(&self, req: S3Request<Body>) -> S3Result<S3Response<(StatusCode, Body)>> {
@@ -81,6 +82,10 @@ where
// check_access before call
async fn check_access(&self, req: &mut S3Request<Body>) -> S3Result<()> {
// TODO: check access by req.credentials
if req.uri.path().starts_with(RPC_PREFIX) {
return Ok(());
}
match req.credentials {
Some(_) => Ok(()),
None => Err(s3_error!(AccessDenied, "Signature is required")),
+97
View File
@@ -0,0 +1,97 @@
use super::router::AdminOperation;
use super::router::Operation;
use super::router::S3Router;
use crate::storage::ecfs::bytes_stream;
use common::error::Result;
use ecstore::disk::DiskAPI;
use ecstore::disk::FileReader;
use ecstore::store::find_local_disk;
use http::StatusCode;
use hyper::Method;
use matchit::Params;
use s3s::dto::StreamingBlob;
use s3s::s3_error;
use s3s::Body;
use s3s::S3Request;
use s3s::S3Response;
use s3s::S3Result;
use serde_urlencoded::from_bytes;
use tokio_util::io::ReaderStream;
use tracing::warn;
pub const RPC_PREFIX: &str = "/rustfs/rpc";
pub fn regist_rpc_route(r: &mut S3Router<AdminOperation>) -> Result<()> {
r.insert(
Method::GET,
format!("{}{}", RPC_PREFIX, "/read_file_stream").as_str(),
AdminOperation(&ReadFile {}),
)?;
Ok(())
}
// /rustfs/rpc/read_file_stream?disk={}&volume={}&path={}&offset={}&length={}"
#[derive(Debug, Default, serde::Deserialize)]
pub struct ReadFileQuery {
disk: String,
volume: String,
path: String,
offset: usize,
length: usize,
}
pub struct ReadFile {}
#[async_trait::async_trait]
impl Operation for ReadFile {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
warn!("handle ReadFile");
let query = {
if let Some(query) = req.uri.query() {
let input: ReadFileQuery =
from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get query failed1"))?;
input
} else {
ReadFileQuery::default()
}
};
let Some(disk) = find_local_disk(&query.disk).await else {
return Err(s3_error!(InvalidArgument, "disk not found"));
};
let file: FileReader = disk
.read_file_stream(&query.volume, &query.path, query.offset, query.length)
.await
.map_err(|e| s3_error!(InternalError, "read file err {}", e))?;
let s = bytes_stream(ReaderStream::new(file), query.length);
Ok(S3Response::new((StatusCode::OK, Body::from(StreamingBlob::wrap(s)))))
// let querys = req.uri.query().map(|q| {
// let mut querys = HashMap::new();
// for (k, v) in url::form_urlencoded::parse(q.as_bytes()) {
// println!("{}={}", k, v);
// querys.insert(k.to_string(), v.to_string());
// }
// querys
// });
// // TODO: file_path from root
// if let Some(file_path) = querys.and_then(|q| q.get("file_path").cloned()) {
// let file = fs::OpenOptions::new()
// .read(true)
// .open(file_path)
// .await
// .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("open file err {}", e)))?;
// let s = bytes_stream(ReaderStream::new(file), 0);
// return Ok(S3Response::new((StatusCode::OK, Body::from(StreamingBlob::wrap(s)))));
// }
// Ok(S3Response::new((StatusCode::BAD_REQUEST, Body::empty())))
}
}
+91 -91
View File
@@ -9,8 +9,7 @@ use ecstore::{
admin_server_info::get_local_server_property,
bucket::{metadata::load_bucket_metadata, metadata_sys},
disk::{
DeleteOptions, DiskAPI, DiskInfoOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, Reader,
UpdateMetadataOpts,
DeleteOptions, DiskAPI, DiskInfoOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, UpdateMetadataOpts,
},
erasure::Writer,
error::Error as EcsError,
@@ -694,103 +693,104 @@ impl Node for NodeService {
}
type ReadAtStream = ResponseStream<ReadAtResponse>;
async fn read_at(&self, request: Request<Streaming<ReadAtRequest>>) -> Result<Response<Self::ReadAtStream>, Status> {
async fn read_at(&self, _request: Request<Streaming<ReadAtRequest>>) -> Result<Response<Self::ReadAtStream>, Status> {
info!("read_at");
unimplemented!("read_at");
let mut in_stream = request.into_inner();
let (tx, rx) = mpsc::channel(128);
// let mut in_stream = request.into_inner();
// let (tx, rx) = mpsc::channel(128);
tokio::spawn(async move {
let mut file_ref = None;
while let Some(result) = in_stream.next().await {
match result {
Ok(v) => {
match file_ref.as_ref() {
Some(_) => (),
None => {
if let Some(disk) = find_local_disk(&v.disk).await {
match disk.read_file(&v.volume, &v.path).await {
Ok(file_reader) => file_ref = Some(file_reader),
Err(err) => {
tx.send(Ok(ReadAtResponse {
success: false,
data: Vec::new(),
error: Some(err_to_proto_err(&err, &format!("read file failed: {}", err))),
read_size: -1,
}))
.await
.expect("working rx");
break;
}
}
} else {
tx.send(Ok(ReadAtResponse {
success: false,
data: Vec::new(),
error: Some(err_to_proto_err(
&EcsError::new(StorageError::InvalidArgument(
Default::default(),
Default::default(),
Default::default(),
)),
"can not find disk",
)),
read_size: -1,
}))
.await
.expect("working rx");
break;
}
}
};
// tokio::spawn(async move {
// let mut file_ref = None;
// while let Some(result) = in_stream.next().await {
// match result {
// Ok(v) => {
// match file_ref.as_ref() {
// Some(_) => (),
// None => {
// if let Some(disk) = find_local_disk(&v.disk).await {
// match disk.read_file(&v.volume, &v.path).await {
// Ok(file_reader) => file_ref = Some(file_reader),
// Err(err) => {
// tx.send(Ok(ReadAtResponse {
// success: false,
// data: Vec::new(),
// error: Some(err_to_proto_err(&err, &format!("read file failed: {}", err))),
// read_size: -1,
// }))
// .await
// .expect("working rx");
// break;
// }
// }
// } else {
// tx.send(Ok(ReadAtResponse {
// success: false,
// data: Vec::new(),
// error: Some(err_to_proto_err(
// &EcsError::new(StorageError::InvalidArgument(
// Default::default(),
// Default::default(),
// Default::default(),
// )),
// "can not find disk",
// )),
// read_size: -1,
// }))
// .await
// .expect("working rx");
// break;
// }
// }
// };
let mut data = vec![0u8; v.length.try_into().unwrap()];
// let mut data = vec![0u8; v.length.try_into().unwrap()];
match file_ref
.as_mut()
.unwrap()
.read_at(v.offset.try_into().unwrap(), &mut data)
.await
{
Ok(read_size) => tx.send(Ok(ReadAtResponse {
success: true,
data,
read_size: read_size.try_into().unwrap(),
error: None,
})),
Err(err) => tx.send(Ok(ReadAtResponse {
success: false,
data: Vec::new(),
error: Some(err_to_proto_err(&err, &format!("read at failed: {}", err))),
read_size: -1,
})),
}
.await
.unwrap();
}
Err(err) => {
if let Some(io_err) = match_for_io_error(&err) {
if io_err.kind() == ErrorKind::BrokenPipe {
// here you can handle special case when client
// disconnected in unexpected way
eprintln!("\tclient disconnected: broken pipe");
break;
}
}
// match file_ref
// .as_mut()
// .unwrap()
// .read_at(v.offset.try_into().unwrap(), &mut data)
// .await
// {
// Ok(read_size) => tx.send(Ok(ReadAtResponse {
// success: true,
// data,
// read_size: read_size.try_into().unwrap(),
// error: None,
// })),
// Err(err) => tx.send(Ok(ReadAtResponse {
// success: false,
// data: Vec::new(),
// error: Some(err_to_proto_err(&err, &format!("read at failed: {}", err))),
// read_size: -1,
// })),
// }
// .await
// .unwrap();
// }
// Err(err) => {
// if let Some(io_err) = match_for_io_error(&err) {
// if io_err.kind() == ErrorKind::BrokenPipe {
// // here you can handle special case when client
// // disconnected in unexpected way
// eprintln!("\tclient disconnected: broken pipe");
// break;
// }
// }
match tx.send(Err(err)).await {
Ok(_) => (),
Err(_err) => break, // response was dropped
}
}
}
}
println!("\tstream ended");
});
// match tx.send(Err(err)).await {
// Ok(_) => (),
// Err(_err) => break, // response was dropped
// }
// }
// }
// }
// println!("\tstream ended");
// });
let out_stream = ReceiverStream::new(rx);
// let out_stream = ReceiverStream::new(rx);
Ok(tonic::Response::new(Box::pin(out_stream)))
// Ok(tonic::Response::new(Box::pin(out_stream)))
}
async fn list_dir(&self, request: Request<ListDirRequest>) -> Result<Response<ListDirResponse>, Status> {