diff --git a/crates/e2e_test/src/protocols/README.md b/crates/e2e_test/src/protocols/README.md index 9506b57b1..9a687acab 100644 --- a/crates/e2e_test/src/protocols/README.md +++ b/crates/e2e_test/src/protocols/README.md @@ -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 diff --git a/crates/e2e_test/src/protocols/test_runner.rs b/crates/e2e_test/src/protocols/test_runner.rs index 631aa28f3..5d2c73ab0 100644 --- a/crates/e2e_test/src/protocols/test_runner.rs +++ b/crates/e2e_test/src/protocols/test_runner.rs @@ -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> { let suite = ProtocolTestSuite::new(); let results = suite.run_test_suite().await; diff --git a/crates/e2e_test/src/protocols/webdav_core.rs b/crates/e2e_test/src/protocols/webdav_core.rs index db5e8506c..bb64127d2 100644 --- a/crates/e2e_test/src/protocols/webdav_core.rs +++ b/crates/e2e_test/src/protocols/webdav_core.rs @@ -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>, + content_type: Option<&str>, +) -> Result { + let uri = url.parse::()?; + 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 +} diff --git a/crates/protocols/src/webdav/driver.rs b/crates/protocols/src/webdav/driver.rs index abbe85cc6..f31654f53 100644 --- a/crates/protocols/src/webdav/driver.rs +++ b/crates/protocols/src/webdav/driver.rs @@ -397,6 +397,20 @@ where session_context: Arc, } +enum ResolvedPath { + File(Box), + Directory { + prefix: String, + metadata: Option>, + }, +} + +enum HeadObjectProbe { + Forbidden, + Missing, + Found(Box), +} + impl Debug for WebDavDriver 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 { + 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 { + 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 { + 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), 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 = 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) } - 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) - } 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) } - } + }; } 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 = 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>, + put_keys: Vec, + delete_keys: Vec, + fail_delete_keys: HashSet, + } + + #[derive(Clone, Default)] + struct RecordingStorage { + state: Arc>, + } + + 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, + ) -> Result { + 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::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 { + 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 { + 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 { + 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 { + unreachable!("head_object is not used in rename regression tests") + } + + async fn head_bucket( + &self, + _bucket: &str, + _access_key: &str, + _secret_key: &str, + ) -> Result { + 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 { + unreachable!("list_objects_v2 is not used in rename regression tests") + } + + async fn list_buckets(&self, _access_key: &str, _secret_key: &str) -> Result { + unreachable!("list_buckets is not used in rename regression tests") + } + + async fn create_bucket( + &self, + _bucket: &str, + _access_key: &str, + _secret_key: &str, + ) -> Result { + unreachable!("create_bucket is not used in rename regression tests") + } + + async fn delete_bucket( + &self, + _bucket: &str, + _access_key: &str, + _secret_key: &str, + ) -> Result { + unreachable!("delete_bucket is not used in rename regression tests") + } + } + + fn recording_driver( + initial_objects: &[(&str, &str, &[u8])], + fail_delete_keys: &[&str], + ) -> (WebDavDriver, 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())) + ); + } }