mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-02 05:14:35 +08:00
fix(ftps): bound upload memory with multipart streaming (#8064)
* fix(ftps): bound upload memory with multipart streaming * fix(ftps): keep multipart upload within driver module
This commit is contained in:
@@ -133,6 +133,8 @@ pub struct DeleteObjectCall {
|
||||
}
|
||||
|
||||
struct Inner {
|
||||
capture_upload_bodies: bool,
|
||||
upload_bodies: Vec<Vec<u8>>,
|
||||
// Response queues. Each method pops from its own queue. Empty queue
|
||||
// plus no default means a configured-miss error.
|
||||
get_object: VecDeque<Result<GetObjectOutput, DummyError>>,
|
||||
@@ -190,6 +192,8 @@ struct Inner {
|
||||
impl Inner {
|
||||
fn new() -> Self {
|
||||
Self {
|
||||
capture_upload_bodies: false,
|
||||
upload_bodies: Vec::new(),
|
||||
get_object: VecDeque::new(),
|
||||
get_object_range: VecDeque::new(),
|
||||
put_object: VecDeque::new(),
|
||||
@@ -265,6 +269,30 @@ impl DummyBackend {
|
||||
}
|
||||
}
|
||||
|
||||
/// Opt in to consuming upload bodies for byte-for-byte protocol tests.
|
||||
pub fn capture_upload_bodies(&self) {
|
||||
self.inner.lock().expect("lock").capture_upload_bodies = true;
|
||||
}
|
||||
|
||||
pub fn upload_bodies(&self) -> Vec<Vec<u8>> {
|
||||
self.inner.lock().expect("lock").upload_bodies.clone()
|
||||
}
|
||||
|
||||
async fn record_upload_body(&self, body: &mut Option<StreamingBlob>) -> Result<(), DummyError> {
|
||||
use futures_util::TryStreamExt;
|
||||
if !self.inner.lock().expect("lock").capture_upload_bodies {
|
||||
return Ok(());
|
||||
}
|
||||
let mut bytes = Vec::new();
|
||||
if let Some(mut stream) = body.take() {
|
||||
while let Some(chunk) = stream.try_next().await.map_err(|e| DummyError::Injected(e.to_string()))? {
|
||||
bytes.extend_from_slice(&chunk);
|
||||
}
|
||||
}
|
||||
self.inner.lock().expect("lock").upload_bodies.push(bytes);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// Queue-configuration helpers. Each test stages the responses it
|
||||
// expects in order. The method pops in FIFO order.
|
||||
|
||||
@@ -629,7 +657,8 @@ impl StorageBackend for DummyBackend {
|
||||
}
|
||||
}
|
||||
|
||||
async fn put_object(&self, input: PutObjectInput, _credentials: &Credentials) -> Result<PutObjectOutput, Self::Error> {
|
||||
async fn put_object(&self, mut input: PutObjectInput, _credentials: &Credentials) -> Result<PutObjectOutput, Self::Error> {
|
||||
self.record_upload_body(&mut input.body).await?;
|
||||
// Decide control flow while holding the lock. Release before
|
||||
// awaiting so the stall path does not hold the Mutex across
|
||||
// an await point.
|
||||
@@ -789,7 +818,8 @@ impl StorageBackend for DummyBackend {
|
||||
}
|
||||
}
|
||||
|
||||
async fn upload_part(&self, input: UploadPartInput, _credentials: &Credentials) -> Result<UploadPartOutput, Self::Error> {
|
||||
async fn upload_part(&self, mut input: UploadPartInput, _credentials: &Credentials) -> Result<UploadPartOutput, Self::Error> {
|
||||
self.record_upload_body(&mut input.body).await?;
|
||||
// Record the call and decide the control flow while holding the
|
||||
// lock. Release the lock before awaiting so the stall path does
|
||||
// not hold the Mutex across an await point.
|
||||
|
||||
@@ -150,10 +150,10 @@ pub fn is_operation_supported(protocol: super::session::Protocol, action: &S3Act
|
||||
S3Action::HeadObject => true, // SIZE command
|
||||
|
||||
// Multipart operations
|
||||
S3Action::CreateMultipartUpload => false,
|
||||
S3Action::UploadPart => false,
|
||||
S3Action::CompleteMultipartUpload => false,
|
||||
S3Action::AbortMultipartUpload => false,
|
||||
S3Action::CreateMultipartUpload => true,
|
||||
S3Action::UploadPart => true,
|
||||
S3Action::CompleteMultipartUpload => true,
|
||||
S3Action::AbortMultipartUpload => true,
|
||||
S3Action::ListMultipartUploads => false,
|
||||
S3Action::ListParts => false,
|
||||
|
||||
|
||||
@@ -16,12 +16,12 @@ use crate::common::client::s3::StorageBackend as S3StorageBackend;
|
||||
use crate::common::gateway::S3Action;
|
||||
use crate::common::gateway::authorize_operation;
|
||||
use async_trait::async_trait;
|
||||
use futures_util::stream;
|
||||
use rustfs_utils::MaskedAccessKey;
|
||||
use rustfs_utils::path;
|
||||
use s3s::dto::*;
|
||||
use std::fmt::Debug;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use tokio::io::AsyncRead;
|
||||
use tracing::{debug, error};
|
||||
use unftp_core::storage::{Error, ErrorKind, Fileinfo, Metadata, Result, StorageBackend};
|
||||
@@ -94,12 +94,12 @@ impl Metadata for FtpsMetadata {
|
||||
/// FTPS storage driver implementation
|
||||
pub struct FtpsDriver<S> {
|
||||
/// Storage backend for S3 operations
|
||||
storage: S,
|
||||
storage: Arc<S>,
|
||||
}
|
||||
|
||||
impl<S> Debug for FtpsDriver<S>
|
||||
where
|
||||
S: S3StorageBackend + Debug,
|
||||
S: S3StorageBackend + Debug + 'static,
|
||||
{
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("FtpsDriver").field("storage", &"StorageBackend").finish()
|
||||
@@ -108,11 +108,13 @@ where
|
||||
|
||||
impl<S> FtpsDriver<S>
|
||||
where
|
||||
S: S3StorageBackend + Debug,
|
||||
S: S3StorageBackend + Debug + 'static,
|
||||
{
|
||||
/// Create a new FTPS driver with the given storage backend
|
||||
pub fn new(storage: S) -> Self {
|
||||
Self { storage }
|
||||
Self {
|
||||
storage: Arc::new(storage),
|
||||
}
|
||||
}
|
||||
|
||||
/// List all buckets (for root path)
|
||||
@@ -232,7 +234,7 @@ where
|
||||
#[async_trait]
|
||||
impl<S> StorageBackend<super::server::FtpsUser> for FtpsDriver<S>
|
||||
where
|
||||
S: S3StorageBackend + Debug,
|
||||
S: S3StorageBackend + Debug + 'static,
|
||||
{
|
||||
type Metadata = FtpsMetadata;
|
||||
|
||||
@@ -544,55 +546,18 @@ where
|
||||
.await
|
||||
.map_err(|_| Error::new(ErrorKind::PermanentFileNotAvailable, "Access denied"))?;
|
||||
|
||||
// Convert AsyncRead to bytes
|
||||
let bytes_vec = {
|
||||
let mut buffer = Vec::new();
|
||||
let mut reader = bytes;
|
||||
tokio::io::copy(&mut reader, &mut buffer)
|
||||
.await
|
||||
.map_err(|e| Error::new(ErrorKind::TransientFileNotAvailable, e.to_string()))?;
|
||||
buffer
|
||||
};
|
||||
|
||||
let file_size = bytes_vec.len();
|
||||
|
||||
let mut put_builder = PutObjectInput::builder();
|
||||
put_builder.set_bucket(bucket.clone());
|
||||
put_builder.set_key(key.clone());
|
||||
put_builder.set_content_length(Some(file_size as i64));
|
||||
|
||||
// Create StreamingBlob with known size
|
||||
let data_bytes = bytes::Bytes::from(bytes_vec);
|
||||
let stream = stream::once(async move { Ok::<bytes::Bytes, std::io::Error>(data_bytes) });
|
||||
let streaming_blob = s3s::dto::StreamingBlob::wrap(stream);
|
||||
put_builder.set_body(Some(streaming_blob));
|
||||
let put_input = put_builder
|
||||
.build()
|
||||
.map_err(|_| Error::new(ErrorKind::PermanentFileNotAvailable, "Failed to build PutObjectInput"))?;
|
||||
|
||||
match self.storage.put_object(put_input, session_context.credentials()).await {
|
||||
Ok(_output) => {
|
||||
Ok(file_size as u64) // Return the size of the uploaded object
|
||||
}
|
||||
Err(e) => {
|
||||
upload::upload(Arc::clone(&self.storage), session_context, bytes, &bucket, &key)
|
||||
.await
|
||||
.inspect_err(|e| {
|
||||
error!(
|
||||
event = EVENT_FTPS_OBJECT_PUT_FAILED,
|
||||
component = LOG_COMPONENT_PROTOCOLS,
|
||||
subsystem = LOG_SUBSYSTEM_FTPS_DRIVER,
|
||||
username = %masked_username,
|
||||
path = %path_str,
|
||||
bucket = %bucket,
|
||||
object = %key,
|
||||
file_size,
|
||||
error = ?e,
|
||||
error = %e,
|
||||
"ftps object put failed"
|
||||
);
|
||||
Err(Error::new(
|
||||
ErrorKind::PermanentFileNotAvailable,
|
||||
format!("Failed to upload object: {:?}", e),
|
||||
))
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
async fn del<P: AsRef<Path> + Send>(&self, user: &super::server::FtpsUser, path: P) -> Result<()> {
|
||||
@@ -927,3 +892,517 @@ mod tests {
|
||||
assert!(parse_s3_path("/bucket/__XLDIR__").is_err());
|
||||
}
|
||||
}
|
||||
|
||||
mod upload {
|
||||
use crate::common::{
|
||||
client::s3::StorageBackend,
|
||||
gateway::{S3Action, authorize_operation},
|
||||
session::SessionContext,
|
||||
};
|
||||
use bytes::Bytes;
|
||||
use s3s::dto::*;
|
||||
use std::{
|
||||
sync::{Arc, LazyLock},
|
||||
time::Duration,
|
||||
};
|
||||
use tokio::{
|
||||
io::{AsyncRead, AsyncReadExt},
|
||||
sync::Semaphore,
|
||||
};
|
||||
use unftp_core::storage::{Error, ErrorKind, Result};
|
||||
|
||||
// One in-flight part per upload; never buffer the next part while the backend consumes this one.
|
||||
const PART_SIZE: usize = 16 * 1024 * 1024;
|
||||
const MAX_PARTS: i32 = 10_000;
|
||||
const ABORT_TIMEOUT: Duration = Duration::from_secs(30);
|
||||
static ABORT_PERMITS: LazyLock<Arc<Semaphore>> = LazyLock::new(|| Arc::new(Semaphore::new(32)));
|
||||
const EVENT_FTPS_UPLOAD_CLEANUP: &str = "ftps_upload_cleanup";
|
||||
const LOG_COMPONENT_PROTOCOLS: &str = "protocols";
|
||||
const LOG_SUBSYSTEM_FTPS_UPLOAD: &str = "ftps_upload";
|
||||
|
||||
fn upload_error(error: impl std::fmt::Display) -> Error {
|
||||
Error::new(ErrorKind::TransientFileNotAvailable, error.to_string())
|
||||
}
|
||||
|
||||
async fn authorize(session: &SessionContext, action: S3Action, bucket: &str, key: &str) -> Result<()> {
|
||||
authorize_operation(session, &action, bucket, Some(key))
|
||||
.await
|
||||
.map_err(|_| Error::new(ErrorKind::PermanentFileNotAvailable, "Access denied"))
|
||||
}
|
||||
|
||||
fn next_part_number(completed: usize) -> Result<i32> {
|
||||
let completed = i32::try_from(completed).map_err(upload_error)?;
|
||||
if completed >= MAX_PARTS {
|
||||
return Err(Error::new(
|
||||
ErrorKind::ExceededStorageAllocationError,
|
||||
"FTPS multipart upload exceeds 10000 parts",
|
||||
));
|
||||
}
|
||||
Ok(completed + 1)
|
||||
}
|
||||
|
||||
async fn read_part(reader: &mut (impl AsyncRead + Unpin)) -> Result<Vec<u8>> {
|
||||
let mut part = Vec::with_capacity(PART_SIZE);
|
||||
reader
|
||||
.take(u64::try_from(PART_SIZE).map_err(upload_error)?)
|
||||
.read_to_end(&mut part)
|
||||
.await
|
||||
.map_err(upload_error)?;
|
||||
Ok(part)
|
||||
}
|
||||
|
||||
fn body(part: Vec<u8>) -> StreamingBlob {
|
||||
StreamingBlob::from_bytes(Bytes::from(part))
|
||||
}
|
||||
|
||||
// Keep the upload ID alive across all cancellable awaits after CreateMultipartUpload.
|
||||
// A process/runtime crash still requires an AbortIncompleteMultipartUpload lifecycle rule.
|
||||
struct PendingUpload<S: StorageBackend + 'static> {
|
||||
storage: Arc<S>,
|
||||
session: SessionContext,
|
||||
input: Option<AbortMultipartUploadInput>,
|
||||
}
|
||||
|
||||
impl<S: StorageBackend + 'static> PendingUpload<S> {
|
||||
async fn abort(&mut self) {
|
||||
if let Some(input) = self.input.as_ref() {
|
||||
abort_upload(&*self.storage, &self.session, input.clone()).await;
|
||||
self.input = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn abort_upload<S: StorageBackend>(storage: &S, session: &SessionContext, input: AbortMultipartUploadInput) {
|
||||
let cleanup = async {
|
||||
authorize(session, S3Action::AbortMultipartUpload, &input.bucket, &input.key).await?;
|
||||
storage
|
||||
.abort_multipart_upload(input, session.credentials())
|
||||
.await
|
||||
.map_err(upload_error)?;
|
||||
Ok::<_, Error>(())
|
||||
};
|
||||
if !matches!(tokio::time::timeout(ABORT_TIMEOUT, cleanup).await, Ok(Ok(()))) {
|
||||
tracing::warn!(
|
||||
event = EVENT_FTPS_UPLOAD_CLEANUP,
|
||||
component = LOG_COMPONENT_PROTOCOLS,
|
||||
subsystem = LOG_SUBSYSTEM_FTPS_UPLOAD,
|
||||
result = "abort_failed",
|
||||
"FTPS incomplete upload requires lifecycle cleanup"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
impl<S: StorageBackend + 'static> Drop for PendingUpload<S> {
|
||||
fn drop(&mut self) {
|
||||
let Some(input) = self.input.take() else {
|
||||
return;
|
||||
};
|
||||
let (Ok(runtime), Ok(permit)) =
|
||||
(tokio::runtime::Handle::try_current(), Arc::clone(&ABORT_PERMITS).try_acquire_owned())
|
||||
else {
|
||||
tracing::warn!(
|
||||
event = EVENT_FTPS_UPLOAD_CLEANUP,
|
||||
component = LOG_COMPONENT_PROTOCOLS,
|
||||
subsystem = LOG_SUBSYSTEM_FTPS_UPLOAD,
|
||||
result = "abort_unavailable",
|
||||
"FTPS incomplete upload requires lifecycle cleanup"
|
||||
);
|
||||
return;
|
||||
};
|
||||
let storage = Arc::clone(&self.storage);
|
||||
let session = self.session.clone();
|
||||
runtime.spawn(async move {
|
||||
let _permit = permit;
|
||||
abort_upload(&*storage, &session, input).await;
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) async fn upload<S: StorageBackend + 'static>(
|
||||
storage: Arc<S>,
|
||||
session: &SessionContext,
|
||||
mut reader: impl AsyncRead + Unpin,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
) -> Result<u64> {
|
||||
let mut part = read_part(&mut reader).await?;
|
||||
if part.len() < PART_SIZE {
|
||||
let size = u64::try_from(part.len()).map_err(upload_error)?;
|
||||
authorize(session, S3Action::PutObject, bucket, key).await?;
|
||||
let input = PutObjectInput::builder()
|
||||
.bucket(bucket.to_owned())
|
||||
.key(key.to_owned())
|
||||
.content_length(Some(i64::try_from(size).map_err(upload_error)?))
|
||||
.body(Some(body(part)))
|
||||
.build()
|
||||
.map_err(upload_error)?;
|
||||
storage.put_object(input, session.credentials()).await.map_err(upload_error)?;
|
||||
return Ok(size);
|
||||
}
|
||||
|
||||
authorize(session, S3Action::CreateMultipartUpload, bucket, key).await?;
|
||||
let input = CreateMultipartUploadInput::builder()
|
||||
.bucket(bucket.to_owned())
|
||||
.key(key.to_owned())
|
||||
.build()
|
||||
.map_err(upload_error)?;
|
||||
let output = storage
|
||||
.create_multipart_upload(input, session.credentials())
|
||||
.await
|
||||
.map_err(upload_error)?;
|
||||
let upload_id = output
|
||||
.upload_id
|
||||
.filter(|id| !id.is_empty())
|
||||
.ok_or_else(|| upload_error("CreateMultipartUpload returned no upload ID"))?;
|
||||
let abort = AbortMultipartUploadInput::builder()
|
||||
.bucket(bucket.to_owned())
|
||||
.key(key.to_owned())
|
||||
.upload_id(upload_id.clone())
|
||||
.build()
|
||||
.map_err(upload_error)?;
|
||||
let mut pending = PendingUpload {
|
||||
storage: Arc::clone(&storage),
|
||||
session: session.clone(),
|
||||
input: Some(abort),
|
||||
};
|
||||
|
||||
let result = async {
|
||||
let mut parts = Vec::new();
|
||||
let mut total = 0u64;
|
||||
while !part.is_empty() {
|
||||
let number = next_part_number(parts.len())?;
|
||||
let length = i64::try_from(part.len()).map_err(upload_error)?;
|
||||
total = total
|
||||
.checked_add(u64::try_from(part.len()).map_err(upload_error)?)
|
||||
.ok_or_else(|| upload_error("FTPS upload size overflow"))?;
|
||||
authorize(session, S3Action::UploadPart, bucket, key).await?;
|
||||
let input = UploadPartInput::builder()
|
||||
.bucket(bucket.to_owned())
|
||||
.key(key.to_owned())
|
||||
.upload_id(upload_id.clone())
|
||||
.part_number(number)
|
||||
.content_length(Some(length))
|
||||
.body(Some(body(part)))
|
||||
.build()
|
||||
.map_err(upload_error)?;
|
||||
let output = storage
|
||||
.upload_part(input, session.credentials())
|
||||
.await
|
||||
.map_err(upload_error)?;
|
||||
let etag = output.e_tag.ok_or_else(|| upload_error("UploadPart returned no ETag"))?;
|
||||
parts.push(CompletedPart {
|
||||
e_tag: Some(etag),
|
||||
part_number: Some(number),
|
||||
..Default::default()
|
||||
});
|
||||
part = read_part(&mut reader).await?;
|
||||
}
|
||||
authorize(session, S3Action::CompleteMultipartUpload, bucket, key).await?;
|
||||
let input = CompleteMultipartUploadInput::builder()
|
||||
.bucket(bucket.to_owned())
|
||||
.key(key.to_owned())
|
||||
.upload_id(upload_id)
|
||||
.multipart_upload(Some(CompletedMultipartUpload { parts: Some(parts) }))
|
||||
.build()
|
||||
.map_err(upload_error)?;
|
||||
storage
|
||||
.complete_multipart_upload(input, session.credentials())
|
||||
.await
|
||||
.map_err(upload_error)?;
|
||||
// No await between a successful commit and disarming cancellation cleanup.
|
||||
pending.input = None;
|
||||
Ok(total)
|
||||
}
|
||||
.await;
|
||||
if result.is_err() {
|
||||
pending.abort().await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::{
|
||||
common::{
|
||||
dummy_storage::{DummyBackend, DummyError},
|
||||
gateway::{is_operation_supported, with_test_auth_override},
|
||||
session::{Protocol, test_session},
|
||||
},
|
||||
ftps::{driver::FtpsDriver, server::FtpsUser},
|
||||
};
|
||||
use std::{
|
||||
pin::Pin,
|
||||
task::{Context, Poll},
|
||||
};
|
||||
use tokio::{io::ReadBuf, sync::Notify};
|
||||
use unftp_core::storage::StorageBackend as _;
|
||||
|
||||
fn user() -> FtpsUser {
|
||||
FtpsUser {
|
||||
username: "upload-test".into(),
|
||||
name: None,
|
||||
session_context: test_session(Protocol::Ftps),
|
||||
}
|
||||
}
|
||||
|
||||
fn multipart_backend() -> DummyBackend {
|
||||
let backend = DummyBackend::new();
|
||||
backend.queue_create_multipart_upload_ok("upload-1");
|
||||
backend.queue_complete_multipart_upload_ok();
|
||||
backend
|
||||
}
|
||||
|
||||
// Refuse to produce part N+1 before part N has reached the backend. The old
|
||||
// read-to-EOF implementation fails this without allocating a huge fixture.
|
||||
struct PacedReader {
|
||||
backend: DummyBackend,
|
||||
position: usize,
|
||||
length: usize,
|
||||
fail_at_end: bool,
|
||||
}
|
||||
impl AsyncRead for PacedReader {
|
||||
fn poll_read(mut self: Pin<&mut Self>, _: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||
if self.position / PART_SIZE > self.backend.upload_part_calls().len() {
|
||||
return Poll::Ready(Err(std::io::Error::other("read ahead of uploaded part")));
|
||||
}
|
||||
if self.position == self.length && self.fail_at_end {
|
||||
return Poll::Ready(Err(std::io::Error::other("interrupted source")));
|
||||
}
|
||||
let count = buf.remaining().min(self.length - self.position).min(8191);
|
||||
let value = u8::try_from(self.position / PART_SIZE).expect("small test part number");
|
||||
// Do not cross a part boundary in a single read, so each part has a distinct byte.
|
||||
let count = count.min(PART_SIZE - self.position % PART_SIZE);
|
||||
buf.put_slice(&vec![value; count]);
|
||||
self.position += count;
|
||||
Poll::Ready(Ok(()))
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ftps_part_numbers_stop_at_s3_limit() {
|
||||
assert_eq!(next_part_number(0).expect("first part"), 1);
|
||||
assert_eq!(next_part_number(9_999).expect("last part"), 10_000);
|
||||
assert!(next_part_number(10_000).is_err());
|
||||
assert!(next_part_number(usize::MAX).is_err());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_part_buffer_does_not_grow_to_detect_eof() {
|
||||
let mut reader = tokio::io::repeat(0).take(u64::try_from(PART_SIZE).expect("size"));
|
||||
let part = read_part(&mut reader).await.expect("read full part");
|
||||
assert_eq!(part.len(), PART_SIZE);
|
||||
assert_eq!(part.capacity(), PART_SIZE);
|
||||
assert!(read_part(&mut reader).await.expect("EOF").is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_small_and_empty_uploads_preserve_bytes() {
|
||||
for payload in [Vec::new(), b"small upload\x00\xff".to_vec()] {
|
||||
let backend = DummyBackend::new();
|
||||
backend.capture_upload_bodies();
|
||||
let driver = FtpsDriver::new(backend.clone());
|
||||
let result = with_test_auth_override(
|
||||
|_, _, _| true,
|
||||
driver.put(&user(), std::io::Cursor::new(payload.clone()), "/bucket/key", 0),
|
||||
)
|
||||
.await
|
||||
.expect("STOR");
|
||||
assert_eq!(result, u64::try_from(payload.len()).expect("size"));
|
||||
assert_eq!(backend.upload_bodies(), vec![payload]);
|
||||
assert!(backend.create_multipart_calls().is_empty());
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_upload_streams_before_eof_and_preserves_part_bytes() {
|
||||
for length in [PART_SIZE - 1, PART_SIZE, PART_SIZE + 1, 2 * PART_SIZE, 2 * PART_SIZE + 7] {
|
||||
let backend = multipart_backend();
|
||||
backend.capture_upload_bodies();
|
||||
for n in 1..=3 {
|
||||
backend.queue_upload_part_ok(format!("etag-{n}"));
|
||||
}
|
||||
let driver = FtpsDriver::new(backend.clone());
|
||||
let reader = PacedReader {
|
||||
backend: backend.clone(),
|
||||
position: 0,
|
||||
length,
|
||||
fail_at_end: false,
|
||||
};
|
||||
let size = with_test_auth_override(|_, _, _| true, driver.put(&user(), reader, "/bucket/key", 0))
|
||||
.await
|
||||
.expect("bounded STOR");
|
||||
assert_eq!(size, u64::try_from(length).expect("size"));
|
||||
let bodies = backend.upload_bodies();
|
||||
assert_eq!(bodies.iter().map(Vec::len).sum::<usize>(), length);
|
||||
for (index, part) in bodies.iter().enumerate() {
|
||||
assert!(part.len() <= PART_SIZE);
|
||||
assert!(part.iter().all(|b| *b == u8::try_from(index).expect("index")));
|
||||
}
|
||||
if length >= PART_SIZE {
|
||||
assert_eq!(backend.complete_multipart_calls()[0].part_count, length.div_ceil(PART_SIZE));
|
||||
for (index, call) in backend.upload_part_calls().iter().enumerate() {
|
||||
assert_eq!(call.part_number, i32::try_from(index + 1).expect("part"));
|
||||
assert_eq!(call.content_length, Some(i64::try_from(bodies[index].len()).expect("length")));
|
||||
}
|
||||
}
|
||||
assert!(backend.abort_multipart_calls().is_empty());
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_upload_aborts_on_source_part_or_complete_failure() {
|
||||
for failure in ["source", "part", "etag", "complete"] {
|
||||
let backend = DummyBackend::new();
|
||||
backend.queue_create_multipart_upload_ok("upload-1");
|
||||
match failure {
|
||||
"part" => backend.queue_upload_part_err(DummyError::Injected("part failed".into())),
|
||||
"etag" => backend.queue_upload_part_ok_without_etag(),
|
||||
_ => backend.queue_upload_part_ok("etag-1"),
|
||||
}
|
||||
backend.queue_complete_multipart_upload_err(DummyError::Injected("complete failed".into()));
|
||||
let driver = FtpsDriver::new(backend.clone());
|
||||
let reader = PacedReader {
|
||||
backend: backend.clone(),
|
||||
position: 0,
|
||||
length: PART_SIZE,
|
||||
fail_at_end: failure == "source",
|
||||
};
|
||||
let result = with_test_auth_override(|_, _, _| true, driver.put(&user(), reader, "/bucket/key", 0)).await;
|
||||
assert!(result.is_err(), "{failure} must not report a committed object");
|
||||
assert_eq!(backend.abort_multipart_calls().len(), 1, "{failure}");
|
||||
assert_eq!(backend.abort_multipart_calls()[0].upload_id, "upload-1");
|
||||
assert_eq!(backend.complete_multipart_calls().len(), usize::from(failure == "complete"));
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_upload_enforces_authorization_at_each_mutation() {
|
||||
for denied in [
|
||||
S3Action::PutObject,
|
||||
S3Action::CreateMultipartUpload,
|
||||
S3Action::UploadPart,
|
||||
S3Action::CompleteMultipartUpload,
|
||||
] {
|
||||
assert!(is_operation_supported(Protocol::Ftps, &denied));
|
||||
let backend = multipart_backend();
|
||||
backend.queue_upload_part_ok("etag-1");
|
||||
let driver = FtpsDriver::new(backend.clone());
|
||||
let reader = tokio::io::repeat(0).take(u64::try_from(PART_SIZE).expect("size"));
|
||||
let denied_action = denied.clone();
|
||||
let result = with_test_auth_override(
|
||||
move |action, _, _| action != &denied_action,
|
||||
driver.put(&user(), reader, "/bucket/key", 0),
|
||||
)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
match denied {
|
||||
S3Action::PutObject | S3Action::CreateMultipartUpload => assert!(backend.create_multipart_calls().is_empty()),
|
||||
S3Action::UploadPart => assert!(backend.upload_part_calls().is_empty()),
|
||||
_ => assert!(backend.complete_multipart_calls().is_empty()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_upload_cleanup_does_not_bypass_abort_permission() {
|
||||
let backend = multipart_backend();
|
||||
backend.queue_upload_part_err(DummyError::Injected("part failed".into()));
|
||||
let driver = FtpsDriver::new(backend.clone());
|
||||
let reader = tokio::io::repeat(0).take(u64::try_from(PART_SIZE).expect("size"));
|
||||
let result = with_test_auth_override(
|
||||
|action, _, _| action != &S3Action::AbortMultipartUpload,
|
||||
driver.put(&user(), reader, "/bucket/key", 0),
|
||||
)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
assert!(backend.abort_multipart_calls().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ftps_cancelled_upload_aborts_pending_parts() {
|
||||
let backend = multipart_backend();
|
||||
let entered = Arc::new(Notify::new());
|
||||
backend.stall_upload_part(Arc::clone(&entered));
|
||||
let driver = FtpsDriver::new(backend.clone());
|
||||
let user = user();
|
||||
with_test_auth_override(|_, _, _| true, async {
|
||||
let reader = tokio::io::repeat(0).take(u64::try_from(PART_SIZE * 2).expect("size"));
|
||||
let mut put = Box::pin(driver.put(&user, reader, "/bucket/key", 0));
|
||||
tokio::select! {
|
||||
_ = entered.notified() => {},
|
||||
result = &mut put => panic!("upload should stall: {result:?}"),
|
||||
}
|
||||
drop(put);
|
||||
tokio::time::timeout(Duration::from_secs(5), async {
|
||||
while backend.abort_multipart_calls().is_empty() {
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("cancellation cleanup");
|
||||
})
|
||||
.await;
|
||||
assert_eq!(backend.abort_multipart_calls().len(), 1);
|
||||
assert!(backend.complete_multipart_calls().is_empty());
|
||||
}
|
||||
// Yield on each read so the memory probe keeps all upload buffers live at once.
|
||||
struct YieldingReader {
|
||||
remaining: u64,
|
||||
yielded: bool,
|
||||
}
|
||||
|
||||
impl AsyncRead for YieldingReader {
|
||||
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||
if self.remaining == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
if !self.yielded {
|
||||
self.yielded = true;
|
||||
cx.waker().wake_by_ref();
|
||||
return Poll::Pending;
|
||||
}
|
||||
self.yielded = false;
|
||||
let count = buf
|
||||
.remaining()
|
||||
.min(1024 * 1024)
|
||||
.min(usize::try_from(self.remaining).unwrap_or(usize::MAX));
|
||||
buf.initialize_unfilled_to(count).fill(0x5a);
|
||||
buf.advance(count);
|
||||
self.remaining -= u64::try_from(count).expect("read size");
|
||||
Poll::Ready(Ok(()))
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[ignore = "manual synthetic concurrent-upload memory measurement"]
|
||||
async fn ftps_large_upload_memory_probe() {
|
||||
let bytes: u64 = std::env::var("RUSTFS_FTPS_TEST_BYTES")
|
||||
.unwrap_or_else(|_| "67108864".into())
|
||||
.parse()
|
||||
.expect("byte count");
|
||||
let concurrency: usize = std::env::var("RUSTFS_FTPS_TEST_CONCURRENCY")
|
||||
.unwrap_or_else(|_| "10".into())
|
||||
.parse()
|
||||
.expect("concurrency");
|
||||
assert!(bytes > 0 && concurrency > 0);
|
||||
let operations = (0..concurrency).map(|_| async {
|
||||
let backend = multipart_backend();
|
||||
for n in 0..bytes.div_ceil(u64::try_from(PART_SIZE).expect("part size")) {
|
||||
backend.queue_upload_part_ok(format!("part-{n}"));
|
||||
}
|
||||
let driver = FtpsDriver::new(backend);
|
||||
let reader = YieldingReader {
|
||||
remaining: bytes,
|
||||
yielded: false,
|
||||
};
|
||||
driver.put(&user(), reader, "/bucket/key", 0).await
|
||||
});
|
||||
let results = with_test_auth_override(|_, _, _| true, futures_util::future::join_all(operations)).await;
|
||||
for result in results {
|
||||
assert_eq!(result.expect("upload"), bytes);
|
||||
}
|
||||
println!("synthetic_upload_bytes={bytes} concurrent_uploads={concurrency}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
# FTPS uploads
|
||||
|
||||
FTPS `STOR` buffers at most one 16 MiB payload chunk per active upload in the
|
||||
protocol driver. It waits for the storage backend to consume each chunk before
|
||||
reading the next. Backend, TLS, connection, and allocator overhead are additional;
|
||||
this is not a process-wide memory limit. Concurrent uploads each have their own
|
||||
buffer.
|
||||
|
||||
Files smaller than 16 MiB, including empty files, use `PutObject`. Files at or
|
||||
above that threshold use sequential S3 multipart uploads, so they can exceed the
|
||||
5 GiB single-`PutObject` limit. Multipart ETags differ from single-PUT ETags and
|
||||
must not be interpreted as a whole-file MD5 checksum. An upload is successful only
|
||||
after the multipart completion succeeds. The fixed part size and S3's 10,000-part
|
||||
limit allow up to 156.25 GiB per FTPS upload; a larger input fails and cleanup is
|
||||
attempted. Resuming or appending with a nonzero offset remains unsupported.
|
||||
|
||||
Each write operation requires `s3:PutObject`. Failed or cancelled transfers also
|
||||
attempt `s3:AbortMultipartUpload` using the authenticated user's permissions;
|
||||
cleanup does not bypass IAM. Grant that permission on the upload prefix to allow
|
||||
immediate cleanup. Cancellation cleanup has a bounded task count and timeout.
|
||||
Configure an `AbortIncompleteMultipartUpload` bucket lifecycle rule as a fallback
|
||||
for denied/failed cleanup, process crashes, or losing the upload ID while upload
|
||||
initiation is in flight. Before completion, received parts do not replace an existing completed object.
|
||||
A lost or failed completion response may have an ambiguous outcome; clients
|
||||
should verify the destination before retrying.
|
||||
@@ -0,0 +1,71 @@
|
||||
# FTPS upload memory regression
|
||||
|
||||
The old `FtpsDriver::put` copied the entire input into a `Vec` before calling
|
||||
`PutObject`. Wrapping that completed allocation in `StreamingBlob` did not make
|
||||
input consumption streaming. It also sent large files through the 5 GiB
|
||||
single-PUT path.
|
||||
|
||||
## Automated regression
|
||||
|
||||
```sh
|
||||
cargo test -p rustfs-protocols --no-default-features --features ftps --lib
|
||||
```
|
||||
|
||||
`ftps_upload_streams_before_eof_and_preserves_part_bytes` uses a reader that
|
||||
refuses to supply the next part until the previous part reaches the backend.
|
||||
Restoring the driver from commit `9c30cc88513c5e1b5443b6ee86bf4bc7965e79e9`
|
||||
while retaining the new test harness makes this test fail with
|
||||
`read ahead of uploaded part`. The multipart implementation passes. Tests also
|
||||
check exact bytes, empty files, part boundaries, part numbering, errors,
|
||||
permissions, cancellation cleanup, and buffer capacity.
|
||||
|
||||
## Manual memory probe
|
||||
|
||||
The ignored test `ftps_large_upload_memory_probe` creates concurrent synthetic
|
||||
readers. Each reader yields between reads so all uploads are active together.
|
||||
It calls the real FTPS storage driver with a scripted storage backend that
|
||||
returns successful S3 responses and discards the payload. It neither stores
|
||||
files nor opens network connections. No large fixture is required.
|
||||
|
||||
Build the test executable first; measuring Cargo would include compiler memory:
|
||||
|
||||
```sh
|
||||
export CARGO_PROFILE_DEV_DEBUG=0 CARGO_PROFILE_TEST_DEBUG=0 CARGO_INCREMENTAL=0
|
||||
cargo test -p rustfs-protocols --no-default-features --features ftps --lib --no-run
|
||||
```
|
||||
|
||||
Use the executable path printed by that command as `test_binary`, then on macOS:
|
||||
|
||||
```sh
|
||||
/usr/bin/time -l "$test_binary" --ignored --exact \
|
||||
ftps::driver::upload::tests::ftps_large_upload_memory_probe --nocapture
|
||||
|
||||
RUSTFS_FTPS_TEST_BYTES=10737418240 RUSTFS_FTPS_TEST_CONCURRENCY=10 \
|
||||
/usr/bin/time -l "$test_binary" --ignored --exact \
|
||||
ftps::driver::upload::tests::ftps_large_upload_memory_probe --nocapture
|
||||
```
|
||||
|
||||
On Linux use `/usr/bin/time -v` (its maximum RSS is reported in KiB). The default
|
||||
probe uses 64 MiB per upload and 10 uploads. These environment variables are
|
||||
**test-only**, not server configuration.
|
||||
|
||||
Observed on macOS arm64 with Rust 1.98.1, debug information disabled, one run per
|
||||
case (maximum resident set size reported in bytes):
|
||||
|
||||
| Driver | Bytes per upload | Concurrent uploads | Maximum RSS |
|
||||
| --- | ---: | ---: | ---: |
|
||||
| Original driver at `9c30cc8`, same probe | 67,108,864 | 10 | 710,033,408 (677.1 MiB) |
|
||||
| Bounded multipart driver | 67,108,864 | 10 | 179,781,632 (171.5 MiB) |
|
||||
| Bounded multipart driver | 10,737,418,240 | 10 | 187,465,728 (178.8 MiB) |
|
||||
|
||||
For the original-driver comparison, only `ftps/driver.rs` was replaced with its
|
||||
base-commit version in the test checkout; the identical probe and dummy backend
|
||||
were retained, the executable rebuilt, and the 64 MiB case rerun. The original
|
||||
driver was not subjected to the 100 GiB aggregate case.
|
||||
|
||||
These measurements isolate protocol-driver buffering, not total server RSS,
|
||||
persisted-data correctness, TLS behavior, or storage throughput. The large case
|
||||
supplies 100 GiB of generated bytes without retaining them. A real server also
|
||||
uses memory for TLS, storage, caches, and allocator overhead. Memory still scales
|
||||
with active upload count. See [FTPS operations](../operations/ftps.md) for limits
|
||||
and cleanup behavior.
|
||||
Reference in New Issue
Block a user