mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 16:28:15 +00:00
Signed-off-by: giter <giter@users.noreply.github.com> Signed-off-by: 安正超 <anzhengchao@gmail.com> Co-authored-by: giter <giter@users.noreply.github.com> Co-authored-by: 安正超 <anzhengchao@gmail.com> Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Co-authored-by: Lijiajie <lijiajie@ffcode.net> Co-authored-by: cxymds <Cxymds@qq.com> Co-authored-by: loverustfs <hello@rustfs.com> Co-authored-by: 唐小鸭 <tangtang1251@qq.com> Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: houseme <4829346+houseme@users.noreply.github.com>
This commit is contained in:
@@ -55,6 +55,7 @@ RUSTFS_BUILD_FEATURES=webdav cargo test --package e2e_test test_protocol_core_su
|
||||
- GET (download file)
|
||||
- PROPFIND on bucket (list objects)
|
||||
- DELETE file
|
||||
- MOVE file (rename object)
|
||||
- DELETE bucket
|
||||
- Authentication failure test
|
||||
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
use crate::common::init_logging;
|
||||
use crate::protocols::ftps_core::test_ftps_core_operations;
|
||||
use crate::protocols::webdav_core::test_webdav_core_operations;
|
||||
use serial_test::serial;
|
||||
use std::time::Instant;
|
||||
use tokio::time::{Duration, sleep};
|
||||
use tracing::{error, info};
|
||||
@@ -160,6 +161,7 @@ impl ProtocolTestSuite {
|
||||
|
||||
/// Test suite
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_protocol_core_suite() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let suite = ProtocolTestSuite::new();
|
||||
let results = suite.run_test_suite().await;
|
||||
|
||||
@@ -14,17 +14,24 @@
|
||||
|
||||
//! Core WebDAV tests
|
||||
|
||||
use crate::common::local_http_client;
|
||||
use crate::common::rustfs_binary_path_with_features;
|
||||
use crate::protocols::test_env::{DEFAULT_ACCESS_KEY, DEFAULT_SECRET_KEY, ProtocolTestEnvironment};
|
||||
use anyhow::Result;
|
||||
use base64::Engine;
|
||||
use http::header::{CONTENT_TYPE, HOST};
|
||||
use reqwest::Client;
|
||||
use rustfs_signer::constants::UNSIGNED_PAYLOAD;
|
||||
use rustfs_signer::sign_v4;
|
||||
use s3s::Body;
|
||||
use serial_test::serial;
|
||||
use tokio::process::Command;
|
||||
use tracing::info;
|
||||
|
||||
// Fixed WebDAV port for testing
|
||||
const WEBDAV_PORT: u16 = 9080;
|
||||
const WEBDAV_ADDRESS: &str = "127.0.0.1:9080";
|
||||
const S3_TEST_ADDRESS: &str = "127.0.0.1:9010";
|
||||
|
||||
/// Create HTTP client with basic auth
|
||||
fn create_client() -> Client {
|
||||
@@ -36,19 +43,114 @@ fn create_client() -> Client {
|
||||
|
||||
/// Get basic auth header value
|
||||
fn basic_auth_header() -> String {
|
||||
let credentials = format!("{}:{}", DEFAULT_ACCESS_KEY, DEFAULT_SECRET_KEY);
|
||||
basic_auth_header_for(DEFAULT_ACCESS_KEY, DEFAULT_SECRET_KEY)
|
||||
}
|
||||
|
||||
fn basic_auth_header_for(access_key: &str, secret_key: &str) -> String {
|
||||
let credentials = format!("{}:{}", access_key, secret_key);
|
||||
let encoded = base64::engine::general_purpose::STANDARD.encode(credentials);
|
||||
format!("Basic {}", encoded)
|
||||
}
|
||||
|
||||
async fn signed_admin_request(
|
||||
method: http::Method,
|
||||
url: &str,
|
||||
body: Option<Vec<u8>>,
|
||||
content_type: Option<&str>,
|
||||
) -> Result<reqwest::Response> {
|
||||
let uri = url.parse::<http::Uri>()?;
|
||||
let authority = uri
|
||||
.authority()
|
||||
.ok_or_else(|| anyhow::anyhow!("request URL missing authority"))?
|
||||
.to_string();
|
||||
let mut request = http::Request::builder().method(method.clone()).uri(uri);
|
||||
request = request.header(HOST, authority);
|
||||
request = request.header("x-amz-content-sha256", UNSIGNED_PAYLOAD);
|
||||
if let Some(content_type) = content_type {
|
||||
request = request.header(CONTENT_TYPE, content_type);
|
||||
}
|
||||
|
||||
let content_len = body.as_ref().map(|body| body.len() as i64).unwrap_or_default();
|
||||
let signed = sign_v4(
|
||||
request.body(Body::empty())?,
|
||||
content_len,
|
||||
DEFAULT_ACCESS_KEY,
|
||||
DEFAULT_SECRET_KEY,
|
||||
"",
|
||||
"us-east-1",
|
||||
);
|
||||
|
||||
let reqwest_method = reqwest::Method::from_bytes(method.as_str().as_bytes())?;
|
||||
let mut request_builder = local_http_client().request(reqwest_method, url);
|
||||
for (name, value) in signed.headers() {
|
||||
request_builder = request_builder.header(name, value);
|
||||
}
|
||||
if let Some(body) = body {
|
||||
request_builder = request_builder.body(body);
|
||||
}
|
||||
|
||||
Ok(request_builder.send().await?)
|
||||
}
|
||||
|
||||
async fn admin_create_user(base_url: &str, username: &str, secret_key: &str) -> Result<()> {
|
||||
let url = format!("{}/rustfs/admin/v3/add-user?accessKey={}", base_url, username);
|
||||
let body = serde_json::json!({
|
||||
"secretKey": secret_key,
|
||||
"status": "enabled"
|
||||
});
|
||||
let response =
|
||||
signed_admin_request(http::Method::PUT, &url, Some(body.to_string().into_bytes()), Some("application/json")).await?;
|
||||
|
||||
if response.status() != reqwest::StatusCode::OK {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
anyhow::bail!("create user failed: {status} {body}");
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn admin_add_canned_policy(base_url: &str, policy_name: &str, policy: &serde_json::Value) -> Result<()> {
|
||||
let url = format!("{}/rustfs/admin/v3/add-canned-policy?name={}", base_url, policy_name);
|
||||
let response =
|
||||
signed_admin_request(http::Method::PUT, &url, Some(policy.to_string().into_bytes()), Some("application/json")).await?;
|
||||
|
||||
if response.status() != reqwest::StatusCode::OK {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
anyhow::bail!("add canned policy failed: {status} {body}");
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn admin_attach_policy_to_user(base_url: &str, policy_name: &str, username: &str) -> Result<()> {
|
||||
let url = format!(
|
||||
"{}/rustfs/admin/v3/set-user-or-group-policy?policyName={}&userOrGroup={}&isGroup=false",
|
||||
base_url, policy_name, username
|
||||
);
|
||||
let response = signed_admin_request(http::Method::PUT, &url, Some(Vec::new()), None).await?;
|
||||
|
||||
if response.status() != reqwest::StatusCode::OK {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
anyhow::bail!("attach policy failed: {status} {body}");
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Test WebDAV: MKCOL (create bucket), PUT, GET, DELETE, PROPFIND operations
|
||||
pub async fn test_webdav_core_operations() -> Result<()> {
|
||||
let env = ProtocolTestEnvironment::new().map_err(|e| anyhow::anyhow!("{}", e))?;
|
||||
let admin_base_url = format!("http://{}", S3_TEST_ADDRESS);
|
||||
|
||||
// Start server manually
|
||||
info!("Starting WebDAV server on {}", WEBDAV_ADDRESS);
|
||||
let binary_path = rustfs_binary_path_with_features(Some("ftps,webdav"));
|
||||
let binary_path = rustfs_binary_path_with_features(Some("webdav"));
|
||||
let mut server_process = Command::new(&binary_path)
|
||||
.arg("--address")
|
||||
.arg(S3_TEST_ADDRESS)
|
||||
.env("RUSTFS_WEBDAV_ENABLE", "true")
|
||||
.env("RUSTFS_WEBDAV_ADDRESS", WEBDAV_ADDRESS)
|
||||
.env("RUSTFS_WEBDAV_TLS_ENABLED", "false") // No TLS for testing
|
||||
@@ -170,6 +272,372 @@ pub async fn test_webdav_core_operations() -> Result<()> {
|
||||
);
|
||||
info!("PASS: Verified file '{}' is deleted", filename);
|
||||
|
||||
// Test MOVE (rename) file
|
||||
info!("Testing WebDAV: PUT file for MOVE test");
|
||||
let move_filename = "move-source.txt";
|
||||
let move_dest_filename = "move-dest.txt";
|
||||
let move_content = "File to be moved!";
|
||||
let resp = client
|
||||
.put(format!("{}/{}/{}", base_url, bucket_name, move_filename))
|
||||
.header("Authorization", &auth_header)
|
||||
.body(move_content)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"PUT for MOVE test should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: PUT file '{}' for MOVE test successful", move_filename);
|
||||
|
||||
// Execute MOVE request
|
||||
info!("Testing WebDAV: MOVE file '{}' to '{}'", move_filename, move_dest_filename);
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MOVE").unwrap(),
|
||||
format!("{}/{}/{}", base_url, bucket_name, move_filename),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.header("Destination", format!("/{}/{}", bucket_name, move_dest_filename))
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 204 || resp.status().as_u16() == 201,
|
||||
"MOVE should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!(
|
||||
"PASS: MOVE file '{}' to '{}' successful (HTTP {})",
|
||||
move_filename,
|
||||
move_dest_filename,
|
||||
resp.status()
|
||||
);
|
||||
|
||||
// Verify source file is gone
|
||||
info!("Testing WebDAV: Verify source '{}' is deleted after MOVE", move_filename);
|
||||
let resp = client
|
||||
.get(format!("{}/{}/{}", base_url, bucket_name, move_filename))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().as_u16() == 404,
|
||||
"GET moved source should return 404, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: Source '{}' is deleted after MOVE", move_filename);
|
||||
|
||||
// Verify destination file exists and content matches
|
||||
info!("Testing WebDAV: Verify destination '{}' has correct content", move_dest_filename);
|
||||
let resp = client
|
||||
.get(format!("{}/{}/{}", base_url, bucket_name, move_dest_filename))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success(),
|
||||
"GET destination after MOVE should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
let moved_content = resp.text().await?;
|
||||
assert_eq!(moved_content, move_content, "Moved file content should match original");
|
||||
info!("PASS: Destination '{}' has correct content after MOVE", move_dest_filename);
|
||||
|
||||
// Test directory creation and rename
|
||||
info!("Testing WebDAV: MKCOL directory");
|
||||
let dir_name = "test-directory";
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MKCOL").unwrap(),
|
||||
format!("{}/{}/{}", base_url, bucket_name, dir_name),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"MKCOL directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: MKCOL directory successful");
|
||||
|
||||
// Upload file into directory
|
||||
let dir_filename = "dir-file.txt";
|
||||
let dir_file_content = "File inside test directory!";
|
||||
let resp = client
|
||||
.put(format!("{}/{}/{}/{}", base_url, bucket_name, dir_name, dir_filename))
|
||||
.header("Authorization", &auth_header)
|
||||
.body(dir_file_content)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"PUT file into directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: PUT file into directory successful");
|
||||
|
||||
// Test PROPFIND on directory
|
||||
info!("Testing WebDAV: PROPFIND directory");
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"PROPFIND").unwrap(),
|
||||
format!("{}/{}/{}", base_url, bucket_name, dir_name),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.header("Depth", "1")
|
||||
.send()
|
||||
.await?;
|
||||
assert!(resp.status().is_success(), "PROPFIND directory should succeed, got: {}", resp.status());
|
||||
let propfind_body = resp.text().await?;
|
||||
assert!(propfind_body.contains(dir_filename), "PROPFIND should list file in directory");
|
||||
info!("PASS: PROPFIND directory successful, file listed correctly");
|
||||
|
||||
// Current WebDAV support exposes collection listing via PROPFIND; GET on a collection is not implemented.
|
||||
info!("Testing WebDAV: GET directory listing fallback behavior");
|
||||
let resp = client
|
||||
.get(format!("{}/{}/{}", base_url, bucket_name, dir_name))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(
|
||||
resp.status().as_u16(),
|
||||
405,
|
||||
"GET on a WebDAV collection should currently return 405, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: GET collection correctly returned 405; PROPFIND remains the listing path");
|
||||
|
||||
// Rename directory
|
||||
info!("Testing WebDAV: MOVE directory");
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MOVE").unwrap(),
|
||||
format!("{}/{}/{}", base_url, bucket_name, dir_name),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.header("Destination", format!("/{}/renamed-dir", bucket_name))
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 204 || resp.status().as_u16() == 201,
|
||||
"MOVE directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: MOVE directory successful");
|
||||
|
||||
// Verify source directory is gone
|
||||
info!("Testing WebDAV: Verify source directory is deleted after MOVE");
|
||||
let resp = client
|
||||
.get(format!("{}/{}/{}", base_url, bucket_name, dir_name))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().as_u16() == 404,
|
||||
"GET moved directory should return 404, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: Source directory is deleted after MOVE");
|
||||
|
||||
// Verify renamed directory exists and file content matches
|
||||
info!("Testing WebDAV: Verify file in renamed directory");
|
||||
let resp = client
|
||||
.get(format!("{}/{}/renamed-dir/{}", base_url, bucket_name, dir_filename))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success(),
|
||||
"GET file in renamed directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
let renamed_dir_content = resp.text().await?;
|
||||
assert_eq!(renamed_dir_content, dir_file_content, "File content in renamed directory should match");
|
||||
info!("PASS: File in renamed directory has correct content");
|
||||
|
||||
// Test nested directory creation and rename
|
||||
info!("Testing WebDAV: MKCOL nested directory");
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MKCOL").unwrap(),
|
||||
format!("{}/{}/renamed-dir/nested-dir", base_url, bucket_name),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"MKCOL nested directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: MKCOL nested directory successful");
|
||||
|
||||
// Upload file into nested directory
|
||||
let nested_file_content = "File in nested directory!";
|
||||
let resp = client
|
||||
.put(format!("{}/{}/renamed-dir/nested-dir/nested-file.txt", base_url, bucket_name))
|
||||
.header("Authorization", &auth_header)
|
||||
.body(nested_file_content)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"PUT file into nested directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
// Rename nested directory
|
||||
info!("Testing WebDAV: MOVE nested directory");
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MOVE").unwrap(),
|
||||
format!("{}/{}/renamed-dir/nested-dir", base_url, bucket_name),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.header("Destination", format!("/{}/renamed-dir/new-nested-dir", bucket_name))
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 204 || resp.status().as_u16() == 201,
|
||||
"MOVE nested directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
info!("PASS: MOVE nested directory successful");
|
||||
|
||||
// Verify nested file after rename
|
||||
let resp = client
|
||||
.get(format!("{}/{}/renamed-dir/new-nested-dir/nested-file.txt", base_url, bucket_name))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success(),
|
||||
"GET file in renamed nested directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
let nested_content = resp.text().await?;
|
||||
assert_eq!(
|
||||
nested_content, nested_file_content,
|
||||
"File content in renamed nested directory should match"
|
||||
);
|
||||
info!("PASS: File in renamed nested directory has correct content");
|
||||
|
||||
// Test directory MOVE authz failure does not create partial destination writes
|
||||
info!("Testing WebDAV: directory MOVE denied by missing DeleteObject must not create partial writes");
|
||||
let restricted_bucket = "webdav-authz-bucket";
|
||||
let restricted_dir = "restricted-src";
|
||||
let restricted_dst = "restricted-dst";
|
||||
let restricted_file = "locked.txt";
|
||||
let restricted_content = "must remain only at source";
|
||||
let restricted_user = "webdav-limited-user";
|
||||
let restricted_secret = "webdav-limited-secret";
|
||||
let restricted_policy_name = "webdav-move-no-delete";
|
||||
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MKCOL").unwrap(),
|
||||
format!("{}/{}", base_url, restricted_bucket),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"MKCOL restricted bucket should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MKCOL").unwrap(),
|
||||
format!("{}/{}/{}", base_url, restricted_bucket, restricted_dir),
|
||||
)
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"MKCOL restricted directory should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
let resp = client
|
||||
.put(format!("{}/{}/{}/{}", base_url, restricted_bucket, restricted_dir, restricted_file))
|
||||
.header("Authorization", &auth_header)
|
||||
.body(restricted_content)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success() || resp.status().as_u16() == 201,
|
||||
"PUT restricted file should succeed, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
admin_create_user(&admin_base_url, restricted_user, restricted_secret).await?;
|
||||
admin_add_canned_policy(
|
||||
&admin_base_url,
|
||||
restricted_policy_name,
|
||||
&serde_json::json!({
|
||||
"Version": "2012-10-17",
|
||||
"Statement": [
|
||||
{
|
||||
"Effect": "Allow",
|
||||
"Action": ["s3:ListBucket"],
|
||||
"Resource": [format!("arn:aws:s3:::{}", restricted_bucket)]
|
||||
},
|
||||
{
|
||||
"Effect": "Allow",
|
||||
"Action": ["s3:GetObject", "s3:PutObject"],
|
||||
"Resource": [format!("arn:aws:s3:::{}/*", restricted_bucket)]
|
||||
}
|
||||
]
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
admin_attach_policy_to_user(&admin_base_url, restricted_policy_name, restricted_user).await?;
|
||||
|
||||
let restricted_auth = basic_auth_header_for(restricted_user, restricted_secret);
|
||||
let resp = client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MOVE").unwrap(),
|
||||
format!("{}/{}/{}", base_url, restricted_bucket, restricted_dir),
|
||||
)
|
||||
.header("Authorization", &restricted_auth)
|
||||
.header("Destination", format!("/{}/{}", restricted_bucket, restricted_dst))
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
!resp.status().is_success(),
|
||||
"MOVE without DeleteObject permission should be rejected, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
let resp = client
|
||||
.get(format!("{}/{}/{}/{}", base_url, restricted_bucket, restricted_dst, restricted_file))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(
|
||||
resp.status().as_u16(),
|
||||
404,
|
||||
"Denied MOVE must not create destination object, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
let resp = client
|
||||
.get(format!("{}/{}/{}/{}", base_url, restricted_bucket, restricted_dir, restricted_file))
|
||||
.header("Authorization", &auth_header)
|
||||
.send()
|
||||
.await?;
|
||||
assert!(
|
||||
resp.status().is_success(),
|
||||
"Source object should remain after denied MOVE, got: {}",
|
||||
resp.status()
|
||||
);
|
||||
assert_eq!(resp.text().await?, restricted_content, "Denied MOVE must leave source content untouched");
|
||||
info!("PASS: denied directory MOVE left source intact and created no destination objects");
|
||||
|
||||
// Test DELETE bucket
|
||||
info!("Testing WebDAV: DELETE bucket '{}'", bucket_name);
|
||||
let resp = client
|
||||
@@ -205,3 +673,9 @@ pub async fn test_webdav_core_operations() -> Result<()> {
|
||||
|
||||
result
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_webdav_core_operations_direct() -> Result<()> {
|
||||
test_webdav_core_operations().await
|
||||
}
|
||||
|
||||
@@ -397,6 +397,20 @@ where
|
||||
session_context: Arc<SessionContext>,
|
||||
}
|
||||
|
||||
enum ResolvedPath {
|
||||
File(Box<HeadObjectOutput>),
|
||||
Directory {
|
||||
prefix: String,
|
||||
metadata: Option<Box<HeadObjectOutput>>,
|
||||
},
|
||||
}
|
||||
|
||||
enum HeadObjectProbe {
|
||||
Forbidden,
|
||||
Missing,
|
||||
Found(Box<HeadObjectOutput>),
|
||||
}
|
||||
|
||||
impl<S> Debug for WebDavDriver<S>
|
||||
where
|
||||
S: S3StorageBackend + Debug + Clone + Send + Sync + 'static,
|
||||
@@ -430,6 +444,210 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
fn credentials(&self) -> (&str, &str) {
|
||||
(
|
||||
&self.session_context.principal.user_identity.credentials.access_key,
|
||||
&self.session_context.principal.user_identity.credentials.secret_key,
|
||||
)
|
||||
}
|
||||
|
||||
fn is_missing_head_object_error(error: &str) -> bool {
|
||||
let lower = error.to_ascii_lowercase();
|
||||
lower.contains("nosuchkey")
|
||||
|| lower.contains("notfound")
|
||||
|| lower.contains("not found")
|
||||
|| lower.contains("status code: 404")
|
||||
}
|
||||
|
||||
async fn prefix_has_entries(&self, bucket: &str, prefix: &str) -> FsResult<bool> {
|
||||
let (access_key, secret_key) = self.credentials();
|
||||
let list_input = ListObjectsV2Input::builder()
|
||||
.bucket(bucket.to_string())
|
||||
.prefix(Some(prefix.to_string()))
|
||||
.max_keys(Some(1))
|
||||
.build()
|
||||
.map_err(|_| FsError::GeneralFailure)?;
|
||||
|
||||
let output = self
|
||||
.storage
|
||||
.list_objects_v2(list_input, access_key, secret_key)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Failed to list objects in {} with prefix '{}': {}", bucket, prefix, e);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
|
||||
Ok(output.contents.map(|c| !c.is_empty()).unwrap_or(false)
|
||||
|| output.common_prefixes.map(|c| !c.is_empty()).unwrap_or(false))
|
||||
}
|
||||
|
||||
async fn copy_object_streaming(&self, src_bucket: &str, src_key: &str, dst_bucket: &str, dst_key: &str) -> FsResult<()> {
|
||||
let (access_key, secret_key) = self.credentials();
|
||||
let get_output = self
|
||||
.storage
|
||||
.get_object(src_bucket, src_key, access_key, secret_key, None)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Failed to get source object '{}' in '{}': {}", src_key, src_bucket, e);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
|
||||
let GetObjectOutput {
|
||||
body,
|
||||
content_length,
|
||||
content_type,
|
||||
..
|
||||
} = get_output;
|
||||
let body = body.ok_or_else(|| {
|
||||
error!("GetObject for source object '{}/{}' returned no body stream", src_bucket, src_key);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
|
||||
let mut put_builder = PutObjectInput::builder()
|
||||
.bucket(dst_bucket.to_string())
|
||||
.key(dst_key.to_string())
|
||||
.body(Some(body));
|
||||
|
||||
if let Some(content_length) = content_length {
|
||||
put_builder = put_builder.content_length(Some(content_length));
|
||||
}
|
||||
|
||||
if let Some(content_type) = content_type {
|
||||
put_builder = put_builder.content_type(Some(content_type));
|
||||
}
|
||||
|
||||
let put_input = put_builder.build().map_err(|_| FsError::GeneralFailure)?;
|
||||
|
||||
self.storage
|
||||
.put_object(put_input, access_key, secret_key)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
"Failed to copy object from '{}/{}' to '{}/{}': {}",
|
||||
src_bucket, src_key, dst_bucket, dst_key, e
|
||||
);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn execute_directory_rename_pairs(
|
||||
&self,
|
||||
src_bucket: &str,
|
||||
dst_bucket: &str,
|
||||
rename_pairs: &[(String, String)],
|
||||
) -> FsResult<()> {
|
||||
let (access_key, secret_key) = self.credentials();
|
||||
|
||||
for (src_obj_key, dst_obj_key) in rename_pairs {
|
||||
self.copy_object_streaming(src_bucket, src_obj_key, dst_bucket, dst_obj_key)
|
||||
.await?;
|
||||
}
|
||||
|
||||
for (src_obj_key, _) in rename_pairs {
|
||||
self.storage
|
||||
.delete_object(src_bucket, src_obj_key, access_key, secret_key)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Failed to delete source object '{}' after directory rename: {}", src_obj_key, e);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn probe_head_object(&self, bucket: &str, key: &str) -> FsResult<HeadObjectProbe> {
|
||||
let (access_key, secret_key) = self.credentials();
|
||||
|
||||
if authorize_operation(&self.session_context, &S3Action::HeadObject, bucket, Some(key))
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
return Ok(HeadObjectProbe::Forbidden);
|
||||
}
|
||||
|
||||
match self.storage.head_object(bucket, key, access_key, secret_key).await {
|
||||
Ok(output) => Ok(HeadObjectProbe::Found(Box::new(output))),
|
||||
Err(e) => {
|
||||
let err_msg = e.to_string();
|
||||
if Self::is_missing_head_object_error(&err_msg) {
|
||||
Ok(HeadObjectProbe::Missing)
|
||||
} else {
|
||||
error!("Failed to probe object '{}/{}': {}", bucket, key, err_msg);
|
||||
Err(FsError::GeneralFailure)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn resolve_path(&self, bucket: &str, key: &str) -> FsResult<ResolvedPath> {
|
||||
let prefix = format!("{}/", key);
|
||||
let mut had_visibility = false;
|
||||
|
||||
match self.probe_head_object(bucket, key).await? {
|
||||
HeadObjectProbe::Found(output) => {
|
||||
let size = output.content_length.unwrap_or(0) as u64;
|
||||
let is_dir_marker = output.content_type.as_deref() == Some("application/x-directory");
|
||||
|
||||
if is_dir_marker {
|
||||
return Ok(ResolvedPath::Directory {
|
||||
prefix,
|
||||
metadata: Some(output),
|
||||
});
|
||||
}
|
||||
|
||||
if size == 0
|
||||
&& authorize_operation(&self.session_context, &S3Action::ListBucket, bucket, Some(&prefix))
|
||||
.await
|
||||
.is_ok()
|
||||
&& self.prefix_has_entries(bucket, &prefix).await?
|
||||
{
|
||||
return Ok(ResolvedPath::Directory {
|
||||
prefix,
|
||||
metadata: Some(output),
|
||||
});
|
||||
}
|
||||
|
||||
return Ok(ResolvedPath::File(output));
|
||||
}
|
||||
HeadObjectProbe::Missing => {
|
||||
had_visibility = true;
|
||||
}
|
||||
HeadObjectProbe::Forbidden => {}
|
||||
}
|
||||
|
||||
match self.probe_head_object(bucket, &prefix).await? {
|
||||
HeadObjectProbe::Found(output) => {
|
||||
return Ok(ResolvedPath::Directory {
|
||||
prefix,
|
||||
metadata: Some(output),
|
||||
});
|
||||
}
|
||||
HeadObjectProbe::Missing => {
|
||||
had_visibility = true;
|
||||
}
|
||||
HeadObjectProbe::Forbidden => {}
|
||||
}
|
||||
|
||||
if authorize_operation(&self.session_context, &S3Action::ListBucket, bucket, Some(&prefix))
|
||||
.await
|
||||
.is_ok()
|
||||
{
|
||||
had_visibility = true;
|
||||
if self.prefix_has_entries(bucket, &prefix).await? {
|
||||
return Ok(ResolvedPath::Directory { prefix, metadata: None });
|
||||
}
|
||||
}
|
||||
|
||||
if had_visibility {
|
||||
Err(FsError::NotFound)
|
||||
} else {
|
||||
Err(FsError::Forbidden)
|
||||
}
|
||||
}
|
||||
|
||||
/// Parse WebDAV path to bucket and object key
|
||||
fn parse_path(&self, path: &DavPath) -> Result<(String, Option<String>), FsError> {
|
||||
let path_str = path.as_url_string();
|
||||
@@ -534,6 +752,24 @@ where
|
||||
Ok(output) => {
|
||||
let mut entries = Vec::new();
|
||||
|
||||
// Collect common prefix base names for filtering
|
||||
let common_prefix_names: std::collections::HashSet<String> = output
|
||||
.common_prefixes
|
||||
.as_ref()
|
||||
.map(|prefixes| {
|
||||
prefixes
|
||||
.iter()
|
||||
.filter_map(|p| p.prefix.as_ref())
|
||||
.map(|p| {
|
||||
std::path::PathBuf::from(p.trim_end_matches('/'))
|
||||
.file_name()
|
||||
.map(|n| n.to_string_lossy().to_string())
|
||||
.unwrap_or_else(|| p.clone())
|
||||
})
|
||||
.collect()
|
||||
})
|
||||
.unwrap_or_default();
|
||||
|
||||
// Add files (objects)
|
||||
if let Some(objects) = output.contents {
|
||||
for obj in objects {
|
||||
@@ -555,6 +791,17 @@ where
|
||||
.unwrap_or_else(|| key.clone());
|
||||
|
||||
let size = obj.size.unwrap_or(0) as u64;
|
||||
|
||||
// Skip directory markers (keys ending with /)
|
||||
if key.ends_with('/') {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Skip 0-byte objects that match a directory name (Windows WebDAV duplicates)
|
||||
if size == 0 && common_prefix_names.contains(&filename) {
|
||||
continue;
|
||||
}
|
||||
|
||||
let modified = obj
|
||||
.last_modified
|
||||
.map(|dt| {
|
||||
@@ -762,18 +1009,8 @@ where
|
||||
}
|
||||
|
||||
if let Some(key) = key {
|
||||
// Get object metadata
|
||||
match self
|
||||
.storage
|
||||
.head_object(
|
||||
&bucket,
|
||||
&key,
|
||||
&self.session_context.principal.user_identity.credentials.access_key,
|
||||
&self.session_context.principal.user_identity.credentials.secret_key,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(output) => {
|
||||
return match self.resolve_path(&bucket, &key).await? {
|
||||
ResolvedPath::File(output) => {
|
||||
let size = output.content_length.unwrap_or(0) as u64;
|
||||
let modified = output
|
||||
.last_modified
|
||||
@@ -792,52 +1029,32 @@ where
|
||||
content_type: output.content_type.map(|c| c.to_string()),
|
||||
}) as Box<dyn DavMetaData>)
|
||||
}
|
||||
Err(e) => {
|
||||
// Check if it might be a "directory" (prefix)
|
||||
let prefix = format!("{}/", key);
|
||||
let list_input = ListObjectsV2Input::builder()
|
||||
.bucket(bucket.clone())
|
||||
.prefix(Some(prefix))
|
||||
.max_keys(Some(1))
|
||||
.build()
|
||||
.map_err(|_| FsError::GeneralFailure)?;
|
||||
ResolvedPath::Directory { metadata, .. } => {
|
||||
let modified = metadata
|
||||
.as_ref()
|
||||
.and_then(|output| output.last_modified.as_ref())
|
||||
.map(|dt| {
|
||||
let offset_dt: time::OffsetDateTime = dt.clone().into();
|
||||
SystemTime::from(offset_dt)
|
||||
})
|
||||
.unwrap_or_else(SystemTime::now);
|
||||
|
||||
match self
|
||||
.storage
|
||||
.list_objects_v2(
|
||||
list_input,
|
||||
&self.session_context.principal.user_identity.credentials.access_key,
|
||||
&self.session_context.principal.user_identity.credentials.secret_key,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(output) => {
|
||||
if output.contents.map(|c| !c.is_empty()).unwrap_or(false)
|
||||
|| output.common_prefixes.map(|c| !c.is_empty()).unwrap_or(false)
|
||||
{
|
||||
// It's a directory
|
||||
Ok(Box::new(WebDavMetaData {
|
||||
size: 0,
|
||||
modified: SystemTime::now(),
|
||||
created: SystemTime::now(),
|
||||
is_dir: true,
|
||||
etag: None,
|
||||
content_type: None,
|
||||
}) as Box<dyn DavMetaData>)
|
||||
} else {
|
||||
debug!("Object not found: {}/{}: {}", bucket, key, e);
|
||||
Err(FsError::NotFound)
|
||||
}
|
||||
}
|
||||
Err(_) => {
|
||||
debug!("Object not found: {}/{}: {}", bucket, key, e);
|
||||
Err(FsError::NotFound)
|
||||
}
|
||||
}
|
||||
Ok(Box::new(WebDavMetaData {
|
||||
size: 0,
|
||||
modified,
|
||||
created: modified,
|
||||
is_dir: true,
|
||||
etag: metadata.as_ref().and_then(|output| output.e_tag.as_ref().map(etag_to_string)),
|
||||
content_type: metadata.and_then(|output| output.content_type.map(|c| c.to_string())),
|
||||
}) as Box<dyn DavMetaData>)
|
||||
}
|
||||
}
|
||||
};
|
||||
} else {
|
||||
// Get bucket metadata
|
||||
authorize_operation(&self.session_context, &S3Action::HeadBucket, &bucket, None)
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
|
||||
match self
|
||||
.storage
|
||||
.head_bucket(
|
||||
@@ -1074,9 +1291,141 @@ where
|
||||
.boxed()
|
||||
}
|
||||
|
||||
fn rename<'a>(&'a self, _from: &'a DavPath, _to: &'a DavPath) -> FsFuture<'a, ()> {
|
||||
// S3 doesn't support native rename, would need copy + delete
|
||||
async move { Err(FsError::NotImplemented) }.boxed()
|
||||
fn rename<'a>(&'a self, from: &'a DavPath, to: &'a DavPath) -> FsFuture<'a, ()> {
|
||||
async move {
|
||||
let (src_bucket, src_key) = self.parse_path(from)?;
|
||||
let (dst_bucket, dst_key) = self.parse_path(to)?;
|
||||
|
||||
if src_bucket.is_empty() || dst_bucket.is_empty() {
|
||||
return Err(FsError::Forbidden);
|
||||
}
|
||||
|
||||
let src_key = src_key.ok_or(FsError::Forbidden)?;
|
||||
let dst_key = dst_key.ok_or(FsError::Forbidden)?;
|
||||
let (access_key, secret_key) = self.credentials();
|
||||
let resolved_src = self.resolve_path(&src_bucket, &src_key).await?;
|
||||
let (src_prefix, include_src_marker) = match resolved_src {
|
||||
ResolvedPath::File(_) => {
|
||||
authorize_operation(&self.session_context, &S3Action::GetObject, &src_bucket, Some(&src_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
authorize_operation(&self.session_context, &S3Action::PutObject, &dst_bucket, Some(&dst_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
authorize_operation(&self.session_context, &S3Action::DeleteObject, &src_bucket, Some(&src_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
|
||||
self.copy_object_streaming(&src_bucket, &src_key, &dst_bucket, &dst_key)
|
||||
.await?;
|
||||
|
||||
self.storage
|
||||
.delete_object(&src_bucket, &src_key, access_key, secret_key)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Failed to delete source object after rename: {}", e);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
|
||||
debug!("Successfully renamed file '{}/{}' to '{}/{}'", src_bucket, src_key, dst_bucket, dst_key);
|
||||
return Ok(());
|
||||
}
|
||||
ResolvedPath::Directory { prefix, .. } => {
|
||||
let include_src_marker =
|
||||
matches!(self.probe_head_object(&src_bucket, &src_key).await?, HeadObjectProbe::Found(_));
|
||||
(prefix, include_src_marker)
|
||||
}
|
||||
};
|
||||
let dst_prefix = format!("{}/", dst_key);
|
||||
|
||||
authorize_operation(&self.session_context, &S3Action::ListBucket, &src_bucket, Some(&src_prefix))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
|
||||
let mut continuation_token: Option<String> = None;
|
||||
let mut renamed_any = false;
|
||||
|
||||
if include_src_marker {
|
||||
authorize_operation(&self.session_context, &S3Action::GetObject, &src_bucket, Some(&src_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
authorize_operation(&self.session_context, &S3Action::PutObject, &dst_bucket, Some(&dst_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
authorize_operation(&self.session_context, &S3Action::DeleteObject, &src_bucket, Some(&src_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
|
||||
self.execute_directory_rename_pairs(&src_bucket, &dst_bucket, &[(src_key.clone(), dst_key.clone())])
|
||||
.await?;
|
||||
renamed_any = true;
|
||||
}
|
||||
|
||||
loop {
|
||||
let mut list_builder = ListObjectsV2Input::builder()
|
||||
.bucket(src_bucket.clone())
|
||||
.prefix(Some(src_prefix.clone()));
|
||||
|
||||
if let Some(ref token) = continuation_token {
|
||||
list_builder = list_builder.continuation_token(Some(token.clone()));
|
||||
}
|
||||
|
||||
let list_input = list_builder.build().map_err(|_| FsError::GeneralFailure)?;
|
||||
let output = self
|
||||
.storage
|
||||
.list_objects_v2(list_input, access_key, secret_key)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Failed to list objects during directory rename: {}", e);
|
||||
FsError::GeneralFailure
|
||||
})?;
|
||||
|
||||
let mut page_pairs: Vec<(String, String)> = Vec::new();
|
||||
if let Some(objects) = output.contents {
|
||||
for obj in objects {
|
||||
if let Some(obj_key) = obj.key {
|
||||
let new_key = obj_key.replacen(&src_prefix, &dst_prefix, 1);
|
||||
page_pairs.push((obj_key, new_key));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if !page_pairs.is_empty() {
|
||||
for (src_obj_key, dst_obj_key) in &page_pairs {
|
||||
authorize_operation(&self.session_context, &S3Action::GetObject, &src_bucket, Some(src_obj_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
authorize_operation(&self.session_context, &S3Action::PutObject, &dst_bucket, Some(dst_obj_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
authorize_operation(&self.session_context, &S3Action::DeleteObject, &src_bucket, Some(src_obj_key))
|
||||
.await
|
||||
.map_err(|_| FsError::Forbidden)?;
|
||||
}
|
||||
|
||||
self.execute_directory_rename_pairs(&src_bucket, &dst_bucket, &page_pairs)
|
||||
.await?;
|
||||
renamed_any = true;
|
||||
}
|
||||
|
||||
if !output.is_truncated.unwrap_or(false) {
|
||||
break;
|
||||
}
|
||||
continuation_token = output.next_continuation_token;
|
||||
}
|
||||
|
||||
if !renamed_any {
|
||||
debug!("Source not found: {}/{}", src_bucket, src_key);
|
||||
return Err(FsError::NotFound);
|
||||
}
|
||||
|
||||
debug!(
|
||||
"Successfully renamed directory '{}/{}' to '{}/{}'",
|
||||
src_bucket, src_key, dst_bucket, dst_key
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
.boxed()
|
||||
}
|
||||
|
||||
fn copy<'a>(&'a self, _from: &'a DavPath, _to: &'a DavPath) -> FsFuture<'a, ()> {
|
||||
@@ -1091,14 +1440,17 @@ mod tests {
|
||||
use crate::common::client::s3::StorageBackend as S3StorageBackend;
|
||||
use crate::common::session::{Protocol, ProtocolPrincipal, SessionContext};
|
||||
use async_trait::async_trait;
|
||||
use bytes::Bytes;
|
||||
use dav_server::davpath::DavPath;
|
||||
use dav_server::fs::FsError;
|
||||
use futures_util::StreamExt;
|
||||
use rustfs_credentials::Credentials;
|
||||
use rustfs_policy::auth::UserIdentity;
|
||||
use s3s::dto::*;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::fmt::{Debug, Formatter};
|
||||
use std::net::{IpAddr, Ipv4Addr};
|
||||
use std::sync::Arc;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
#[derive(Clone)]
|
||||
struct DummyStorage;
|
||||
@@ -1221,6 +1573,190 @@ mod tests {
|
||||
WebDavDriver::new(DummyStorage, Arc::new(session_context))
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct RecordingStorageState {
|
||||
objects: HashMap<(String, String), Vec<u8>>,
|
||||
put_keys: Vec<String>,
|
||||
delete_keys: Vec<String>,
|
||||
fail_delete_keys: HashSet<String>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct RecordingStorage {
|
||||
state: Arc<Mutex<RecordingStorageState>>,
|
||||
}
|
||||
|
||||
impl Debug for RecordingStorage {
|
||||
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
|
||||
f.write_str("RecordingStorage")
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl S3StorageBackend for RecordingStorage {
|
||||
type Error = std::io::Error;
|
||||
|
||||
async fn get_object(
|
||||
&self,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
_start_pos: Option<u64>,
|
||||
) -> Result<GetObjectOutput, Self::Error> {
|
||||
let data = self
|
||||
.state
|
||||
.lock()
|
||||
.expect("recording storage lock poisoned")
|
||||
.objects
|
||||
.get(&(bucket.to_string(), key.to_string()))
|
||||
.cloned()
|
||||
.ok_or_else(|| std::io::Error::new(std::io::ErrorKind::NotFound, "missing object"))?;
|
||||
|
||||
let content_length = data.len() as i64;
|
||||
let body =
|
||||
StreamingBlob::wrap(futures_util::stream::once(async move { Ok::<Bytes, std::io::Error>(Bytes::from(data)) }));
|
||||
|
||||
Ok(GetObjectOutput {
|
||||
body: Some(body),
|
||||
content_length: Some(content_length),
|
||||
..Default::default()
|
||||
})
|
||||
}
|
||||
|
||||
async fn get_object_range(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_key: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
_start_pos: u64,
|
||||
_length: u64,
|
||||
) -> Result<GetObjectOutput, Self::Error> {
|
||||
unreachable!("range reads are not used in rename regression tests")
|
||||
}
|
||||
|
||||
async fn put_object(
|
||||
&self,
|
||||
mut input: PutObjectInput,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<PutObjectOutput, Self::Error> {
|
||||
let bucket = input.bucket.clone();
|
||||
let key = input.key.clone();
|
||||
let mut bytes = Vec::new();
|
||||
|
||||
if let Some(mut body) = input.body.take() {
|
||||
while let Some(chunk) = body.next().await {
|
||||
let chunk = chunk.map_err(|e| std::io::Error::other(e.to_string()))?;
|
||||
bytes.extend_from_slice(&chunk);
|
||||
}
|
||||
}
|
||||
|
||||
let mut state = self.state.lock().expect("recording storage lock poisoned");
|
||||
state.put_keys.push(key.clone());
|
||||
state.objects.insert((bucket, key), bytes);
|
||||
|
||||
Ok(PutObjectOutput::default())
|
||||
}
|
||||
|
||||
async fn delete_object(
|
||||
&self,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<DeleteObjectOutput, Self::Error> {
|
||||
let mut state = self.state.lock().expect("recording storage lock poisoned");
|
||||
state.delete_keys.push(key.to_string());
|
||||
if state.fail_delete_keys.contains(key) {
|
||||
return Err(std::io::Error::other("injected delete failure"));
|
||||
}
|
||||
state.objects.remove(&(bucket.to_string(), key.to_string()));
|
||||
Ok(DeleteObjectOutput::default())
|
||||
}
|
||||
|
||||
async fn head_object(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_key: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<HeadObjectOutput, Self::Error> {
|
||||
unreachable!("head_object is not used in rename regression tests")
|
||||
}
|
||||
|
||||
async fn head_bucket(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<HeadBucketOutput, Self::Error> {
|
||||
unreachable!("head_bucket is not used in rename regression tests")
|
||||
}
|
||||
|
||||
async fn list_objects_v2(
|
||||
&self,
|
||||
_input: ListObjectsV2Input,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<ListObjectsV2Output, Self::Error> {
|
||||
unreachable!("list_objects_v2 is not used in rename regression tests")
|
||||
}
|
||||
|
||||
async fn list_buckets(&self, _access_key: &str, _secret_key: &str) -> Result<ListBucketsOutput, Self::Error> {
|
||||
unreachable!("list_buckets is not used in rename regression tests")
|
||||
}
|
||||
|
||||
async fn create_bucket(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<CreateBucketOutput, Self::Error> {
|
||||
unreachable!("create_bucket is not used in rename regression tests")
|
||||
}
|
||||
|
||||
async fn delete_bucket(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_access_key: &str,
|
||||
_secret_key: &str,
|
||||
) -> Result<DeleteBucketOutput, Self::Error> {
|
||||
unreachable!("delete_bucket is not used in rename regression tests")
|
||||
}
|
||||
}
|
||||
|
||||
fn recording_driver(
|
||||
initial_objects: &[(&str, &str, &[u8])],
|
||||
fail_delete_keys: &[&str],
|
||||
) -> (WebDavDriver<RecordingStorage>, RecordingStorage) {
|
||||
let storage = RecordingStorage::default();
|
||||
let identity = UserIdentity::new(Credentials {
|
||||
access_key: "ak".to_string(),
|
||||
secret_key: "sk".to_string(),
|
||||
..Default::default()
|
||||
});
|
||||
let session_context = SessionContext::new(
|
||||
ProtocolPrincipal::new(Arc::new(identity)),
|
||||
Protocol::WebDav,
|
||||
IpAddr::V4(Ipv4Addr::LOCALHOST),
|
||||
);
|
||||
|
||||
let state = RecordingStorageState {
|
||||
objects: initial_objects
|
||||
.iter()
|
||||
.map(|(bucket, key, body)| (((*bucket).to_string(), (*key).to_string()), body.to_vec()))
|
||||
.collect(),
|
||||
fail_delete_keys: fail_delete_keys.iter().map(|key| (*key).to_string()).collect(),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
*storage.state.lock().expect("recording storage lock poisoned") = state;
|
||||
|
||||
(WebDavDriver::new(storage.clone(), Arc::new(session_context)), storage)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_path_decodes_url_encoded_object_names() {
|
||||
let driver = driver();
|
||||
@@ -1241,4 +1777,98 @@ mod tests {
|
||||
|
||||
assert_eq!(err, FsError::GeneralFailure);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_path_handles_directory_paths_with_trailing_slash() {
|
||||
let driver = driver();
|
||||
let path = DavPath::new("/bucket/folder/").expect("path should parse");
|
||||
|
||||
let (bucket, key) = driver.parse_path(&path).expect("path should decode");
|
||||
|
||||
assert_eq!(bucket, "bucket");
|
||||
assert_eq!(key.as_deref(), Some("folder"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_path_handles_chinese_directory_names() {
|
||||
let driver = driver();
|
||||
let path = DavPath::new("/bucket/%E6%96%B0%E5%BB%BA%E6%96%87%E4%BB%B6%E5%A4%B9%20(4)").expect("path should parse");
|
||||
|
||||
let (bucket, key) = driver.parse_path(&path).expect("path should decode");
|
||||
|
||||
assert_eq!(bucket, "bucket");
|
||||
assert_eq!(key.as_deref(), Some("新建文件夹 (4)"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_path_handles_nested_paths() {
|
||||
let driver = driver();
|
||||
let path = DavPath::new("/bucket/dir/subdir/file.txt").expect("path should parse");
|
||||
|
||||
let (bucket, key) = driver.parse_path(&path).expect("path should decode");
|
||||
|
||||
assert_eq!(bucket, "bucket");
|
||||
assert_eq!(key.as_deref(), Some("dir/subdir/file.txt"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_path_returns_none_key_for_bucket_root() {
|
||||
let driver = driver();
|
||||
let path = DavPath::new("/bucket/").expect("path should parse");
|
||||
|
||||
let (bucket, key) = driver.parse_path(&path).expect("path should decode");
|
||||
|
||||
assert_eq!(bucket, "bucket");
|
||||
assert!(key.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_path_handles_url_encoded_spaces_in_object_name() {
|
||||
let driver = driver();
|
||||
let path = DavPath::new("/bucket/file%20with%20spaces.txt").expect("path should parse");
|
||||
|
||||
let (bucket, key) = driver.parse_path(&path).expect("path should decode");
|
||||
|
||||
assert_eq!(bucket, "bucket");
|
||||
assert_eq!(key.as_deref(), Some("file with spaces.txt"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn directory_rename_returns_error_when_delete_fails_after_successful_copy() {
|
||||
let (driver, storage) = recording_driver(
|
||||
&[
|
||||
("bucket", "src/file-a.txt", b"file-a"),
|
||||
("bucket", "src/file-b.txt", b"file-b"),
|
||||
],
|
||||
&["src/file-a.txt"],
|
||||
);
|
||||
|
||||
let err = driver
|
||||
.execute_directory_rename_pairs(
|
||||
"bucket",
|
||||
"bucket",
|
||||
&[
|
||||
("src/file-a.txt".to_string(), "dst/file-a.txt".to_string()),
|
||||
("src/file-b.txt".to_string(), "dst/file-b.txt".to_string()),
|
||||
],
|
||||
)
|
||||
.await
|
||||
.expect_err("delete failure should be surfaced");
|
||||
|
||||
assert_eq!(err, FsError::GeneralFailure);
|
||||
|
||||
let state = storage.state.lock().expect("recording storage lock poisoned");
|
||||
assert_eq!(state.put_keys, vec!["dst/file-a.txt".to_string(), "dst/file-b.txt".to_string()]);
|
||||
assert_eq!(state.delete_keys, vec!["src/file-a.txt".to_string()]);
|
||||
assert!(
|
||||
state
|
||||
.objects
|
||||
.contains_key(&("bucket".to_string(), "dst/file-a.txt".to_string()))
|
||||
);
|
||||
assert!(
|
||||
state
|
||||
.objects
|
||||
.contains_key(&("bucket".to_string(), "dst/file-b.txt".to_string()))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user