fix: fix FTPS/SFTP download issues and optimize S3Client caching (#1353)

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: loverustfs <hello@rustfs.com>
This commit is contained in:
yxrxy
2026-01-04 17:28:18 +08:00
committed by GitHub
parent ffbcd3852f
commit 38c2d74d36
4 changed files with 133 additions and 125 deletions
+17 -18
View File
@@ -34,16 +34,20 @@ use s3s::dto::{GetObjectInput, PutObjectInput};
use std::fmt::Debug;
use std::path::{Path, PathBuf};
use tokio::io::AsyncRead;
use tokio_util::io::StreamReader;
use tracing::{debug, error, info, trace};
/// FTPS storage driver implementation
#[derive(Debug)]
pub struct FtpsDriver {}
#[derive(Debug, Clone)]
pub struct FtpsDriver {
fs: crate::storage::ecfs::FS,
}
impl FtpsDriver {
/// Create a new FTPS driver
pub fn new() -> Self {
Self {}
let fs = crate::storage::ecfs::FS {};
Self { fs }
}
/// Validate FTP feature support
@@ -68,9 +72,7 @@ impl FtpsDriver {
/// Create ProtocolS3Client for the given user
fn create_s3_client_for_user(&self, user: &super::server::FtpsUser) -> Result<ProtocolS3Client> {
let session_context = &user.session_context;
let fs = crate::storage::ecfs::FS {};
let s3_client = ProtocolS3Client::new(fs, session_context.access_key().to_string());
let s3_client = ProtocolS3Client::new(self.fs.clone(), session_context.access_key().to_string());
Ok(s3_client)
}
@@ -455,31 +457,28 @@ impl StorageBackend<super::server::FtpsUser> for FtpsDriver {
let mut builder = GetObjectInput::builder();
builder.set_bucket(bucket);
builder.set_key(object_key);
let mut input = builder
if start_pos > 0
&& let Ok(range) = s3s::dto::Range::parse(&format!("bytes={}-", start_pos))
{
builder.set_range(Some(range));
}
let input = builder
.build()
.map_err(|_| Error::new(ErrorKind::PermanentFileNotAvailable, "Failed to build GetObjectInput"))?;
if start_pos > 0 {
input.range = Some(
s3s::dto::Range::parse(&format!("bytes={}-", start_pos))
.map_err(|_| Error::new(ErrorKind::PermanentFileNotAvailable, "Invalid range format"))?,
);
}
match s3_client.get_object(input).await {
Ok(output) => {
if let Some(body) = output.body {
// Map the s3s/Box<dyn StdError> error to std::io::Error
let stream = body.map_err(std::io::Error::other);
// Wrap the stream in StreamReader to make it a tokio::io::AsyncRead
let reader = tokio_util::io::StreamReader::new(stream);
let reader = StreamReader::new(stream);
Ok(Box::new(reader))
} else {
Err(Error::new(ErrorKind::PermanentFileNotAvailable, "Empty object body"))
}
}
Err(e) => {
error!("Failed to get object: {}", e);
let protocol_error = map_s3_error_to_ftps(&e);
Err(Error::new(ErrorKind::PermanentFileNotAvailable, protocol_error))
}
+29 -30
View File
@@ -21,6 +21,7 @@ use futures::TryStreamExt;
use russh_sftp::protocol::{Attrs, Data, File, FileAttributes, Handle, Name, OpenFlags, Status, StatusCode, Version};
use russh_sftp::server::Handler;
use rustfs_utils::path;
use s3s::S3ErrorCode;
use s3s::dto::{DeleteBucketInput, DeleteObjectInput, GetObjectInput, ListObjectsV2Input, PutObjectInput, StreamingBlob};
use std::collections::HashMap;
use std::future::Future;
@@ -30,6 +31,7 @@ use std::sync::atomic::{AtomicU32, Ordering};
use tokio::fs::{File as TokioFile, OpenOptions};
use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
use tokio::sync::RwLock;
use tokio_util::io::StreamReader;
use tracing::{debug, error, trace};
use uuid::Uuid;
@@ -74,24 +76,25 @@ pub struct SftpHandler {
next_handle_id: Arc<AtomicU32>,
temp_dir: PathBuf,
current_dir: Arc<RwLock<String>>,
fs: crate::storage::ecfs::FS,
}
impl SftpHandler {
pub fn new(session_context: SessionContext) -> Self {
let fs = crate::storage::ecfs::FS {};
Self {
session_context,
handles: Arc::new(RwLock::new(HashMap::new())),
next_handle_id: Arc::new(AtomicU32::new(INITIAL_HANDLE_ID)),
temp_dir: std::env::temp_dir(),
current_dir: Arc::new(RwLock::new(ROOT_PATH.to_string())),
fs,
}
}
fn create_s3_client(&self) -> Result<ProtocolS3Client, StatusCode> {
// Create FS instance (empty struct that accesses global ECStore)
let fs = crate::storage::ecfs::FS {};
let client = ProtocolS3Client::new(fs, self.session_context.access_key().to_string());
Ok(client)
fn create_s3_client(&self) -> ProtocolS3Client {
ProtocolS3Client::new(self.fs.clone(), self.session_context.access_key().to_string())
}
fn parse_path(&self, path_str: &str) -> Result<(String, Option<String>), StatusCode> {
@@ -188,7 +191,7 @@ impl SftpHandler {
.await
.map_err(|_| StatusCode::PermissionDenied)?;
let s3_client = self.create_s3_client()?;
let s3_client = self.create_s3_client();
match action {
S3Action::HeadBucket => {
@@ -365,13 +368,7 @@ impl Handler for SftpHandler {
return Err(StatusCode::Failure);
}
let s3_client = match this.create_s3_client() {
Ok(c) => c,
Err(e) => {
let _ = tokio::fs::remove_file(&temp_file_path).await;
return Err(e);
}
};
let s3_client = this.create_s3_client();
let stream = tokio_util::io::ReaderStream::new(file);
let body = StreamingBlob::wrap(stream);
@@ -427,31 +424,35 @@ impl Handler for SftpHandler {
}
};
let s3_client = this.create_s3_client();
let range_end = offset + (len as u64) - 1;
let mut builder = GetObjectInput::builder();
builder.set_bucket(bucket);
builder.set_key(key);
if let Ok(range) = s3s::dto::Range::parse(&format!("bytes={}-{}", offset, range_end)) {
if offset > 0
&& let Ok(range) = s3s::dto::Range::parse(&format!("bytes={}-{}", offset, range_end))
{
builder.set_range(Some(range));
}
let s3_client = this.create_s3_client()?;
let input = builder.build().map_err(|_| StatusCode::Failure)?;
match s3_client.get_object(input).await {
Ok(output) => {
let mut data = Vec::with_capacity(len as usize);
let mut data = Vec::new();
if let Some(body) = output.body {
let stream = body.map_err(std::io::Error::other);
let mut reader = tokio_util::io::StreamReader::new(stream);
let _ = reader.read_to_end(&mut data).await;
let mut reader = StreamReader::new(stream);
reader.read_to_end(&mut data).await.map_err(|_| StatusCode::Failure)?;
}
Ok(Data { id, data })
}
Err(e) => {
debug!("S3 Read failed: {}", e);
Ok(Data { id, data: Vec::new() })
}
Err(e) => match e.code() {
S3ErrorCode::InvalidRange => Err(StatusCode::Eof),
_ => Err(map_s3_error_to_sftp_status(&e)),
},
}
}
}
@@ -538,9 +539,7 @@ impl Handler for SftpHandler {
.map_err(|_| StatusCode::PermissionDenied)?;
// List all buckets
let s3_client = this.create_s3_client().inspect_err(|&e| {
error!("SFTP Opendir - failed to create S3 client: {}", e);
})?;
let s3_client = this.create_s3_client();
let input = s3s::dto::ListBucketsInput::builder()
.build()
@@ -604,7 +603,7 @@ impl Handler for SftpHandler {
}
builder.set_delimiter(Some("/".to_string()));
let s3_client = this.create_s3_client()?;
let s3_client = this.create_s3_client();
let input = builder.build().map_err(|_| StatusCode::Failure)?;
let mut files = Vec::new();
@@ -721,7 +720,7 @@ impl Handler for SftpHandler {
..Default::default()
};
let s3_client = this.create_s3_client()?;
let s3_client = this.create_s3_client();
s3_client.delete_object(input).await.map_err(|e| {
error!("SFTP REMOVE - failed to delete object: {}", e);
StatusCode::Failure
@@ -742,7 +741,7 @@ impl Handler for SftpHandler {
.await
.map_err(|_| StatusCode::PermissionDenied)?;
let s3_client = this.create_s3_client()?;
let s3_client = this.create_s3_client();
// Check if bucket is empty
let list_input = ListObjectsV2Input {
@@ -818,7 +817,7 @@ impl Handler for SftpHandler {
.await
.map_err(|_| StatusCode::PermissionDenied)?;
let s3_client = this.create_s3_client()?;
let s3_client = this.create_s3_client();
let empty_stream = futures::stream::empty::<Result<bytes::Bytes, std::io::Error>>();
let body = StreamingBlob::wrap(empty_stream);
let input = PutObjectInput {
@@ -854,7 +853,7 @@ impl Handler for SftpHandler {
.await
.map_err(|_| StatusCode::PermissionDenied)?;
let s3_client = this.create_s3_client()?;
let s3_client = this.create_s3_client();
let input = s3s::dto::CreateBucketInput {
bucket,
..Default::default()